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, Clone)]
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: crate::feed::bound_guid(e.guid),
url: e.url,
title: Some(e.title),
author: None,
published: e.published,
content_html: e.summary,
fetched_at: None,
keep_stored_content: false,
}
}
}
#[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
.map(|n| crate::feed::bound_text(n, crate::feed::MAX_TITLE_BYTES)),
url: crate::feed::bound_text(url, crate::feed::MAX_URL_BYTES),
},
))
}
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: crate::feed::bound_text(doc.title, crate::feed::MAX_TITLE_BYTES),
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))
.map(|u| crate::feed::bound_text(u, crate::feed::MAX_URL_BYTES)),
summary: non_blank(doc.description)
.or_else(|| non_blank(doc.text_content))
.map(|raw| {
crate::feed::plain_text_to_html_bounded(
&raw,
crate::feed::MAX_CONTENT_HTML_BYTES,
)
}),
})
})
.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(usize),
Sibling,
Orphan,
Malformed,
}
fn classify_document(
record: &crate::atproto::RecordEntry,
wanted: &std::collections::HashMap<String, usize>,
known: &std::collections::HashSet<&str>,
) -> DocumentFate {
match serde_json::from_value::<DocumentValue>(record.value.clone()) {
Ok(doc) if wanted.contains_key(&doc.site) => DocumentFate::Keep(wanted[&doc.site]),
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> {
if uri.collection != nsid::STANDARD_PUBLICATION {
return Err(
NotAPublication(format!("{uri} is not a {} URI", nsid::STANDARD_PUBLICATION)).into(),
);
}
fetch_repo(
http,
plc_directory,
&uri.authority,
std::slice::from_ref(&uri.rkey),
)
.await?
.pop()
.unwrap_or_else(|| Err(NotAPublication(format!("{uri} was not read")).into()))
}
#[derive(Debug)]
pub struct NotAPublication(pub String);
impl std::fmt::Display for NotAPublication {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
impl std::error::Error for NotAPublication {}
pub async fn fetch_repo(
http: &reqwest::Client,
plc_directory: &str,
did: &str,
rkeys: &[String],
) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
fetch_repo_capped(
http,
plc_directory,
did,
rkeys,
crate::atproto::MAX_LARGE_RECORDS,
crate::atproto::MAX_LIST_BYTES,
)
.await
}
pub(crate) async fn fetch_repo_capped(
http: &reqwest::Client,
plc_directory: &str,
did: &str,
requested: &[String],
per_site_cap: usize,
budget_bytes: usize,
) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
use anyhow::Context;
if requested.is_empty() {
return Ok(Vec::new());
}
let mut rkeys: Vec<String> = Vec::new();
for r in requested {
if !rkeys.contains(r) {
rkeys.push(r.clone());
}
}
let rkeys = rkeys.as_slice();
let pds = crate::atproto::resolve_did_to_pds(http, plc_directory, did)
.await
.with_context(|| format!("resolving the PDS for {did}"))?;
let client = crate::atproto::PdsClient::anonymous(http.clone(), pds, did.to_string());
let mut budget = crate::atproto::ByteBudget::new(budget_bytes);
let (publications, skipped_publications) = client
.list_all_records_skipping_within(nsid::STANDARD_PUBLICATION, &mut budget)
.await
.with_context(|| format!("listing publications for {did}"))?;
if skipped_publications > 0 {
tracing::warn!(
repo = %did,
skipped = skipped_publications,
"skipped malformed publication records in this repo"
);
}
let wanted: Vec<Option<(String, Publication)>> = rkeys
.iter()
.map(|rkey| publication_from_records(rkey, &publications))
.collect();
let index_of: std::collections::HashMap<String, usize> = wanted
.iter()
.enumerate()
.filter_map(|(i, w)| w.as_ref().map(|(site, _)| (site.clone(), i)))
.collect();
let not_found = |rkey: &str| {
anyhow::Error::new(NotAPublication(format!(
"at://{did}/{}/{rkey} is not a readable site.standard.publication",
nsid::STANDARD_PUBLICATION
)))
};
if index_of.is_empty() {
return Ok(requested.iter().map(|r| Err(not_found(r))).collect());
}
let known: std::collections::HashSet<&str> =
publications.iter().map(|p| p.uri.as_str()).collect();
let mut kept_per = vec![0usize; rkeys.len()];
let mut capped = vec![false; rkeys.len()];
let mut orphaned = 0usize;
let documents = client
.list_recent_matching_within(
nsid::STANDARD_DOCUMENT,
per_site_cap.saturating_mul(index_of.len()),
&mut budget,
DOCUMENT_PAGE_SIZE,
|record| match classify_document(record, &index_of, &known) {
DocumentFate::Keep(i) if kept_per[i] < per_site_cap => {
kept_per[i] += 1;
true
}
DocumentFate::Keep(i) => {
capped[i] = true;
false
}
DocumentFate::Orphan => {
orphaned += 1;
false
}
DocumentFate::Sibling | DocumentFate::Malformed => false,
},
)
.await
.with_context(|| format!("listing documents for {did}"))?;
let mut per_site: Vec<Vec<crate::atproto::RecordEntry>> = vec![Vec::new(); rkeys.len()];
for record in documents.records {
let site = serde_json::from_value::<DocumentValue>(record.value.clone()).map(|d| d.site);
if let Some(&i) = site.ok().as_deref().and_then(|s| index_of.get(s)) {
per_site[i].push(record);
}
}
let reads: Vec<anyhow::Result<PublicationRead>> = wanted
.into_iter()
.zip(per_site)
.enumerate()
.map(|(i, (w, records))| match w {
None => Err(not_found(&rkeys[i])),
Some((site, publication)) => {
let walk = crate::atproto::RecordWalk {
records,
complete: documents.complete && !capped[i],
malformed: documents.malformed,
};
Ok(read_from(publication, &site, walk, orphaned))
}
})
.collect();
let mut reads: Vec<Option<anyhow::Result<PublicationRead>>> =
reads.into_iter().map(Some).collect();
let positions: Vec<usize> = requested
.iter()
.map(|r| {
rkeys
.iter()
.position(|k| k == r)
.expect("every requested rkey is in rkeys")
})
.collect();
Ok(positions
.iter()
.enumerate()
.map(|(n, &u)| {
let repeated_later = positions[n + 1..].contains(&u);
match (&reads[u], repeated_later) {
(Some(Ok(read)), true) => Ok(read.clone()),
(Some(Err(_)), _) | (None, _) => Err(not_found(&requested[n])),
(Some(Ok(_)), false) => match reads[u].take() {
Some(Ok(read)) => Ok(read),
_ => Err(not_found(&requested[n])),
},
}
})
.collect())
}
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 documents.malformed > 0 {
tracing::warn!(
site = %canonical_site,
skipped = documents.malformed,
"skipped malformed document records for this publication"
);
}
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)]
pub(crate) 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,
malformed: 0,
};
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",
}),
)
}
pub(crate) async fn serve_repo(
did: &'static str,
records: Vec<(&'static str, &'static str, serde_json::Value)>,
) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
use axum::extract::{Query, Request};
use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let port = addr.port();
let (plc_host, pds_host) = (
format!("plc-{port}.repo.test"),
format!("pds-{port}.repo.test"),
);
for host in [&plc_host, &pds_host] {
crate::net::test_host_override(host, addr);
}
let hits = Arc::new(AtomicUsize::new(0));
let counter = Arc::clone(&hits);
let records = Arc::new(records);
let app = axum::Router::new().fallback(
move |Query(q): Query<HashMap<String, String>>, req: Request| {
let records = Arc::clone(&records);
let counter = Arc::clone(&counter);
async move {
let path = req.uri().path().to_string();
if path == format!("/{did}") {
return axum::Json(json!({
"id": did,
"service": [{
"id": "#atproto_pds",
"type": "AtprotoPersonalDataServer",
"serviceEndpoint": format!("http://{pds_host}:{port}"),
}],
}));
}
assert_eq!(path, "/xrpc/com.atproto.repo.listRecords", "unexpected request");
counter.fetch_add(1, Ordering::SeqCst);
let collection = q.get("collection").cloned().unwrap_or_default();
let limit: usize = q.get("limit").and_then(|l| l.parse().ok()).unwrap_or(50);
let offset: usize = q.get("cursor").and_then(|c| c.parse().ok()).unwrap_or(0);
let page: Vec<_> = records
.iter()
.filter(|(c, _, _)| *c == collection)
.skip(offset)
.take(limit)
.map(|(c, rkey, value)| {
if rkey.is_empty() {
return json!({ "cid": "bafy", "value": value });
}
json!({ "uri": format!("at://{did}/{c}/{rkey}"), "cid": "bafy", "value": value })
})
.collect();
let next = (offset + page.len()).to_string();
axum::Json(json!({ "records": page, "cursor": next }))
}
},
);
tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
(format!("http://{plc_host}:{port}"), hits)
}
#[tokio::test]
async fn fetch_reads_a_publication_from_a_mocked_repo() {
const OWN: &str = "did:plc:fetchmock";
let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
let sibling = format!("at://{OWN}/{}/other", nsid::STANDARD_PUBLICATION);
let doc = |site: &str, title: &str| {
json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
"path": "/p", "site": site })
};
let (plc, hits) = serve_repo(
OWN,
vec![
(
nsid::STANDARD_PUBLICATION,
"mine",
json!({ "name": "Mine", "url": "https://mine.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"other",
json!({ "name": "Other", "url": "https://other.example" }),
),
(nsid::STANDARD_DOCUMENT, "3l2fmaaaaaa2a", doc(&site, "kept")),
(
nsid::STANDARD_DOCUMENT,
"3l2fmaaaaaa2b",
doc(&sibling, "sibling's"),
),
],
)
.await;
let client = crate::feed::build_client().unwrap();
let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
.await
.unwrap();
assert!(read.complete, "the walk did not finish");
assert_eq!(read.publication.name.as_deref(), Some("Mine"));
let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
assert_eq!(
titles,
vec!["kept"],
"the site filter let a sibling through"
);
assert_eq!(
hits.load(std::sync::atomic::Ordering::SeqCst),
4,
"two collections, each a full page then an empty one"
);
}
#[tokio::test]
async fn a_failed_publication_read_is_a_poll_failure() {
let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
let url = format!(
"at://did:plc:unreachableaaaaaaaaaaaaa/{}/x",
nsid::STANDARD_PUBLICATION
);
crate::store::upsert_feed(
&pool,
&crate::store::NewFeed {
url: url.clone(),
..Default::default()
},
)
.await
.unwrap();
let feed = crate::store::get_feed_by_url(&pool, &url)
.await
.unwrap()
.unwrap();
let mut config = crate::config::Config::default();
config.oauth.plc_directory = "http://plc.nowhere.invalid".into();
let client = crate::feed::build_client().unwrap();
let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
.await
.expect("a source failure surfaced as a store error");
assert!(
matches!(
outcome,
crate::feed::PollOutcome::Failed {
kind: crate::feed::FailureKind::Fetch,
..
}
),
"an unreachable publication was not a fetch failure: {outcome:?}"
);
}
#[tokio::test]
async fn a_deleted_publication_record_is_not_an_unreachable_publisher() {
const GONE: &str = "did:plc:goneaaaaaaaaaaaaaaaaaaaa";
let site = format!("at://{GONE}/{}/deleted", nsid::STANDARD_PUBLICATION);
let (plc, _) = serve_repo(
GONE,
vec![(
nsid::STANDARD_PUBLICATION,
"another",
json!({ "name": "Other", "url": "https://o.example" }),
)],
)
.await;
let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
crate::store::upsert_feed(
&pool,
&crate::store::NewFeed {
url: site.clone(),
..Default::default()
},
)
.await
.unwrap();
let feed = crate::store::get_feed_by_url(&pool, &site)
.await
.unwrap()
.unwrap();
let mut config = crate::config::Config::default();
config.oauth.plc_directory = plc;
let client = crate::feed::build_client().unwrap();
let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
.await
.unwrap();
assert!(
matches!(
outcome,
crate::feed::PollOutcome::Failed {
kind: crate::feed::FailureKind::Parse,
..
}
),
"a deleted publication was filed as a network failure: {outcome:?}"
);
}
async fn serve_answering(
did: &'static str,
plc_status: u16,
pds: (u16, &'static str),
delay: std::time::Duration,
endless: bool,
) -> String {
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let port = addr.port();
let (plc_host, pds_host) = (
format!("plc-{port}.answer.test"),
format!("pds-{port}.answer.test"),
);
crate::net::test_host_override(&plc_host, addr);
crate::net::test_host_override(&pds_host, addr);
let pages = Arc::new(AtomicUsize::new(0));
let endpoint = format!("http://{pds_host}:{port}");
let app = axum::Router::new().fallback(move |req: axum::extract::Request| {
let pages = Arc::clone(&pages);
let endpoint = endpoint.clone();
async move {
use axum::response::IntoResponse;
if req.uri().path() == format!("/{did}") {
let status = axum::http::StatusCode::from_u16(plc_status).unwrap();
let doc = json!({ "id": did, "service": [{ "id": "#atproto_pds",
"type": "AtprotoPersonalDataServer", "serviceEndpoint": endpoint }] });
return (status, axum::Json(doc)).into_response();
}
tokio::time::sleep(delay).await;
if endless {
let n = pages.fetch_add(1, Ordering::SeqCst);
let body = json!({ "records": [{
"uri": format!("at://{did}/{}/3lend{n:08}", nsid::STANDARD_DOCUMENT),
"cid": "b",
"value": { "title": "x", "path": "/x", "publishedAt": "2026-07-11T00:00:00Z",
"site": format!("at://{did}/{}/other", nsid::STANDARD_PUBLICATION) } }],
"cursor": format!("c{n}") });
return axum::Json(body).into_response();
}
let status = axum::http::StatusCode::from_u16(pds.0).unwrap();
(status, [("content-type", "application/json")], pds.1).into_response()
}
});
tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
format!("http://{plc_host}:{port}")
}
async fn poll_publication_at(
did: &str,
plc: String,
deadline: Option<std::time::Duration>,
) -> crate::feed::PollOutcome {
let site = format!("at://{did}/{}/mine", nsid::STANDARD_PUBLICATION);
let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
crate::store::upsert_feed(
&pool,
&crate::store::NewFeed {
url: site.clone(),
..Default::default()
},
)
.await
.unwrap();
let feed = crate::store::get_feed_by_url(&pool, &site)
.await
.unwrap()
.unwrap();
let mut config = crate::config::Config::default();
config.oauth.plc_directory = plc;
if let Some(d) = deadline {
config.publication_read_deadline = d;
}
let client = crate::feed::build_client().unwrap();
crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
.await
.unwrap()
}
fn kind_of(outcome: &crate::feed::PollOutcome) -> Option<crate::feed::FailureKind> {
match outcome {
crate::feed::PollOutcome::Failed { kind, .. } => Some(*kind),
_ => None,
}
}
#[tokio::test]
async fn a_publication_read_has_an_overall_deadline() {
const DID: &str = "did:plc:slowrepoaaaaaaaaaaaaaaaa";
let plc = serve_answering(
DID,
200,
(200, ""),
std::time::Duration::from_millis(50),
true,
)
.await;
let started = std::time::Instant::now();
let outcome =
poll_publication_at(DID, plc, Some(std::time::Duration::from_millis(300))).await;
assert!(
started.elapsed() < std::time::Duration::from_secs(3),
"the read ran {:?}",
started.elapsed()
);
assert_eq!(
kind_of(&outcome),
Some(crate::feed::FailureKind::Fetch),
"{outcome:?}"
);
}
#[tokio::test]
async fn a_publication_failure_is_filed_under_what_happened() {
let zero = std::time::Duration::ZERO;
const GONE: &str = "did:plc:tombstonedaaaaaaaaaaaaaa";
let plc = serve_answering(GONE, 404, (200, ""), zero, false).await;
let outcome = poll_publication_at(GONE, plc, None).await;
assert_eq!(
kind_of(&outcome),
Some(crate::feed::FailureKind::Status),
"PLC 404: {outcome:?}"
);
const NOREPO: &str = "did:plc:norepoaaaaaaaaaaaaaaaaaa";
let plc = serve_answering(
NOREPO,
200,
(400, r#"{"error":"RepoNotFound"}"#),
zero,
false,
)
.await;
let outcome = poll_publication_at(NOREPO, plc, None).await;
assert_eq!(
kind_of(&outcome),
Some(crate::feed::FailureKind::Status),
"RepoNotFound: {outcome:?}"
);
const GARBLED: &str = "did:plc:garbledaaaaaaaaaaaaaaaaa";
let plc = serve_answering(GARBLED, 200, (200, r#"{"records":"x"}"#), zero, false).await;
let outcome = poll_publication_at(GARBLED, plc, None).await;
assert_eq!(
kind_of(&outcome),
Some(crate::feed::FailureKind::Parse),
"garbled body: {outcome:?}"
);
}
#[tokio::test]
async fn an_empty_publication_is_a_healthy_poll() {
const EMPTY: &str = "did:plc:emptypubaaaaaaaaaaaaaaaa";
let site = format!("at://{EMPTY}/{}/quiet", nsid::STANDARD_PUBLICATION);
let (plc, _) = serve_repo(
EMPTY,
vec![(
nsid::STANDARD_PUBLICATION,
"quiet",
json!({ "name": "Quiet", "url": "https://quiet.example" }),
)],
)
.await;
let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
crate::store::upsert_feed(
&pool,
&crate::store::NewFeed {
url: site.clone(),
..Default::default()
},
)
.await
.unwrap();
let feed = crate::store::get_feed_by_url(&pool, &site)
.await
.unwrap()
.unwrap();
let mut config = crate::config::Config::default();
config.oauth.plc_directory = plc;
let client = crate::feed::build_client().unwrap();
let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
.await
.unwrap();
assert!(
matches!(
outcome,
crate::feed::PollOutcome::Updated { new_entries: 0 }
),
"an empty publication was not a healthy poll: {outcome:?}"
);
}
#[tokio::test]
async fn fetch_skips_malformed_records_beside_good_ones() {
const OWN: &str = "did:plc:malformedrepo";
let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
let doc = |title: &str| {
json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
"path": "/p", "site": site })
};
let (plc, _) = serve_repo(
OWN,
vec![
(
nsid::STANDARD_PUBLICATION,
"",
json!({ "name": "Broken", "url": "https://x.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"mine",
json!({ "name": "Mine", "url": "https://mine.example" }),
),
(nsid::STANDARD_DOCUMENT, "", doc("unreadable")),
(nsid::STANDARD_DOCUMENT, "3l2mfaaaaaa2a", doc("kept")),
],
)
.await;
let client = crate::feed::build_client().unwrap();
let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
.await
.expect("a malformed record stalled a stranger's publication");
assert!(read.complete);
let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
assert_eq!(titles, vec!["kept"]);
}
const SHARED: &str = "did:plc:sharedrepoaaaaaaaaaaaaaa";
fn shared_doc(site_rkey: &str, title: &str, at: &str) -> serde_json::Value {
json!({ "title": title, "publishedAt": at, "path": format!("/{title}"),
"site": format!("at://{SHARED}/{}/{site_rkey}", nsid::STANDARD_PUBLICATION) })
}
#[tokio::test]
async fn two_publications_in_one_repo_cost_one_walk() {
let (plc, hits) = serve_repo(
SHARED,
vec![
(
nsid::STANDARD_PUBLICATION,
"alpha",
json!({ "name": "Alpha", "url": "https://alpha.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"beta",
json!({ "name": "Beta", "url": "https://beta.example" }),
),
(
nsid::STANDARD_DOCUMENT,
"3l2shaaaaaa2a",
shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
),
(
nsid::STANDARD_DOCUMENT,
"3l2shaaaaaa2b",
shared_doc("beta", "b1", "2026-07-10T00:00:00Z"),
),
(
nsid::STANDARD_DOCUMENT,
"3l2shaaaaaa2c",
shared_doc("alpha", "a2", "2026-07-09T00:00:00Z"),
),
],
)
.await;
let client = crate::feed::build_client().unwrap();
let reads = fetch_repo(
&client,
&plc,
SHARED,
&["alpha".to_string(), "beta".to_string()],
)
.await
.unwrap();
assert_eq!(
hits.load(std::sync::atomic::Ordering::SeqCst),
4,
"two collections walked once each (a page then an empty page), not once per publication"
);
let titles = |r: &anyhow::Result<PublicationRead>| {
let mut t: Vec<String> = r
.as_ref()
.unwrap()
.entries
.iter()
.map(|e| e.title.clone())
.collect();
t.sort();
t
};
assert_eq!(titles(&reads[0]), vec!["a1", "a2"]);
assert_eq!(titles(&reads[1]), vec!["b1"]);
}
#[tokio::test]
async fn a_busy_publication_does_not_starve_its_quiet_sibling() {
let mut records = vec![
(
nsid::STANDARD_PUBLICATION,
"busy",
json!({ "name": "Busy", "url": "https://busy.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"quiet",
json!({ "name": "Quiet", "url": "https://quiet.example" }),
),
];
for (rkey, title) in [
("3l2bsaaaaaa2a", "b1"),
("3l2bsaaaaaa2b", "b2"),
("3l2bsaaaaaa2c", "b3"),
("3l2bsaaaaaa2d", "b4"),
("3l2bsaaaaaa2e", "b5"),
] {
records.push((
nsid::STANDARD_DOCUMENT,
rkey,
shared_doc("busy", title, "2026-07-11T00:00:00Z"),
));
}
records.push((
nsid::STANDARD_DOCUMENT,
"3l2bsaaaaaa2f",
shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
));
let (plc, _) = serve_repo(SHARED, records).await;
let client = crate::feed::build_client().unwrap();
let reads = fetch_repo_capped(
&client,
&plc,
SHARED,
&["busy".to_string(), "quiet".to_string()],
2,
crate::atproto::MAX_LIST_BYTES,
)
.await
.unwrap();
let busy = reads[0].as_ref().unwrap();
let quiet = reads[1].as_ref().unwrap();
assert_eq!(
busy.entries.len(),
2,
"the busy publication was not capped at its own cap"
);
assert!(!busy.complete, "a capped publication was reported complete");
assert_eq!(quiet.entries.len(), 1, "the quiet sibling was starved");
assert!(quiet.complete);
}
#[tokio::test]
#[ignore = "known limitation: a one-repo group shares one byte budget (#229)"]
async fn big_siblings_do_not_spend_a_quiet_publications_share() {
let body = "w".repeat(20 * 1024);
let mut records = vec![
(
nsid::STANDARD_PUBLICATION,
"big1",
json!({ "name": "Big 1", "url": "https://b1.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"big2",
json!({ "name": "Big 2", "url": "https://b2.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"quiet",
json!({ "name": "Quiet", "url": "https://q.example" }),
),
];
for i in 0..200 {
let rkey: &'static str = Box::leak(format!("3l2big{i:06}").into_boxed_str());
let site = if i % 2 == 0 { "big1" } else { "big2" };
let mut doc = shared_doc(site, &format!("d{i}"), "2026-07-11T00:00:00Z");
doc["textContent"] = json!(body);
records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
}
records.push((
nsid::STANDARD_DOCUMENT,
"3l2zzzzzzzzzz",
shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
));
let (plc, _) = serve_repo(SHARED, records).await;
let client = crate::feed::build_client().unwrap();
let rkeys: Vec<String> = ["big1", "big2", "quiet"]
.iter()
.map(|s| s.to_string())
.collect();
let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 4 * 1024 * 1024)
.await
.unwrap();
let quiet = reads[2].as_ref().unwrap();
assert_eq!(
quiet.entries.len(),
1,
"the quiet publication was starved of bytes by its siblings"
);
assert!(
!reads[0].as_ref().unwrap().complete,
"a publication over its share was reported complete"
);
}
#[tokio::test]
async fn a_busy_publication_beside_idle_siblings_reads_as_it_would_alone() {
let body = "w".repeat(20 * 1024);
let mut records: Vec<(&'static str, &'static str, serde_json::Value)> = Vec::new();
let rkeys: Vec<String> = (0..16).map(|i| format!("p{i:02}")).collect();
for r in &rkeys {
let r: &'static str = Box::leak(r.clone().into_boxed_str());
records.push((
nsid::STANDARD_PUBLICATION,
r,
json!({ "name": r, "url": "https://p.example" }),
));
}
for i in 0..100 {
let rkey: &'static str = Box::leak(format!("3l2bus{i:06}").into_boxed_str());
let mut doc = shared_doc("p00", &format!("d{i}"), "2026-07-11T00:00:00Z");
doc["textContent"] = json!(body);
records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
}
let (plc, _) = serve_repo(SHARED, records).await;
let client = crate::feed::build_client().unwrap();
let budget = 4 * 1024 * 1024;
let alone = fetch_repo_capped(&client, &plc, SHARED, &rkeys[..1], 2_000, budget)
.await
.unwrap();
let grouped = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, budget)
.await
.unwrap();
let (a, g) = (alone[0].as_ref().unwrap(), grouped[0].as_ref().unwrap());
assert_eq!((a.entries.len(), a.complete), (100, true));
assert_eq!(
(g.entries.len(), g.complete),
(100, true),
"grouping read it worse than alone"
);
}
#[tokio::test]
async fn one_large_document_does_not_fail_its_publication_in_a_group() {
let mut big = shared_doc("big", "huge", "2026-07-11T00:00:00Z");
big["textContent"] = json!("w".repeat(500 * 1024));
let records = vec![
(
nsid::STANDARD_PUBLICATION,
"big",
json!({ "name": "Big", "url": "https://b.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"other",
json!({ "name": "Other", "url": "https://o.example" }),
),
(nsid::STANDARD_DOCUMENT, "3l2hugeaaaa2a", big),
];
let (plc, _) = serve_repo(SHARED, records).await;
let client = crate::feed::build_client().unwrap();
let rkeys = vec!["big".to_string(), "other".to_string()];
let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 1024 * 1024)
.await
.unwrap();
assert_eq!(
reads[0].as_ref().unwrap().entries.len(),
1,
"a large document was dropped"
);
}
#[tokio::test]
async fn every_requested_rkey_gets_a_result() {
let (plc, _) = serve_repo(
SHARED,
vec![(
nsid::STANDARD_PUBLICATION,
"alpha",
json!({ "name": "Alpha", "url": "https://alpha.example" }),
)],
)
.await;
let client = crate::feed::build_client().unwrap();
let reads = fetch_repo(
&client,
&plc,
SHARED,
&["nope".to_string(), "nope".to_string()],
)
.await
.unwrap();
assert_eq!(reads.len(), 2);
assert!(reads.iter().all(|r| r.is_err()));
}
#[tokio::test]
async fn repeated_and_empty_requests_are_handled() {
let (plc, hits) = serve_repo(
SHARED,
vec![
(
nsid::STANDARD_PUBLICATION,
"alpha",
json!({ "name": "Alpha", "url": "https://alpha.example" }),
),
(
nsid::STANDARD_DOCUMENT,
"3l2rpaaaaaa2a",
shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
),
],
)
.await;
let client = crate::feed::build_client().unwrap();
let none = fetch_repo(&client, &plc, SHARED, &[]).await.unwrap();
assert!(none.is_empty());
assert_eq!(
hits.load(std::sync::atomic::Ordering::SeqCst),
0,
"an empty request reached the network"
);
let twice = fetch_repo(
&client,
&plc,
SHARED,
&["alpha".to_string(), "alpha".to_string()],
)
.await
.unwrap();
for read in &twice {
assert_eq!(
read.as_ref().unwrap().entries.len(),
1,
"a repeated rkey read as empty"
);
}
}
#[tokio::test]
async fn a_group_stores_each_publications_documents_under_its_own_feed() {
let (plc, _) = serve_repo(
SHARED,
vec![
(
nsid::STANDARD_PUBLICATION,
"alpha",
json!({ "name": "Alpha", "url": "https://alpha.example" }),
),
(
nsid::STANDARD_PUBLICATION,
"beta",
json!({ "name": "Beta", "url": "https://beta.example" }),
),
(
nsid::STANDARD_DOCUMENT,
"3l2grpaaaaa2a",
shared_doc("alpha", "only-alpha", "2026-07-11T00:00:00Z"),
),
(
nsid::STANDARD_DOCUMENT,
"3l2grpaaaaa2b",
shared_doc("beta", "only-beta", "2026-07-10T00:00:00Z"),
),
],
)
.await;
let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
let mut feeds = Vec::new();
for rkey in ["alpha", "beta"] {
let url = format!("at://{SHARED}/{}/{rkey}", nsid::STANDARD_PUBLICATION);
crate::store::upsert_feed(
&pool,
&crate::store::NewFeed {
url: url.clone(),
..Default::default()
},
)
.await
.unwrap();
feeds.push(
crate::store::get_feed_by_url(&pool, &url)
.await
.unwrap()
.unwrap(),
);
}
let mut config = crate::config::Config::default();
config.oauth.plc_directory = plc;
let client = crate::feed::build_client().unwrap();
let outcomes = crate::feed::poll_publication_group(&pool, &client, &config, &feeds).await;
assert_eq!(outcomes.len(), 2);
for (feed, want) in feeds.iter().zip(["only-alpha", "only-beta"]) {
let titles: Vec<String> =
sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ?")
.bind(feed.id)
.fetch_all(&pool)
.await
.unwrap();
assert_eq!(
titles,
vec![want.to_string()],
"{} got another publication's documents",
feed.url
);
}
}
#[tokio::test]
async fn a_publication_subscription_delivers_entries_end_to_end() {
const A0: &str = "did:plc:acceptanceaaaaaaaaaaaaaa";
let site = format!("at://{A0}/{}/a0pub", nsid::STANDARD_PUBLICATION);
let document = |title: &str, path: &str| {
json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
"path": path, "site": site, "textContent": "body" })
};
let (plc, _hits) = serve_repo(
A0,
vec![
(
nsid::STANDARD_PUBLICATION,
"a0pub",
json!({ "name": "A0 Journal", "url": "https://a0.example" }),
),
(
nsid::STANDARD_DOCUMENT,
"3l2a0aaaaaa2a",
document("First post", "/first"),
),
(
nsid::STANDARD_DOCUMENT,
"3l2a0aaaaaa2b",
document("Second post", "/second"),
),
],
)
.await;
let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
crate::store::upsert_feed(
&pool,
&crate::store::NewFeed {
url: site.clone(),
..Default::default()
},
)
.await
.unwrap();
let mut config = crate::config::Config::default();
config.oauth.plc_directory = plc;
let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
let feed =
crate::store::due_feeds_of_kind(&pool, &now, crate::feed::FeedKind::Publication, 50)
.await
.unwrap()
.into_iter()
.find(|f| f.url == site)
.expect("the publication is not handed to the publication poller");
let client = crate::feed::build_client().unwrap();
let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
.await
.unwrap();
assert!(
matches!(
outcome,
crate::feed::PollOutcome::Updated { new_entries: 2 }
),
"expected two new entries, got {outcome:?}"
);
let titles: Vec<String> =
sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ? ORDER BY title")
.bind(feed.id)
.fetch_all(&pool)
.await
.unwrap();
assert_eq!(titles, vec!["First post", "Second post"]);
}
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_documents_text_fields_are_bounded_before_they_are_stored() {
let site = canonical("pub");
let big = "x".repeat(8 * 1024 * 1024);
let records = vec![
publication("pub", "https://scanash.com"),
rec(
nsid::STANDARD_DOCUMENT,
"3l2bigaaaaa2a",
json!({ "title": big, "publishedAt": "2026-07-11T00:00:00Z",
"path": format!("/{}", "p".repeat(20_000)), "site": site,
"textContent": "<".repeat(3 * 1024 * 1024) }),
),
];
let (_, publication) = publication_from_records("pub", &records).unwrap();
let entries = entries_from_records(&site, &publication, &records);
let e = &entries[0];
assert!(
e.title.len() <= crate::feed::MAX_TITLE_BYTES,
"title: {}",
e.title.len()
);
let url = e
.url
.as_ref()
.expect("an overlong path is truncated, not dropped");
assert!(
url.len() <= crate::feed::MAX_URL_BYTES,
"url: {}",
url.len()
);
let summary = e.summary.as_ref().unwrap();
assert!(
summary.len() <= crate::feed::MAX_CONTENT_HTML_BYTES,
"the ESCAPED summary is what is stored: {}",
summary.len()
);
}
#[test]
fn a_documents_uri_is_bounded_as_an_entry_id() {
let site = canonical("pub");
let records = vec![
publication("pub", "https://scanash.com"),
rec(
nsid::STANDARD_DOCUMENT,
&"k".repeat(100_000),
json!({ "title": "t", "publishedAt": "2026-07-11T00:00:00Z",
"path": "/p", "site": site }),
),
];
let (_, publication) = publication_from_records("pub", &records).unwrap();
let entries = entries_from_records(&site, &publication, &records);
let stored: crate::store::NewEntry = entries[0].clone().into();
assert!(
stored.guid.len() <= crate::feed::MAX_GUID_BYTES,
"guid: {}",
stored.guid.len()
);
}
#[test]
fn a_publications_own_name_is_bounded() {
let records = vec![rec(
nsid::STANDARD_PUBLICATION,
"pub",
json!({ "name": "n".repeat(100_000), "url": "https://scanash.com" }),
)];
let (_, publication) = publication_from_records("pub", &records).unwrap();
assert!(publication.name.unwrap().len() <= crate::feed::MAX_TITLE_BYTES);
}
#[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 wanted: std::collections::HashMap<String, usize> = [(mine.clone(), 0)].into();
let fate = |d: &crate::atproto::RecordEntry| classify_document(d, &wanted, &known);
assert_eq!(
fate(&document("d1", &mine, "Mine", "/1")),
DocumentFate::Keep(0)
);
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:?}");
}
}
}