1use serde::Deserialize;
30
31use crate::lexicon::nsid;
32
33const DOCUMENT_PAGE_SIZE: u32 = 25;
41
42#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct AtUri {
50 pub authority: String,
51 pub collection: String,
52 pub rkey: String,
53}
54
55impl AtUri {
56 pub fn parse(uri: &str) -> Option<Self> {
58 let rest = uri.strip_prefix(crate::atproto::AT_URI_PREFIX)?;
59 let mut parts = rest.split('/');
60 let (authority, collection, rkey) = (parts.next()?, parts.next()?, parts.next()?);
61 if parts.next().is_some()
62 || authority.is_empty()
63 || collection.is_empty()
64 || rkey.is_empty()
65 {
66 return None;
67 }
68 Some(Self {
69 authority: authority.to_string(),
70 collection: collection.to_string(),
71 rkey: rkey.to_string(),
72 })
73 }
74}
75
76impl std::fmt::Display for AtUri {
77 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78 write!(
79 f,
80 "at://{}/{}/{}",
81 self.authority, self.collection, self.rkey
82 )
83 }
84}
85
86#[derive(Debug, Clone)]
94pub struct PublicationRead {
95 pub publication: Publication,
96 pub entries: Vec<Entry>,
97 pub complete: bool,
98 pub cut_short_by_group: bool,
113}
114
115#[derive(Debug, Clone, PartialEq, Eq)]
118pub struct Publication {
119 pub name: Option<String>,
120 pub url: String,
121}
122
123#[derive(Debug, Clone, PartialEq, Eq)]
125pub struct Entry {
126 pub guid: String,
132 pub title: String,
133 pub published: Option<String>,
142 pub url: Option<String>,
147 pub summary: Option<String>,
151}
152
153impl From<Entry> for crate::store::NewEntry {
154 fn from(e: Entry) -> Self {
157 crate::store::NewEntry {
158 guid: crate::feed::bound_guid(e.guid),
161 url: e.url,
162 title: Some(e.title),
163 author: None,
164 published: e.published,
165 content_html: e.summary,
166 fetched_at: None,
167 keep_stored_content: false,
168 }
169 }
170}
171
172#[derive(Debug, Deserialize)]
173struct PublicationValue {
174 name: Option<String>,
175 url: String,
176}
177
178#[derive(Debug, Deserialize)]
179struct DocumentValue {
180 title: String,
181 #[serde(rename = "publishedAt")]
186 published_at: Option<String>,
187 path: Option<String>,
191 site: String,
197 #[serde(rename = "textContent")]
198 text_content: Option<String>,
199 description: Option<String>,
200}
201
202pub fn publication_from_records(
210 rkey: &str,
211 records: &[crate::atproto::RecordEntry],
212) -> Option<(String, Publication)> {
213 let entry = records
214 .iter()
215 .find(|r| AtUri::parse(&r.uri).is_some_and(|u| u.rkey == rkey))?;
216 let value: PublicationValue = serde_json::from_value(entry.value.clone()).ok()?;
217 let url = crate::net::safe_link(&value.url)?;
222 Some((
223 entry.uri.clone(),
224 Publication {
225 name: value
226 .name
227 .map(|n| crate::feed::bound_text(n, crate::feed::MAX_TITLE_BYTES)),
228 url: crate::feed::bound_text(url, crate::feed::MAX_URL_BYTES),
229 },
230 ))
231}
232
233pub fn entries_from_records(
236 canonical_site: &str,
237 publication: &Publication,
238 records: &[crate::atproto::RecordEntry],
239) -> Vec<Entry> {
240 let base = url::Url::parse(&publication.url).ok().map(|mut u| {
246 if !u.path().ends_with('/') {
247 u.set_path(&format!("{}/", u.path()));
248 }
249 u
250 });
251 let now = chrono::Utc::now();
257 let ceiling = now + chrono::Duration::seconds(crate::atproto::CLOCK_SKEW_GRACE_SECS);
263 records
264 .iter()
265 .filter_map(|record| {
266 let doc: DocumentValue = serde_json::from_value(record.value.clone()).ok()?;
269 if doc.site != canonical_site {
270 return None;
271 }
272 Some(Entry {
273 guid: record.uri.clone(),
274 title: crate::feed::bound_text(doc.title, crate::feed::MAX_TITLE_BYTES),
275 published: doc
291 .published_at
292 .as_deref()
293 .and_then(|raw| chrono::DateTime::parse_from_rfc3339(raw).ok())
294 .map(|d| d.with_timezone(&chrono::Utc))
295 .filter(|d| *d <= ceiling)
296 .or_else(|| {
297 AtUri::parse(&record.uri)
298 .and_then(|uri| crate::atproto::tid_timestamp(&uri.rkey))
299 })
300 .map(crate::feed::fmt_time),
301 url: non_blank(doc.path)
306 .as_deref()
307 .and_then(|path| join_path(base.as_ref(), path))
308 .map(|u| crate::feed::bound_text(u, crate::feed::MAX_URL_BYTES)),
309 summary: non_blank(doc.description)
313 .or_else(|| non_blank(doc.text_content))
314 .map(|raw| {
315 crate::feed::plain_text_to_html_bounded(
316 &raw,
317 crate::feed::MAX_CONTENT_HTML_BYTES,
318 )
319 }),
320 })
321 })
322 .collect()
323}
324
325fn ingest_floor(
373 retention_days: u32,
374 retention_hard_days: u32,
375 now: chrono::DateTime<chrono::Utc>,
376) -> Option<String> {
377 let days = if retention_days > 0 {
378 retention_days
379 } else if retention_hard_days > 0 {
380 retention_hard_days
381 } else {
382 return None;
383 };
384 let window = chrono::Duration::try_days(days.into())?;
385 now.checked_sub_signed(window).map(crate::feed::fmt_time)
386}
387
388pub async fn store_publication(
423 pool: &sqlx::SqlitePool,
424 url: &str,
425 read: PublicationRead,
426 max_entries_per_feed: i64,
427 retention_days: u32,
428 retention_hard_days: u32,
429) -> anyhow::Result<crate::feed::PollOutcome> {
430 let offered = read.entries.len();
431
432 if !read.complete && offered == 0 {
443 return Ok(crate::feed::PollOutcome::Failed {
444 backoff: crate::feed::backoff_for(1),
445 kind: crate::feed::FailureKind::Body,
446 detail: crate::feed::failure_detail(
447 "the publication read stopped before its first document",
448 ),
449 });
450 }
451
452 let floor = ingest_floor(retention_days, retention_hard_days, chrono::Utc::now());
454 let rows: Vec<crate::store::NewEntry> = read
455 .entries
456 .into_iter()
457 .filter(|e| match (&floor, &e.published) {
458 (Some(floor), Some(published)) => published.as_str() >= floor.as_str(),
461 _ => true,
465 })
466 .map(Into::into)
467 .collect();
468
469 if offered > rows.len() {
475 tracing::info!(
476 feed = %url,
477 offered,
478 stored = rows.len(),
479 "the retention floor dropped documents older than the window"
480 );
481 }
482
483 let feed_id = crate::store::upsert_feed(
484 pool,
485 &crate::store::NewFeed {
486 url: url.to_string(),
487 title: read.publication.name.clone(),
488 site_url: Some(read.publication.url.clone()),
489 last_polled: Some(crate::feed::fmt_time(chrono::Utc::now())),
490 ..Default::default()
491 },
492 )
493 .await?;
494
495 let new_entries =
496 crate::store::insert_entries(pool, feed_id, &rows, max_entries_per_feed).await?;
497 Ok(crate::feed::PollOutcome::Updated { new_entries })
498}
499
500fn non_blank(s: Option<String>) -> Option<String> {
501 s.filter(|v| !v.trim().is_empty())
502}
503
504fn join_path(base: Option<&url::Url>, path: &str) -> Option<String> {
512 let base = base?;
526 match base.join(path) {
527 Ok(joined) if joined.origin() == base.origin() => crate::net::safe_link(joined.as_str()),
530 _ => None,
536 }
537}
538
539#[derive(Debug, Clone, Copy, PartialEq, Eq)]
544enum DocumentFate {
545 Keep(usize),
547 Sibling,
550 Orphan,
554 Malformed,
556}
557
558fn classify_document(
559 record: &crate::atproto::RecordEntry,
560 wanted: &std::collections::HashMap<String, usize>,
561 known: &std::collections::HashSet<&str>,
562) -> DocumentFate {
563 match serde_json::from_value::<DocumentValue>(record.value.clone()) {
564 Ok(doc) if wanted.contains_key(&doc.site) => DocumentFate::Keep(wanted[&doc.site]),
565 Ok(doc) if known.contains(doc.site.as_str()) => DocumentFate::Sibling,
566 Ok(_) => DocumentFate::Orphan,
567 Err(_) => DocumentFate::Malformed,
568 }
569}
570
571pub async fn fetch(
592 http: &reqwest::Client,
593 plc_directory: &str,
594 uri: &AtUri,
595) -> anyhow::Result<PublicationRead> {
596 if uri.collection != nsid::STANDARD_PUBLICATION {
601 return Err(
602 NotAPublication(format!("{uri} is not a {} URI", nsid::STANDARD_PUBLICATION)).into(),
603 );
604 }
605 fetch_repo(
606 http,
607 plc_directory,
608 &uri.authority,
609 std::slice::from_ref(&uri.rkey),
610 )
611 .await?
612 .pop()
613 .unwrap_or_else(|| Err(NotAPublication(format!("{uri} was not read")).into()))
614}
615
616#[derive(Debug)]
624pub struct NotAPublication(pub String);
625
626impl std::fmt::Display for NotAPublication {
627 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
628 f.write_str(&self.0)
629 }
630}
631
632impl std::error::Error for NotAPublication {}
633
634pub async fn fetch_repo(
647 http: &reqwest::Client,
648 plc_directory: &str,
649 did: &str,
650 rkeys: &[String],
651) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
652 fetch_repo_capped(
653 http,
654 plc_directory,
655 did,
656 rkeys,
657 crate::atproto::MAX_LARGE_RECORDS,
658 crate::atproto::MAX_LIST_BYTES,
659 )
660 .await
661}
662
663pub(crate) async fn fetch_repo_capped(
668 http: &reqwest::Client,
669 plc_directory: &str,
670 did: &str,
671 requested: &[String],
672 per_site_cap: usize,
673 budget_bytes: usize,
674) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
675 use anyhow::Context;
676
677 if requested.is_empty() {
680 return Ok(Vec::new());
681 }
682 let mut rkeys: Vec<String> = Vec::new();
685 for r in requested {
686 if !rkeys.contains(r) {
687 rkeys.push(r.clone());
688 }
689 }
690 let rkeys = rkeys.as_slice();
691
692 let pds = crate::atproto::resolve_did_to_pds(http, plc_directory, did)
693 .await
694 .with_context(|| format!("resolving the PDS for {did}"))?;
695 let client = crate::atproto::PdsClient::anonymous(http.clone(), pds, did.to_string());
696
697 let mut budget = crate::atproto::ByteBudget::new(budget_bytes);
701 let (publications, skipped_publications) = client
702 .list_all_records_skipping_within(nsid::STANDARD_PUBLICATION, &mut budget)
703 .await
704 .with_context(|| format!("listing publications for {did}"))?;
705 if skipped_publications > 0 {
706 tracing::warn!(
707 repo = %did,
708 skipped = skipped_publications,
709 "skipped malformed publication records in this repo"
710 );
711 }
712
713 let wanted: Vec<Option<(String, Publication)>> = rkeys
715 .iter()
716 .map(|rkey| publication_from_records(rkey, &publications))
717 .collect();
718 let index_of: std::collections::HashMap<String, usize> = wanted
719 .iter()
720 .enumerate()
721 .filter_map(|(i, w)| w.as_ref().map(|(site, _)| (site.clone(), i)))
722 .collect();
723 let not_found = |rkey: &str| {
724 anyhow::Error::new(NotAPublication(format!(
725 "at://{did}/{}/{rkey} is not a readable site.standard.publication",
726 nsid::STANDARD_PUBLICATION
727 )))
728 };
729 if index_of.is_empty() {
730 return Ok(requested.iter().map(|r| Err(not_found(r))).collect());
731 }
732
733 let known: std::collections::HashSet<&str> =
739 publications.iter().map(|p| p.uri.as_str()).collect();
740 let mut kept_per = vec![0usize; rkeys.len()];
741 let mut capped = vec![false; rkeys.len()];
742 let mut orphaned = 0usize;
743 let documents = client
744 .list_recent_matching_within(
745 nsid::STANDARD_DOCUMENT,
746 per_site_cap.saturating_mul(index_of.len()),
747 &mut budget,
748 DOCUMENT_PAGE_SIZE,
749 |record| match classify_document(record, &index_of, &known) {
750 DocumentFate::Keep(i) if kept_per[i] < per_site_cap => {
751 kept_per[i] += 1;
752 true
753 }
754 DocumentFate::Keep(i) => {
755 capped[i] = true;
756 false
757 }
758 DocumentFate::Orphan => {
762 orphaned += 1;
763 false
764 }
765 DocumentFate::Sibling | DocumentFate::Malformed => false,
766 },
767 )
768 .await
769 .with_context(|| format!("listing documents for {did}"))?;
770
771 let grouped = index_of.len() > 1;
773 let mut per_site: Vec<Vec<crate::atproto::RecordEntry>> = vec![Vec::new(); rkeys.len()];
774 for record in documents.records {
775 let site = serde_json::from_value::<DocumentValue>(record.value.clone()).map(|d| d.site);
776 if let Some(&i) = site.ok().as_deref().and_then(|s| index_of.get(s)) {
777 per_site[i].push(record);
778 }
779 }
780 let reads: Vec<anyhow::Result<PublicationRead>> = wanted
781 .into_iter()
782 .zip(per_site)
783 .enumerate()
784 .map(|(i, (w, records))| match w {
785 None => Err(not_found(&rkeys[i])),
786 Some((site, publication)) => {
787 let walk = crate::atproto::RecordWalk {
788 records,
789 complete: documents.complete && !capped[i],
790 malformed: documents.malformed,
791 out_of_budget: documents.out_of_budget,
792 };
793 let mut read = read_from(publication, &site, walk, orphaned);
794 read.cut_short_by_group = grouped && documents.out_of_budget && !capped[i];
801 Ok(read)
802 }
803 })
804 .collect();
805 let mut reads: Vec<Option<anyhow::Result<PublicationRead>>> =
808 reads.into_iter().map(Some).collect();
809 let positions: Vec<usize> = requested
810 .iter()
811 .map(|r| {
812 rkeys
813 .iter()
814 .position(|k| k == r)
815 .expect("every requested rkey is in rkeys")
816 })
817 .collect();
818 Ok(positions
819 .iter()
820 .enumerate()
821 .map(|(n, &u)| {
822 let repeated_later = positions[n + 1..].contains(&u);
823 match (&reads[u], repeated_later) {
824 (Some(Ok(read)), true) => Ok(read.clone()),
825 (Some(Err(_)), _) | (None, _) => Err(not_found(&requested[n])),
826 (Some(Ok(_)), false) => match reads[u].take() {
827 Some(Ok(read)) => Ok(read),
828 _ => Err(not_found(&requested[n])),
829 },
830 }
831 })
832 .collect())
833}
834
835fn read_from(
842 publication: Publication,
843 canonical_site: &str,
844 documents: crate::atproto::RecordWalk,
845 orphaned: usize,
846) -> PublicationRead {
847 let entries = entries_from_records(canonical_site, &publication, &documents.records);
848
849 if !documents.complete {
859 tracing::warn!(
860 site = %canonical_site,
861 kept = entries.len(),
862 "stopped reading this publication before its documents ran out"
863 );
864 }
865 if documents.malformed > 0 {
874 tracing::warn!(
875 site = %canonical_site,
876 skipped = documents.malformed,
877 "skipped malformed document records for this publication"
878 );
879 }
880 if orphaned > 0 {
881 tracing::warn!(
882 site = %canonical_site,
883 orphaned,
884 "documents in this repo reference no publication in it — a `site` spelling nothing matches"
885 );
886 }
887 PublicationRead {
888 publication,
889 entries,
890 complete: documents.complete,
891 cut_short_by_group: false,
893 }
894}
895
896#[cfg(test)]
897pub(crate) mod tests {
898 use super::*;
899 use crate::atproto::RecordEntry;
900 use serde_json::json;
901
902 const DID: &str = "did:plc:ohutz6x5acjmpuulp3x7wxxc";
903
904 #[test]
910 fn a_read_carries_whether_the_walk_finished() {
911 let publication = Publication {
912 name: Some("Scan's Lab".to_string()),
913 url: "https://example.com/blog/".to_string(),
914 };
915 for complete in [true, false] {
916 let walk = crate::atproto::RecordWalk {
917 records: Vec::new(),
918 complete,
919 malformed: 0,
920 out_of_budget: false,
921 };
922 let read = read_from(publication.clone(), "at://d/c/r", walk, 0);
923 assert_eq!(
924 read.complete, complete,
925 "the walk said complete={complete} and the read said {}",
926 read.complete
927 );
928 }
929 }
930
931 fn read_of(entries: Vec<Entry>, complete: bool) -> PublicationRead {
934 PublicationRead {
935 publication: Publication {
936 name: Some("Scan's Lab".to_string()),
937 url: "https://example.com/blog/".to_string(),
938 },
939 entries,
940 complete,
941 cut_short_by_group: false,
942 }
943 }
944
945 fn entry_dated(guid: &str, published: Option<&str>) -> Entry {
946 Entry {
947 guid: guid.to_string(),
948 title: "T".to_string(),
949 published: published.map(str::to_string),
950 url: Some("https://example.com/blog/a".to_string()),
951 summary: None,
952 }
953 }
954
955 fn days_ago(n: i64) -> String {
956 crate::feed::fmt_time(chrono::Utc::now() - chrono::Duration::days(n))
957 }
958
959 const PUB_URL: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
960
961 async fn pool() -> sqlx::SqlitePool {
962 crate::store::init_url("sqlite::memory:").await.unwrap()
963 }
964
965 #[tokio::test]
966 async fn a_complete_read_stores_its_entries_and_reports_them() {
967 let pool = pool().await;
968 let read = read_of(
969 vec![
970 entry_dated("at://d/c/1", Some(&days_ago(1))),
971 entry_dated("at://d/c/2", Some(&days_ago(2))),
972 ],
973 true,
974 );
975 let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
976 .await
977 .unwrap();
978 assert!(
979 matches!(
980 outcome,
981 crate::feed::PollOutcome::Updated { new_entries: 2 }
982 ),
983 "expected two new entries, got {outcome:?}"
984 );
985 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
986 .fetch_one(&pool)
987 .await
988 .unwrap();
989 assert_eq!(n, 2, "the entries were not stored");
990 }
991
992 #[tokio::test]
1000 async fn an_incomplete_read_that_offered_nothing_is_a_failure() {
1001 let pool = pool().await;
1002 let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
1003 .await
1004 .unwrap();
1005 assert!(
1006 matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
1007 "a truncated read that produced nothing is not a healthy poll: {outcome:?}"
1008 );
1009 }
1010
1011 #[tokio::test]
1012 async fn an_incomplete_read_whose_entries_the_floor_dropped_is_not_a_failure() {
1013 let pool = pool().await;
1014 let read = read_of(
1015 vec![entry_dated("at://d/c/old", Some(&days_ago(900)))],
1016 false,
1017 );
1018 let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
1019 .await
1020 .unwrap();
1021 assert!(
1022 !matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
1023 "the read offered an entry; the floor dropping it is not a failed poll: {outcome:?}"
1024 );
1025 }
1026
1027 #[tokio::test]
1029 async fn a_failed_read_does_not_stamp_last_polled() {
1030 let pool = pool().await;
1031 let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
1032 .await
1033 .unwrap();
1034 assert!(matches!(outcome, crate::feed::PollOutcome::Failed { .. }));
1035 let stamped: Option<String> =
1036 sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
1037 .bind(PUB_URL)
1038 .fetch_optional(&pool)
1039 .await
1040 .unwrap()
1041 .flatten();
1042 assert_eq!(
1043 stamped, None,
1044 "a failed poll stamped last_polled, so the feed reads as freshly polled"
1045 );
1046 }
1047
1048 #[tokio::test]
1049 async fn an_entry_already_older_than_the_window_is_not_stored() {
1050 let pool = pool().await;
1051 let read = read_of(
1052 vec![
1053 entry_dated("at://d/c/fresh", Some(&days_ago(1))),
1054 entry_dated("at://d/c/ancient", Some(&days_ago(900))),
1055 ],
1056 true,
1057 );
1058 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1059 .await
1060 .unwrap();
1061 let guids: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
1062 .fetch_all(&pool)
1063 .await
1064 .unwrap();
1065 assert_eq!(
1066 guids,
1067 vec!["at://d/c/fresh".to_string()],
1068 "an entry the next sweep would delete was stored anyway"
1069 );
1070 }
1071
1072 #[tokio::test]
1080 async fn the_floor_follows_the_hard_ceiling_when_the_window_is_disabled() {
1081 let pool = pool().await;
1082 let read = read_of(
1083 vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
1084 true,
1085 );
1086 store_publication(&pool, PUB_URL, read, 0, 0, 180)
1087 .await
1088 .unwrap();
1089 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1090 .fetch_one(&pool)
1091 .await
1092 .unwrap();
1093 assert_eq!(
1094 n, 0,
1095 "an entry the hard ceiling will delete was stored, so it will resurrect"
1096 );
1097 }
1098
1099 #[tokio::test]
1106 async fn the_floor_follows_the_window_when_both_are_set() {
1107 let pool = pool().await;
1108 let read = read_of(
1109 vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
1110 true,
1111 );
1112 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1113 .await
1114 .unwrap();
1115 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1116 .fetch_one(&pool)
1117 .await
1118 .unwrap();
1119 assert_eq!(
1120 n, 0,
1121 "an entry inside the ceiling but outside the window was stored, so it will cycle"
1122 );
1123 }
1124
1125 #[test]
1139 fn a_publications_floor_is_its_archive_ceiling_not_the_rss_window() {
1140 let config = crate::config::Config::default();
1141 let now = chrono::DateTime::parse_from_rfc3339("2026-09-29T12:00:00Z")
1142 .unwrap()
1143 .with_timezone(&chrono::Utc);
1144
1145 let (days, hard) = config.retention_for(crate::feed::FeedKind::Publication);
1146 let floor = ingest_floor(days, hard, now).expect("a publication has an ingest floor");
1147 let expected = crate::feed::fmt_time(
1148 now - chrono::Duration::days(config.publication_retention_days.into()),
1149 );
1150 assert_eq!(
1151 floor, expected,
1152 "a publication's floor must be its archive ceiling ({} days), because \
1153 that is the only sweep pass that can delete its rows",
1154 config.publication_retention_days,
1155 );
1156
1157 let rss_floor = ingest_floor(config.retention_days, config.retention_hard_days, now)
1160 .expect("an RSS feed has an ingest floor");
1161 assert!(
1162 floor < rss_floor,
1163 "the publication floor ({floor}) is no older than the RSS one \
1164 ({rss_floor}), so a months-old document would still be dropped",
1165 );
1166 let a_real_publications_newest_document =
1167 crate::feed::fmt_time(now - chrono::Duration::days(109));
1168 assert!(
1169 a_real_publications_newest_document.as_str() >= floor.as_str(),
1170 "the newest document a real publication offered would be refused at \
1171 ingest: {a_real_publications_newest_document} against a floor of {floor}",
1172 );
1173 assert!(
1174 a_real_publications_newest_document.as_str() < rss_floor.as_str(),
1175 "this assertion is only meaningful while the RSS window WOULD have \
1176 dropped it, and it no longer does",
1177 );
1178 }
1179
1180 #[test]
1187 fn the_ingest_floor_is_the_window_the_sweep_would_use() {
1188 let now = chrono::DateTime::parse_from_rfc3339("2026-09-26T12:00:00Z")
1189 .unwrap()
1190 .with_timezone(&chrono::Utc);
1191
1192 assert_eq!(
1193 ingest_floor(14, 180, now).as_deref(),
1194 Some("2026-09-12T12:00:00Z"),
1195 "with both set, the floor is the WINDOW — the thing that deletes first",
1196 );
1197 assert_eq!(
1198 ingest_floor(180, 30, now).as_deref(),
1199 Some("2026-03-30T12:00:00Z"),
1200 "a ceiling INSIDE the window is one `prune_old_entries` ignores, so it \
1201 must not lower the floor — `min` here discarded five months of archive \
1202 that nothing would have deleted",
1203 );
1204 assert_eq!(
1205 ingest_floor(0, 30, now).as_deref(),
1206 Some("2026-08-27T12:00:00Z"),
1207 "with no window the ceiling stands alone, and it still deletes",
1208 );
1209 assert_eq!(
1210 ingest_floor(0, 0, now),
1211 None,
1212 "with no retention at all there is nothing to floor against",
1213 );
1214 assert_eq!(
1219 ingest_floor(u32::MAX, u32::MAX, now),
1220 None,
1221 "an unrepresentable window produced a floor, so the comparison is \
1222 resting on how `fmt_time` renders an out-of-range year",
1223 );
1224 }
1225
1226 #[tokio::test]
1241 async fn a_ceiling_inside_the_window_does_not_lower_the_floor() {
1242 let pool = pool().await;
1243 let read = read_of(
1244 vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
1245 true,
1246 );
1247 store_publication(&pool, PUB_URL, read, 0, 180, 30)
1248 .await
1249 .unwrap();
1250 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1251 .fetch_one(&pool)
1252 .await
1253 .unwrap();
1254 assert_eq!(
1255 n, 1,
1256 "an entry inside the 180-day window was dropped because of a 30-day \
1257 ceiling the sweep ignores — five months of archive discarded at ingest \
1258 that nothing would have deleted",
1259 );
1260 }
1261
1262 #[tokio::test]
1267 async fn a_complete_read_of_nothing_is_not_a_failure() {
1268 let pool = pool().await;
1269 let outcome = store_publication(&pool, PUB_URL, read_of(vec![], true), 0, 14, 180)
1270 .await
1271 .unwrap();
1272 assert!(
1273 matches!(
1274 outcome,
1275 crate::feed::PollOutcome::Updated { new_entries: 0 }
1276 ),
1277 "a complete read of an empty publication was not a healthy poll: {outcome:?}"
1278 );
1279 }
1280
1281 #[tokio::test]
1285 async fn a_successful_read_stamps_last_polled() {
1286 let pool = pool().await;
1287 let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
1288 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1289 .await
1290 .unwrap();
1291 let stamped: Option<String> =
1292 sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
1293 .bind(PUB_URL)
1294 .fetch_one(&pool)
1295 .await
1296 .unwrap();
1297 assert!(
1298 stamped.is_some(),
1299 "a successful read left `last_polled` NULL, so the publication stays \
1300 due forever and `/stats` never shows it as polled",
1301 );
1302 }
1303
1304 #[tokio::test]
1311 async fn an_absurd_retention_window_keeps_everything_rather_than_nothing() {
1312 let pool = pool().await;
1313 let read = read_of(
1314 vec![entry_dated("at://d/c/ancient", Some(&days_ago(10_000)))],
1315 true,
1316 );
1317 store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
1318 .await
1319 .expect("a huge window is a wide floor, not a crash");
1320 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1321 .fetch_one(&pool)
1322 .await
1323 .unwrap();
1324 assert_eq!(
1325 n, 1,
1326 "a window of u32::MAX days dropped a 27-year-old entry, so the \
1327 saturation went the wrong way",
1328 );
1329 }
1330
1331 #[tokio::test]
1336 async fn the_per_feed_cap_is_the_one_the_caller_passed() {
1337 let pool = pool().await;
1338 let read = read_of(
1339 (0..5)
1340 .map(|i| entry_dated(&format!("at://d/c/{i}"), Some(&days_ago(i + 1))))
1341 .collect(),
1342 true,
1343 );
1344 store_publication(&pool, PUB_URL, read, 2, 0, 0)
1345 .await
1346 .unwrap();
1347 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1348 .fetch_one(&pool)
1349 .await
1350 .unwrap();
1351 assert_eq!(
1352 n, 2,
1353 "five entries under a cap of two left {n} rows, so the caller's cap is \
1354 not the one being applied",
1355 );
1356 }
1357
1358 #[tokio::test]
1359 async fn no_retention_at_all_means_no_ingest_floor() {
1360 let pool = pool().await;
1361 let read = read_of(
1362 vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
1363 true,
1364 );
1365 store_publication(&pool, PUB_URL, read, 0, 0, 0)
1366 .await
1367 .unwrap();
1368 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1369 .fetch_one(&pool)
1370 .await
1371 .unwrap();
1372 assert_eq!(
1373 n, 1,
1374 "with nothing deleting it, an old entry is worth keeping"
1375 );
1376 }
1377
1378 #[tokio::test]
1380 async fn an_absurd_retention_window_does_not_panic() {
1381 let pool = pool().await;
1382 let read = read_of(vec![entry_dated("at://d/c/x", Some(&days_ago(1)))], true);
1383 store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
1384 .await
1385 .expect("a huge window is a wide floor, not a crash");
1386 }
1387
1388 #[tokio::test]
1392 async fn the_feed_row_learns_the_publications_name_and_site() {
1393 let pool = pool().await;
1394 let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
1395 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1396 .await
1397 .unwrap();
1398 let (title, site): (Option<String>, Option<String>) =
1399 sqlx::query_as("SELECT title, site_url FROM feeds WHERE url = ?1")
1400 .bind(PUB_URL)
1401 .fetch_one(&pool)
1402 .await
1403 .unwrap();
1404 assert_eq!(
1405 title.as_deref(),
1406 Some("Scan's Lab"),
1407 "the name never reached the row"
1408 );
1409 assert_eq!(
1410 site.as_deref(),
1411 Some("https://example.com/blog/"),
1412 "the homepage never reached the row"
1413 );
1414 }
1415
1416 #[tokio::test]
1420 async fn an_undated_entry_is_stored_rather_than_dropped() {
1421 let pool = pool().await;
1422 let read = read_of(vec![entry_dated("at://d/c/undated", None)], true);
1423 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1424 .await
1425 .unwrap();
1426 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1427 .fetch_one(&pool)
1428 .await
1429 .unwrap();
1430 assert_eq!(n, 1, "an entry with no date was discarded");
1431 }
1432
1433 const PAST_TID: &str = "3jzfcijpj2z2a";
1441 const PAST_TID_WRITTEN_AT: &str = "2023-06-30T15:03:01Z";
1442
1443 fn rec(collection: &str, rkey: &str, value: serde_json::Value) -> RecordEntry {
1444 RecordEntry {
1445 uri: format!("at://{DID}/{collection}/{rkey}"),
1446 cid: None,
1447 value,
1448 }
1449 }
1450
1451 fn publication(rkey: &str, url: &str) -> RecordEntry {
1452 rec(
1453 nsid::STANDARD_PUBLICATION,
1454 rkey,
1455 json!({ "name": "Scan's Lab", "url": url }),
1456 )
1457 }
1458
1459 fn document(rkey: &str, site: &str, title: &str, path: &str) -> RecordEntry {
1460 rec(
1461 nsid::STANDARD_DOCUMENT,
1462 rkey,
1463 json!({
1464 "title": title,
1465 "publishedAt": "2026-07-11T00:00:00Z",
1466 "path": path,
1467 "site": site,
1468 "textContent": "body",
1469 }),
1470 )
1471 }
1472
1473 pub(crate) async fn serve_repo(
1480 did: &'static str,
1481 records: Vec<(&'static str, &'static str, serde_json::Value)>,
1482 ) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
1483 serve_repo_slow_after(did, records, usize::MAX, std::time::Duration::ZERO).await
1484 }
1485
1486 pub(crate) const FAIL_AFTER: std::time::Duration = std::time::Duration::MAX;
1491
1492 pub(crate) async fn serve_repo_slow_after(
1493 did: &'static str,
1494 records: Vec<(&'static str, &'static str, serde_json::Value)>,
1495 after: usize,
1496 delay: std::time::Duration,
1497 ) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
1498 use axum::extract::{Query, Request};
1499 use std::collections::HashMap;
1500 use std::sync::atomic::{AtomicUsize, Ordering};
1501 use std::sync::Arc;
1502
1503 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1504 let addr = listener.local_addr().unwrap();
1505 let port = addr.port();
1506 let (plc_host, pds_host) = (
1509 format!("plc-{port}.repo.test"),
1510 format!("pds-{port}.repo.test"),
1511 );
1512 for host in [&plc_host, &pds_host] {
1513 crate::net::test_host_override(host, addr);
1514 }
1515 let hits = Arc::new(AtomicUsize::new(0));
1516 let counter = Arc::clone(&hits);
1517 let records = Arc::new(records);
1518 let app = axum::Router::new().fallback(
1519 move |Query(q): Query<HashMap<String, String>>, req: Request| {
1520 let records = Arc::clone(&records);
1521 let counter = Arc::clone(&counter);
1522 async move {
1523 let path = req.uri().path().to_string();
1524 if path == format!("/{did}") {
1525 return axum::Json(json!({
1526 "id": did,
1527 "service": [{
1528 "id": "#atproto_pds",
1529 "type": "AtprotoPersonalDataServer",
1530 "serviceEndpoint": format!("http://{pds_host}:{port}"),
1531 }],
1532 }));
1533 }
1534 assert_eq!(path, "/xrpc/com.atproto.repo.listRecords", "unexpected request");
1535 if counter.fetch_add(1, Ordering::SeqCst) >= after {
1536 if delay == FAIL_AFTER {
1537 return axum::Json(json!({ "error": "RepoDeactivated" }));
1538 }
1539 tokio::time::sleep(delay).await;
1540 }
1541 let collection = q.get("collection").cloned().unwrap_or_default();
1542 let limit: usize = q.get("limit").and_then(|l| l.parse().ok()).unwrap_or(50);
1545 let offset: usize = q.get("cursor").and_then(|c| c.parse().ok()).unwrap_or(0);
1546 let page: Vec<_> = records
1547 .iter()
1548 .filter(|(c, _, _)| *c == collection)
1549 .skip(offset)
1550 .take(limit)
1551 .map(|(c, rkey, value)| {
1552 if rkey.is_empty() {
1555 return json!({ "cid": "bafy", "value": value });
1556 }
1557 json!({ "uri": format!("at://{did}/{c}/{rkey}"), "cid": "bafy", "value": value })
1558 })
1559 .collect();
1560 let next = (offset + page.len()).to_string();
1561 axum::Json(json!({ "records": page, "cursor": next }))
1562 }
1563 },
1564 );
1565 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1566 (format!("http://{plc_host}:{port}"), hits)
1567 }
1568
1569 #[tokio::test]
1574 async fn fetch_reads_a_publication_from_a_mocked_repo() {
1575 const OWN: &str = "did:plc:fetchmock";
1576 let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
1577 let sibling = format!("at://{OWN}/{}/other", nsid::STANDARD_PUBLICATION);
1578 let doc = |site: &str, title: &str| {
1579 json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
1580 "path": "/p", "site": site })
1581 };
1582 let (plc, hits) = serve_repo(
1583 OWN,
1584 vec![
1585 (
1586 nsid::STANDARD_PUBLICATION,
1587 "mine",
1588 json!({ "name": "Mine", "url": "https://mine.example" }),
1589 ),
1590 (
1591 nsid::STANDARD_PUBLICATION,
1592 "other",
1593 json!({ "name": "Other", "url": "https://other.example" }),
1594 ),
1595 (nsid::STANDARD_DOCUMENT, "3l2fmaaaaaa2a", doc(&site, "kept")),
1596 (
1597 nsid::STANDARD_DOCUMENT,
1598 "3l2fmaaaaaa2b",
1599 doc(&sibling, "sibling's"),
1600 ),
1601 ],
1602 )
1603 .await;
1604 let client = crate::feed::build_client().unwrap();
1605 let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
1606 .await
1607 .unwrap();
1608 assert!(read.complete, "the walk did not finish");
1609 assert_eq!(read.publication.name.as_deref(), Some("Mine"));
1610 let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
1611 assert_eq!(
1612 titles,
1613 vec!["kept"],
1614 "the site filter let a sibling through"
1615 );
1616 assert_eq!(
1617 hits.load(std::sync::atomic::Ordering::SeqCst),
1618 4,
1619 "two collections, each a full page then an empty one"
1620 );
1621 }
1622
1623 #[tokio::test]
1627 async fn a_failed_publication_read_is_a_poll_failure() {
1628 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1629 let url = format!(
1632 "at://did:plc:unreachableaaaaaaaaaaaaa/{}/x",
1633 nsid::STANDARD_PUBLICATION
1634 );
1635 crate::store::upsert_feed(
1636 &pool,
1637 &crate::store::NewFeed {
1638 url: url.clone(),
1639 ..Default::default()
1640 },
1641 )
1642 .await
1643 .unwrap();
1644 let feed = crate::store::get_feed_by_url(&pool, &url)
1645 .await
1646 .unwrap()
1647 .unwrap();
1648 let mut config = crate::config::Config::default();
1649 config.oauth.plc_directory = "http://plc.nowhere.invalid".into();
1650 let client = crate::feed::build_client().unwrap();
1651 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1652 .await
1653 .expect("a source failure surfaced as a store error");
1654 assert!(
1655 matches!(
1656 outcome,
1657 crate::feed::PollOutcome::Failed {
1658 kind: crate::feed::FailureKind::Fetch,
1659 ..
1660 }
1661 ),
1662 "an unreachable publication was not a fetch failure: {outcome:?}"
1663 );
1664 }
1665
1666 #[tokio::test]
1670 async fn a_deleted_publication_record_is_not_an_unreachable_publisher() {
1671 const GONE: &str = "did:plc:goneaaaaaaaaaaaaaaaaaaaa";
1672 let site = format!("at://{GONE}/{}/deleted", nsid::STANDARD_PUBLICATION);
1673 let (plc, _) = serve_repo(
1675 GONE,
1676 vec![(
1677 nsid::STANDARD_PUBLICATION,
1678 "another",
1679 json!({ "name": "Other", "url": "https://o.example" }),
1680 )],
1681 )
1682 .await;
1683 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1684 crate::store::upsert_feed(
1685 &pool,
1686 &crate::store::NewFeed {
1687 url: site.clone(),
1688 ..Default::default()
1689 },
1690 )
1691 .await
1692 .unwrap();
1693 let feed = crate::store::get_feed_by_url(&pool, &site)
1694 .await
1695 .unwrap()
1696 .unwrap();
1697 let mut config = crate::config::Config::default();
1698 config.oauth.plc_directory = plc;
1699 let client = crate::feed::build_client().unwrap();
1700 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1701 .await
1702 .unwrap();
1703 assert!(
1704 matches!(
1705 outcome,
1706 crate::feed::PollOutcome::Failed {
1707 kind: crate::feed::FailureKind::Parse,
1708 ..
1709 }
1710 ),
1711 "a deleted publication was filed as a network failure: {outcome:?}"
1712 );
1713 }
1714
1715 async fn serve_answering(
1720 did: &'static str,
1721 plc_status: u16,
1722 pds: (u16, &'static str),
1723 delay: std::time::Duration,
1724 endless: bool,
1725 ) -> String {
1726 use std::sync::atomic::{AtomicUsize, Ordering};
1727 use std::sync::Arc;
1728 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1729 let addr = listener.local_addr().unwrap();
1730 let port = addr.port();
1731 let (plc_host, pds_host) = (
1732 format!("plc-{port}.answer.test"),
1733 format!("pds-{port}.answer.test"),
1734 );
1735 crate::net::test_host_override(&plc_host, addr);
1736 crate::net::test_host_override(&pds_host, addr);
1737 let pages = Arc::new(AtomicUsize::new(0));
1738 let endpoint = format!("http://{pds_host}:{port}");
1739 let app = axum::Router::new().fallback(move |req: axum::extract::Request| {
1740 let pages = Arc::clone(&pages);
1741 let endpoint = endpoint.clone();
1742 async move {
1743 use axum::response::IntoResponse;
1744 if req.uri().path() == format!("/{did}") {
1745 let status = axum::http::StatusCode::from_u16(plc_status).unwrap();
1746 let doc = json!({ "id": did, "service": [{ "id": "#atproto_pds",
1747 "type": "AtprotoPersonalDataServer", "serviceEndpoint": endpoint }] });
1748 return (status, axum::Json(doc)).into_response();
1749 }
1750 tokio::time::sleep(delay).await;
1751 if endless {
1752 let n = pages.fetch_add(1, Ordering::SeqCst);
1753 let body = json!({ "records": [{
1754 "uri": format!("at://{did}/{}/3lend{n:08}", nsid::STANDARD_DOCUMENT),
1755 "cid": "b",
1756 "value": { "title": "x", "path": "/x", "publishedAt": "2026-07-11T00:00:00Z",
1757 "site": format!("at://{did}/{}/other", nsid::STANDARD_PUBLICATION) } }],
1758 "cursor": format!("c{n}") });
1759 return axum::Json(body).into_response();
1760 }
1761 let status = axum::http::StatusCode::from_u16(pds.0).unwrap();
1762 (status, [("content-type", "application/json")], pds.1).into_response()
1763 }
1764 });
1765 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1766 format!("http://{plc_host}:{port}")
1767 }
1768
1769 async fn poll_publication_at(
1770 did: &str,
1771 plc: String,
1772 deadline: Option<std::time::Duration>,
1773 ) -> crate::feed::PollOutcome {
1774 let site = format!("at://{did}/{}/mine", nsid::STANDARD_PUBLICATION);
1775 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1776 crate::store::upsert_feed(
1777 &pool,
1778 &crate::store::NewFeed {
1779 url: site.clone(),
1780 ..Default::default()
1781 },
1782 )
1783 .await
1784 .unwrap();
1785 let feed = crate::store::get_feed_by_url(&pool, &site)
1786 .await
1787 .unwrap()
1788 .unwrap();
1789 let mut config = crate::config::Config::default();
1790 config.oauth.plc_directory = plc;
1791 if let Some(d) = deadline {
1792 config.publication_read_deadline = d;
1793 }
1794 let client = crate::feed::build_client().unwrap();
1795 crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1796 .await
1797 .unwrap()
1798 }
1799
1800 fn kind_of(outcome: &crate::feed::PollOutcome) -> Option<crate::feed::FailureKind> {
1801 match outcome {
1802 crate::feed::PollOutcome::Failed { kind, .. } => Some(*kind),
1803 _ => None,
1804 }
1805 }
1806
1807 #[tokio::test]
1810 async fn a_publication_read_has_an_overall_deadline() {
1811 const DID: &str = "did:plc:slowrepoaaaaaaaaaaaaaaaa";
1812 let plc = serve_answering(
1813 DID,
1814 200,
1815 (200, ""),
1816 std::time::Duration::from_millis(50),
1817 true,
1818 )
1819 .await;
1820 let started = std::time::Instant::now();
1821 let outcome =
1822 poll_publication_at(DID, plc, Some(std::time::Duration::from_millis(300))).await;
1823 assert!(
1824 started.elapsed() < std::time::Duration::from_secs(3),
1825 "the read ran {:?}",
1826 started.elapsed()
1827 );
1828 assert_eq!(
1829 kind_of(&outcome),
1830 Some(crate::feed::FailureKind::Fetch),
1831 "{outcome:?}"
1832 );
1833 }
1834
1835 #[tokio::test]
1839 async fn a_publication_failure_is_filed_under_what_happened() {
1840 let zero = std::time::Duration::ZERO;
1841 const GONE: &str = "did:plc:tombstonedaaaaaaaaaaaaaa";
1842 let plc = serve_answering(GONE, 404, (200, ""), zero, false).await;
1843 let outcome = poll_publication_at(GONE, plc, None).await;
1844 assert_eq!(
1845 kind_of(&outcome),
1846 Some(crate::feed::FailureKind::Status),
1847 "PLC 404: {outcome:?}"
1848 );
1849
1850 const NOREPO: &str = "did:plc:norepoaaaaaaaaaaaaaaaaaa";
1851 let plc = serve_answering(
1852 NOREPO,
1853 200,
1854 (400, r#"{"error":"RepoNotFound"}"#),
1855 zero,
1856 false,
1857 )
1858 .await;
1859 let outcome = poll_publication_at(NOREPO, plc, None).await;
1860 assert_eq!(
1861 kind_of(&outcome),
1862 Some(crate::feed::FailureKind::Status),
1863 "RepoNotFound: {outcome:?}"
1864 );
1865
1866 const GARBLED: &str = "did:plc:garbledaaaaaaaaaaaaaaaaa";
1867 let plc = serve_answering(GARBLED, 200, (200, r#"{"records":"x"}"#), zero, false).await;
1868 let outcome = poll_publication_at(GARBLED, plc, None).await;
1869 assert_eq!(
1870 kind_of(&outcome),
1871 Some(crate::feed::FailureKind::Parse),
1872 "garbled body: {outcome:?}"
1873 );
1874 }
1875
1876 #[tokio::test]
1880 async fn a_refused_listing_is_filed_under_what_the_pds_sent() {
1881 use crate::feed::FailureKind;
1882 let zero = std::time::Duration::ZERO;
1883 let oversized: &'static str =
1886 Box::leak("x".repeat(crate::net::MAX_BODY_BYTES + 1).into_boxed_str());
1887 let mut misfiled = Vec::new();
1888 for (did, (status, body), endless, want, label) in [
1889 (
1890 "did:plc:emptybodyaaaaaaaaaaaaaaa",
1891 (200, ""),
1892 false,
1893 FailureKind::Parse,
1894 "a 200 with an empty body",
1895 ),
1896 (
1897 "did:plc:norecordsaaaaaaaaaaaaaaa",
1898 (200, "{}"),
1899 false,
1900 FailureKind::Parse,
1901 "a 200 with no records field",
1902 ),
1903 (
1904 "did:plc:envelopeaaaaaaaaaaaaaaaa",
1905 (200, r#"{"error":"RepoDeactivated"}"#),
1906 false,
1907 FailureKind::Status,
1908 "a 200 carrying an error envelope",
1909 ),
1910 (
1911 "did:plc:endlessaaaaaaaaaaaaaaaaa",
1912 (200, ""),
1913 true,
1914 FailureKind::Body,
1915 "more pages than MAX_LIST_PAGES",
1916 ),
1917 (
1918 "did:plc:oversizedaaaaaaaaaaaaaaa",
1919 (200, oversized),
1920 false,
1921 FailureKind::Body,
1922 "a body over read_capped's cap",
1923 ),
1924 ] {
1925 let plc = serve_answering(did, 200, (status, body), zero, endless).await;
1926 let outcome = poll_publication_at(did, plc, None).await;
1927 if kind_of(&outcome) != Some(want) {
1928 misfiled.push(format!("{label}: want {want:?}, got {outcome:?}"));
1929 }
1930 }
1931 assert!(misfiled.is_empty(), "misfiled:\n{}", misfiled.join("\n"));
1932 }
1933
1934 #[tokio::test]
1937 async fn an_empty_publication_is_a_healthy_poll() {
1938 const EMPTY: &str = "did:plc:emptypubaaaaaaaaaaaaaaaa";
1939 let site = format!("at://{EMPTY}/{}/quiet", nsid::STANDARD_PUBLICATION);
1940 let (plc, _) = serve_repo(
1941 EMPTY,
1942 vec![(
1943 nsid::STANDARD_PUBLICATION,
1944 "quiet",
1945 json!({ "name": "Quiet", "url": "https://quiet.example" }),
1946 )],
1947 )
1948 .await;
1949 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1950 crate::store::upsert_feed(
1951 &pool,
1952 &crate::store::NewFeed {
1953 url: site.clone(),
1954 ..Default::default()
1955 },
1956 )
1957 .await
1958 .unwrap();
1959 let feed = crate::store::get_feed_by_url(&pool, &site)
1960 .await
1961 .unwrap()
1962 .unwrap();
1963 let mut config = crate::config::Config::default();
1964 config.oauth.plc_directory = plc;
1965 let client = crate::feed::build_client().unwrap();
1966 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1967 .await
1968 .unwrap();
1969 assert!(
1970 matches!(
1971 outcome,
1972 crate::feed::PollOutcome::Updated { new_entries: 0 }
1973 ),
1974 "an empty publication was not a healthy poll: {outcome:?}"
1975 );
1976 }
1977
1978 #[tokio::test]
1982 async fn fetch_skips_malformed_records_beside_good_ones() {
1983 const OWN: &str = "did:plc:malformedrepo";
1984 let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
1985 let doc = |title: &str| {
1986 json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
1987 "path": "/p", "site": site })
1988 };
1989 let (plc, _) = serve_repo(
1990 OWN,
1991 vec![
1992 (
1993 nsid::STANDARD_PUBLICATION,
1994 "",
1995 json!({ "name": "Broken", "url": "https://x.example" }),
1996 ),
1997 (
1998 nsid::STANDARD_PUBLICATION,
1999 "mine",
2000 json!({ "name": "Mine", "url": "https://mine.example" }),
2001 ),
2002 (nsid::STANDARD_DOCUMENT, "", doc("unreadable")),
2003 (nsid::STANDARD_DOCUMENT, "3l2mfaaaaaa2a", doc("kept")),
2004 ],
2005 )
2006 .await;
2007 let client = crate::feed::build_client().unwrap();
2008 let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
2009 .await
2010 .expect("a malformed record stalled a stranger's publication");
2011 assert!(read.complete);
2012 let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
2013 assert_eq!(titles, vec!["kept"]);
2014 }
2015
2016 const SHARED: &str = "did:plc:sharedrepoaaaaaaaaaaaaaa";
2019
2020 fn shared_doc(site_rkey: &str, title: &str, at: &str) -> serde_json::Value {
2021 json!({ "title": title, "publishedAt": at, "path": format!("/{title}"),
2022 "site": format!("at://{SHARED}/{}/{site_rkey}", nsid::STANDARD_PUBLICATION) })
2023 }
2024
2025 #[tokio::test]
2029 async fn two_publications_in_one_repo_cost_one_walk() {
2030 let (plc, hits) = serve_repo(
2031 SHARED,
2032 vec![
2033 (
2034 nsid::STANDARD_PUBLICATION,
2035 "alpha",
2036 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2037 ),
2038 (
2039 nsid::STANDARD_PUBLICATION,
2040 "beta",
2041 json!({ "name": "Beta", "url": "https://beta.example" }),
2042 ),
2043 (
2044 nsid::STANDARD_DOCUMENT,
2045 "3l2shaaaaaa2a",
2046 shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
2047 ),
2048 (
2049 nsid::STANDARD_DOCUMENT,
2050 "3l2shaaaaaa2b",
2051 shared_doc("beta", "b1", "2026-07-10T00:00:00Z"),
2052 ),
2053 (
2054 nsid::STANDARD_DOCUMENT,
2055 "3l2shaaaaaa2c",
2056 shared_doc("alpha", "a2", "2026-07-09T00:00:00Z"),
2057 ),
2058 ],
2059 )
2060 .await;
2061 let client = crate::feed::build_client().unwrap();
2062 let reads = fetch_repo(
2063 &client,
2064 &plc,
2065 SHARED,
2066 &["alpha".to_string(), "beta".to_string()],
2067 )
2068 .await
2069 .unwrap();
2070 assert_eq!(
2071 hits.load(std::sync::atomic::Ordering::SeqCst),
2072 4,
2073 "two collections walked once each (a page then an empty page), not once per publication"
2074 );
2075 let titles = |r: &anyhow::Result<PublicationRead>| {
2076 let mut t: Vec<String> = r
2077 .as_ref()
2078 .unwrap()
2079 .entries
2080 .iter()
2081 .map(|e| e.title.clone())
2082 .collect();
2083 t.sort();
2084 t
2085 };
2086 assert_eq!(titles(&reads[0]), vec!["a1", "a2"]);
2087 assert_eq!(titles(&reads[1]), vec!["b1"]);
2088 }
2089
2090 #[tokio::test]
2093 async fn a_busy_publication_does_not_starve_its_quiet_sibling() {
2094 let mut records = vec![
2095 (
2096 nsid::STANDARD_PUBLICATION,
2097 "busy",
2098 json!({ "name": "Busy", "url": "https://busy.example" }),
2099 ),
2100 (
2101 nsid::STANDARD_PUBLICATION,
2102 "quiet",
2103 json!({ "name": "Quiet", "url": "https://quiet.example" }),
2104 ),
2105 ];
2106 for (rkey, title) in [
2107 ("3l2bsaaaaaa2a", "b1"),
2108 ("3l2bsaaaaaa2b", "b2"),
2109 ("3l2bsaaaaaa2c", "b3"),
2110 ("3l2bsaaaaaa2d", "b4"),
2111 ("3l2bsaaaaaa2e", "b5"),
2112 ] {
2113 records.push((
2114 nsid::STANDARD_DOCUMENT,
2115 rkey,
2116 shared_doc("busy", title, "2026-07-11T00:00:00Z"),
2117 ));
2118 }
2119 records.push((
2120 nsid::STANDARD_DOCUMENT,
2121 "3l2bsaaaaaa2f",
2122 shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
2123 ));
2124 let (plc, _) = serve_repo(SHARED, records).await;
2125 let client = crate::feed::build_client().unwrap();
2126 let reads = fetch_repo_capped(
2127 &client,
2128 &plc,
2129 SHARED,
2130 &["busy".to_string(), "quiet".to_string()],
2131 2,
2132 crate::atproto::MAX_LIST_BYTES,
2133 )
2134 .await
2135 .unwrap();
2136 let busy = reads[0].as_ref().unwrap();
2137 let quiet = reads[1].as_ref().unwrap();
2138 assert_eq!(
2139 busy.entries.len(),
2140 2,
2141 "the busy publication was not capped at its own cap"
2142 );
2143 assert!(!busy.complete, "a capped publication was reported complete");
2144 assert_eq!(quiet.entries.len(), 1, "the quiet sibling was starved");
2145 assert!(quiet.complete);
2146 }
2147
2148 fn starved_group_records() -> Vec<(&'static str, &'static str, serde_json::Value)> {
2160 let body = "w".repeat(20 * 1024);
2161 let mut records = vec![
2162 (
2163 nsid::STANDARD_PUBLICATION,
2164 "big1",
2165 json!({ "name": "Big 1", "url": "https://b1.example" }),
2166 ),
2167 (
2168 nsid::STANDARD_PUBLICATION,
2169 "big2",
2170 json!({ "name": "Big 2", "url": "https://b2.example" }),
2171 ),
2172 (
2173 nsid::STANDARD_PUBLICATION,
2174 "quiet",
2175 json!({ "name": "Quiet", "url": "https://q.example" }),
2176 ),
2177 ];
2178 for i in 0..200 {
2179 let rkey: &'static str = Box::leak(format!("3l2big{i:06}").into_boxed_str());
2180 let site = if i % 2 == 0 { "big1" } else { "big2" };
2181 let mut doc = shared_doc(site, &format!("d{i}"), "2026-07-11T00:00:00Z");
2182 doc["textContent"] = json!(body);
2183 records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2184 }
2185 records.push((
2186 nsid::STANDARD_DOCUMENT,
2187 "3l2zzzzzzzzzz",
2188 shared_doc("quiet", "q1", "2026-07-11T00:00:00Z"),
2189 ));
2190 records
2191 }
2192
2193 const STARVED_BUDGET: usize = 4 * 1024 * 1024;
2194
2195 fn starved_rkeys() -> Vec<String> {
2196 ["big1", "big2", "quiet"]
2197 .iter()
2198 .map(|s| s.to_string())
2199 .collect()
2200 }
2201
2202 async fn shared_feeds(pool: &sqlx::SqlitePool, rkeys: &[String]) -> Vec<crate::store::Feed> {
2204 let mut feeds = Vec::new();
2205 for rkey in rkeys {
2206 let url = format!("at://{SHARED}/{}/{rkey}", nsid::STANDARD_PUBLICATION);
2207 crate::store::upsert_feed(
2208 pool,
2209 &crate::store::NewFeed {
2210 url: url.clone(),
2211 ..Default::default()
2212 },
2213 )
2214 .await
2215 .unwrap();
2216 feeds.push(
2217 crate::store::get_feed_by_url(pool, &url)
2218 .await
2219 .unwrap()
2220 .unwrap(),
2221 );
2222 }
2223 feeds
2224 }
2225
2226 async fn stored_count(pool: &sqlx::SqlitePool, feed_id: i64) -> i64 {
2227 sqlx::query_scalar("SELECT COUNT(*) FROM entries WHERE feed_id = ?")
2228 .bind(feed_id)
2229 .fetch_one(pool)
2230 .await
2231 .unwrap()
2232 }
2233
2234 fn hits_of(hits: &std::sync::Arc<std::sync::atomic::AtomicUsize>) -> usize {
2235 hits.load(std::sync::atomic::Ordering::SeqCst)
2236 }
2237
2238 #[tokio::test]
2242 async fn a_group_read_flags_the_member_its_walk_cut_short() {
2243 let (plc, _) = serve_repo(SHARED, starved_group_records()).await;
2244 let client = crate::feed::build_client().unwrap();
2245 let reads = fetch_repo_capped(
2246 &client,
2247 &plc,
2248 SHARED,
2249 &starved_rkeys(),
2250 2_000,
2251 STARVED_BUDGET,
2252 )
2253 .await
2254 .unwrap();
2255 let quiet = reads[2].as_ref().unwrap();
2256 assert_eq!(
2257 (
2258 quiet.entries.len(),
2259 quiet.complete,
2260 quiet.cut_short_by_group
2261 ),
2262 (0, false, true),
2263 "the quiet publication's group read was not flagged as cut short by its group"
2264 );
2265 }
2266
2267 #[tokio::test]
2274 async fn a_walk_that_stops_at_the_combined_cap_flags_no_one() {
2275 let mut records = vec![
2276 (
2277 nsid::STANDARD_PUBLICATION,
2278 "alpha",
2279 json!({ "name": "Alpha", "url": "https://a.example" }),
2280 ),
2281 (
2282 nsid::STANDARD_PUBLICATION,
2283 "beta",
2284 json!({ "name": "Beta", "url": "https://b.example" }),
2285 ),
2286 (
2287 nsid::STANDARD_PUBLICATION,
2288 "other",
2289 json!({ "name": "Other", "url": "https://o.example" }),
2290 ),
2291 (
2292 nsid::STANDARD_DOCUMENT,
2293 "3l2capaaaaa00",
2294 shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
2295 ),
2296 (
2297 nsid::STANDARD_DOCUMENT,
2298 "3l2capaaaaa01",
2299 shared_doc("beta", "b1", "2026-07-11T00:00:00Z"),
2300 ),
2301 ];
2302 for i in 2..DOCUMENT_PAGE_SIZE as usize {
2304 let rkey: &'static str = Box::leak(format!("3l2capaaaaa{i:02}").into_boxed_str());
2305 records.push((
2306 nsid::STANDARD_DOCUMENT,
2307 rkey,
2308 shared_doc("other", &format!("o{i}"), "2026-07-11T00:00:00Z"),
2309 ));
2310 }
2311 records.push((
2313 nsid::STANDARD_DOCUMENT,
2314 "3l2capaaaaz99",
2315 shared_doc("alpha", "a2", "2026-07-11T00:00:00Z"),
2316 ));
2317 let (plc, _) = serve_repo(SHARED, records).await;
2318 let client = crate::feed::build_client().unwrap();
2319 let rkeys = vec!["alpha".to_string(), "beta".to_string()];
2320 let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 1, STARVED_BUDGET)
2321 .await
2322 .unwrap();
2323 for (name, read) in rkeys.iter().zip(&reads) {
2324 let read = read.as_ref().unwrap();
2325 assert_eq!(read.entries.len(), 1, "{name}");
2326 assert!(
2327 !read.cut_short_by_group,
2328 "{name} was flagged for a re-read that would stop at the same place"
2329 );
2330 }
2331 }
2332
2333 #[tokio::test]
2340 async fn the_publication_its_group_starved_is_re_read_first() {
2341 let client = crate::feed::build_client().unwrap();
2342 let (plc, hits) = serve_repo(SHARED, starved_group_records()).await;
2343 fetch_repo_capped(
2344 &client,
2345 &plc,
2346 SHARED,
2347 &starved_rkeys(),
2348 2_000,
2349 STARVED_BUDGET,
2350 )
2351 .await
2352 .unwrap();
2353 let group_read = hits_of(&hits);
2354 let (plc, hits) = serve_repo(SHARED, starved_group_records()).await;
2355 fetch_repo_capped(
2356 &client,
2357 &plc,
2358 SHARED,
2359 &starved_rkeys()[2..],
2360 2_000,
2361 STARVED_BUDGET,
2362 )
2363 .await
2364 .unwrap();
2365 let quiet_alone = hits_of(&hits);
2366
2367 let (plc, _) = serve_repo_slow_after(
2370 SHARED,
2371 starved_group_records(),
2372 group_read + quiet_alone,
2373 std::time::Duration::from_secs(60),
2374 )
2375 .await;
2376 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2377 let feeds = shared_feeds(&pool, &starved_rkeys()).await;
2378 let mut config = crate::config::Config::default();
2379 config.oauth.plc_directory = plc;
2380 config.publication_read_deadline = std::time::Duration::from_secs(3);
2381 let outcomes = crate::feed::poll_publication_group_with(
2382 &pool,
2383 &client,
2384 &config,
2385 &feeds,
2386 2_000,
2387 STARVED_BUDGET,
2388 std::time::Instant::now(),
2389 )
2390 .await;
2391 assert!(
2392 matches!(
2393 outcomes[2].as_ref().unwrap(),
2394 crate::feed::PollOutcome::Updated { .. }
2395 ),
2396 "a big sibling's re-read spent the deadline and the starved one stayed: {:?}",
2397 outcomes[2]
2398 );
2399 }
2400
2401 #[tokio::test]
2405 async fn big_siblings_do_not_spend_a_quiet_publications_share() {
2406 let (plc, _) = serve_repo(SHARED, starved_group_records()).await;
2407 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2408 let feeds = shared_feeds(&pool, &starved_rkeys()).await;
2409 let mut config = crate::config::Config::default();
2410 config.oauth.plc_directory = plc;
2411 let client = crate::feed::build_client().unwrap();
2412 let outcomes = crate::feed::poll_publication_group_with(
2413 &pool,
2414 &client,
2415 &config,
2416 &feeds,
2417 2_000,
2418 STARVED_BUDGET,
2419 std::time::Instant::now(),
2420 )
2421 .await;
2422 assert!(
2423 matches!(
2424 outcomes[2].as_ref().unwrap(),
2425 crate::feed::PollOutcome::Updated { .. }
2426 ),
2427 "the quiet publication was starved by its siblings: {:?}",
2428 outcomes[2]
2429 );
2430 assert_eq!(stored_count(&pool, feeds[2].id).await, 1);
2431 }
2432
2433 #[tokio::test]
2436 async fn a_big_sibling_cut_short_by_its_group_reads_as_it_would_alone() {
2437 let (plc, _) = serve_repo(SHARED, starved_group_records()).await;
2438 let client = crate::feed::build_client().unwrap();
2439 let rkeys = starved_rkeys();
2440 let alone = fetch_repo_capped(&client, &plc, SHARED, &rkeys[..1], 2_000, STARVED_BUDGET)
2441 .await
2442 .unwrap();
2443 let alone = alone[0].as_ref().unwrap();
2444 let grouped = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, STARVED_BUDGET)
2445 .await
2446 .unwrap();
2447 let grouped = grouped[0].as_ref().unwrap();
2448 assert!(
2449 grouped.entries.len() < alone.entries.len() && grouped.cut_short_by_group,
2450 "the scenario no longer cuts the big publication short in its group \
2451 ({} grouped, {} alone)",
2452 grouped.entries.len(),
2453 alone.entries.len()
2454 );
2455
2456 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2457 let feeds = shared_feeds(&pool, &rkeys).await;
2458 let mut config = crate::config::Config::default();
2459 config.oauth.plc_directory = plc;
2460 let outcomes = crate::feed::poll_publication_group_with(
2461 &pool,
2462 &client,
2463 &config,
2464 &feeds,
2465 2_000,
2466 STARVED_BUDGET,
2467 std::time::Instant::now(),
2468 )
2469 .await;
2470 assert!(matches!(
2471 outcomes[0].as_ref().unwrap(),
2472 crate::feed::PollOutcome::Updated { .. }
2473 ));
2474 assert_eq!(
2475 stored_count(&pool, feeds[0].id).await,
2476 alone.entries.len() as i64,
2477 "the group stored the big publication worse than it reads alone"
2478 );
2479 }
2480
2481 #[tokio::test]
2485 async fn a_member_not_cut_short_by_its_group_is_not_read_again() {
2486 let mut records = vec![
2487 (
2488 nsid::STANDARD_PUBLICATION,
2489 "busy",
2490 json!({ "name": "Busy", "url": "https://busy.example" }),
2491 ),
2492 (
2493 nsid::STANDARD_PUBLICATION,
2494 "quiet",
2495 json!({ "name": "Quiet", "url": "https://quiet.example" }),
2496 ),
2497 ];
2498 for i in 0..5 {
2499 let rkey: &'static str = Box::leak(format!("3l2nrr{i:06}").into_boxed_str());
2500 records.push((
2501 nsid::STANDARD_DOCUMENT,
2502 rkey,
2503 shared_doc("busy", &format!("b{i}"), "2026-07-11T00:00:00Z"),
2504 ));
2505 }
2506 records.push((
2507 nsid::STANDARD_DOCUMENT,
2508 "3l2nrrzzzzzzz",
2509 shared_doc("quiet", "q1", "2026-07-11T00:00:00Z"),
2510 ));
2511 let (plc, hits) = serve_repo(SHARED, records).await;
2512 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2513 let rkeys = vec!["busy".to_string(), "quiet".to_string()];
2514 let feeds = shared_feeds(&pool, &rkeys).await;
2515 let mut config = crate::config::Config::default();
2516 config.oauth.plc_directory = plc;
2517 let client = crate::feed::build_client().unwrap();
2518 let outcomes = crate::feed::poll_publication_group_with(
2519 &pool,
2520 &client,
2521 &config,
2522 &feeds,
2523 2,
2524 crate::atproto::MAX_LIST_BYTES,
2525 std::time::Instant::now(),
2526 )
2527 .await;
2528 assert_eq!(
2529 hits_of(&hits),
2530 4,
2531 "a publication not cut short by its group was read again"
2532 );
2533 assert_eq!(stored_count(&pool, feeds[0].id).await, 2);
2534 assert_eq!(stored_count(&pool, feeds[1].id).await, 1);
2535 assert!(outcomes.iter().all(|o| matches!(
2536 o.as_ref().unwrap(),
2537 crate::feed::PollOutcome::Updated { .. }
2538 )));
2539 }
2540
2541 #[tokio::test]
2545 async fn a_member_at_its_own_cap_is_not_read_again_when_its_group_is_cut_short() {
2546 let mut records = vec![
2547 (
2548 nsid::STANDARD_PUBLICATION,
2549 "capped",
2550 json!({ "name": "Capped", "url": "https://c.example" }),
2551 ),
2552 (
2553 nsid::STANDARD_PUBLICATION,
2554 "big",
2555 json!({ "name": "Big", "url": "https://b.example" }),
2556 ),
2557 ];
2558 for i in 0..60 {
2559 let rkey: &'static str = Box::leak(format!("3l2cap{i:06}").into_boxed_str());
2560 records.push((
2561 nsid::STANDARD_DOCUMENT,
2562 rkey,
2563 shared_doc("capped", &format!("c{i}"), "2026-07-11T00:00:00Z"),
2564 ));
2565 }
2566 let body = "w".repeat(20 * 1024);
2567 for i in 0..100 {
2568 let rkey: &'static str = Box::leak(format!("3l2cbg{i:06}").into_boxed_str());
2569 let mut doc = shared_doc("big", &format!("d{i}"), "2026-07-11T00:00:00Z");
2570 doc["textContent"] = json!(body);
2571 records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2572 }
2573 let (plc, hits) = serve_repo(SHARED, records).await;
2574 let client = crate::feed::build_client().unwrap();
2575 let rkeys = vec!["capped".to_string(), "big".to_string()];
2576 let (cap, budget) = (50, 1024 * 1024);
2577 let grouped = fetch_repo_capped(&client, &plc, SHARED, &rkeys, cap, budget)
2578 .await
2579 .unwrap();
2580 let (c, b) = (grouped[0].as_ref().unwrap(), grouped[1].as_ref().unwrap());
2581 assert_eq!(
2582 (c.entries.len(), c.cut_short_by_group, b.cut_short_by_group),
2583 (cap, false, true),
2584 "the scenario no longer caps one member and cuts the other short"
2585 );
2586 let group_read = hits_of(&hits);
2587 fetch_repo_capped(&client, &plc, SHARED, &rkeys[1..], cap, budget)
2588 .await
2589 .unwrap();
2590 let big_alone = hits_of(&hits) - group_read;
2591
2592 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2593 let feeds = shared_feeds(&pool, &rkeys).await;
2594 let mut config = crate::config::Config::default();
2595 config.oauth.plc_directory = plc;
2596 let before = hits_of(&hits);
2597 crate::feed::poll_publication_group_with(
2598 &pool,
2599 &client,
2600 &config,
2601 &feeds,
2602 cap,
2603 budget,
2604 std::time::Instant::now(),
2605 )
2606 .await;
2607 assert_eq!(
2608 hits_of(&hits) - before,
2609 group_read + big_alone,
2610 "a member at its own cap was read again"
2611 );
2612 }
2613
2614 #[tokio::test]
2617 async fn a_lone_publication_is_never_read_again() {
2618 let mut records = starved_group_records();
2619 records.retain(|(c, rkey, v)| {
2620 (*c == nsid::STANDARD_PUBLICATION && *rkey == "big1")
2621 || v["site"].as_str().is_some_and(|s| s.ends_with("/big1"))
2622 });
2623 let (plc, hits) = serve_repo(SHARED, records).await;
2624 let client = crate::feed::build_client().unwrap();
2625 let rkeys = vec!["big1".to_string()];
2626 let budget = 1024 * 1024;
2627 let alone = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, budget)
2628 .await
2629 .unwrap();
2630 let alone = alone[0].as_ref().unwrap();
2631 assert!(
2632 !alone.complete && !alone.cut_short_by_group,
2633 "a lone publication's own incomplete read was blamed on a group"
2634 );
2635 let one_read = hits_of(&hits);
2636
2637 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2638 let feeds = shared_feeds(&pool, &rkeys).await;
2639 let mut config = crate::config::Config::default();
2640 config.oauth.plc_directory = plc;
2641 crate::feed::poll_publication_group_with(
2642 &pool,
2643 &client,
2644 &config,
2645 &feeds,
2646 2_000,
2647 budget,
2648 std::time::Instant::now(),
2649 )
2650 .await;
2651 assert_eq!(
2652 hits_of(&hits) - one_read,
2653 one_read,
2654 "a lone publication was read again"
2655 );
2656 }
2657
2658 #[tokio::test]
2661 async fn no_re_read_once_the_read_deadline_is_spent() {
2662 let (plc, hits) = serve_repo(SHARED, starved_group_records()).await;
2663 let client = crate::feed::build_client().unwrap();
2664 fetch_repo_capped(
2665 &client,
2666 &plc,
2667 SHARED,
2668 &starved_rkeys(),
2669 2_000,
2670 STARVED_BUDGET,
2671 )
2672 .await
2673 .unwrap();
2674 let group_read = hits_of(&hits);
2675
2676 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2677 let feeds = shared_feeds(&pool, &starved_rkeys()).await;
2678 let mut config = crate::config::Config::default();
2679 config.oauth.plc_directory = plc;
2680 let deadline = std::time::Duration::from_secs(10);
2681 config.publication_read_deadline = deadline;
2682 let spent = std::time::Instant::now().checked_sub(deadline).unwrap();
2685 let outcomes = crate::feed::poll_publication_group_with(
2686 &pool,
2687 &client,
2688 &config,
2689 &feeds,
2690 2_000,
2691 STARVED_BUDGET,
2692 spent,
2693 )
2694 .await;
2695 assert_eq!(
2696 hits_of(&hits) - group_read,
2697 group_read,
2698 "a re-read ran after the read deadline was spent"
2699 );
2700 assert!(
2701 matches!(
2702 outcomes[2].as_ref().unwrap(),
2703 crate::feed::PollOutcome::Failed { .. }
2704 ),
2705 "the group's outcome did not stand: {:?}",
2706 outcomes[2]
2707 );
2708 }
2709
2710 #[tokio::test]
2713 async fn a_re_read_that_overruns_the_deadline_stops_the_re_reads() {
2714 let (plc, hits) = serve_repo(SHARED, starved_group_records()).await;
2715 let client = crate::feed::build_client().unwrap();
2716 fetch_repo_capped(
2717 &client,
2718 &plc,
2719 SHARED,
2720 &starved_rkeys(),
2721 2_000,
2722 STARVED_BUDGET,
2723 )
2724 .await
2725 .unwrap();
2726 let group_read = hits_of(&hits);
2727
2728 let (plc, hits) = serve_repo_slow_after(
2730 SHARED,
2731 starved_group_records(),
2732 group_read,
2733 std::time::Duration::from_secs(60),
2734 )
2735 .await;
2736 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2737 let feeds = shared_feeds(&pool, &starved_rkeys()).await;
2738 let mut config = crate::config::Config::default();
2739 config.oauth.plc_directory = plc;
2740 config.publication_read_deadline = std::time::Duration::from_secs(3);
2741 let started = std::time::Instant::now();
2742 let outcomes = crate::feed::poll_publication_group_with(
2743 &pool,
2744 &client,
2745 &config,
2746 &feeds,
2747 2_000,
2748 STARVED_BUDGET,
2749 started,
2750 )
2751 .await;
2752 assert!(
2753 started.elapsed() < std::time::Duration::from_secs(10),
2754 "the re-reads were not bounded by the read deadline"
2755 );
2756 assert_eq!(
2757 hits_of(&hits),
2758 group_read + 1,
2759 "re-reads went on after one overran the deadline"
2760 );
2761 assert!(
2762 matches!(
2763 outcomes[2].as_ref().unwrap(),
2764 crate::feed::PollOutcome::Failed { .. }
2765 ),
2766 "the group's outcome did not stand: {:?}",
2767 outcomes[2]
2768 );
2769 }
2770
2771 #[tokio::test]
2776 async fn a_re_read_that_fails_keeps_the_group_outcome() {
2777 let (plc, hits) = serve_repo(SHARED, starved_group_records()).await;
2778 let client = crate::feed::build_client().unwrap();
2779 fetch_repo_capped(
2780 &client,
2781 &plc,
2782 SHARED,
2783 &starved_rkeys(),
2784 2_000,
2785 STARVED_BUDGET,
2786 )
2787 .await
2788 .unwrap();
2789 let group_read = hits_of(&hits);
2790
2791 let (plc, _) =
2792 serve_repo_slow_after(SHARED, starved_group_records(), group_read, FAIL_AFTER).await;
2793 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2794 let feeds = shared_feeds(&pool, &starved_rkeys()).await;
2795 let mut config = crate::config::Config::default();
2796 config.oauth.plc_directory = plc;
2797 let outcomes = crate::feed::poll_publication_group_with(
2798 &pool,
2799 &client,
2800 &config,
2801 &feeds,
2802 2_000,
2803 STARVED_BUDGET,
2804 std::time::Instant::now(),
2805 )
2806 .await;
2807 assert!(
2808 matches!(
2809 outcomes[0].as_ref().unwrap(),
2810 crate::feed::PollOutcome::Updated { .. }
2811 ),
2812 "a failed re-read turned the group's Updated into: {:?}",
2813 outcomes[0]
2814 );
2815 }
2816
2817 #[tokio::test]
2822 async fn a_busy_publication_beside_idle_siblings_reads_as_it_would_alone() {
2823 let body = "w".repeat(20 * 1024);
2824 let mut records: Vec<(&'static str, &'static str, serde_json::Value)> = Vec::new();
2825 let rkeys: Vec<String> = (0..16).map(|i| format!("p{i:02}")).collect();
2826 for r in &rkeys {
2827 let r: &'static str = Box::leak(r.clone().into_boxed_str());
2828 records.push((
2829 nsid::STANDARD_PUBLICATION,
2830 r,
2831 json!({ "name": r, "url": "https://p.example" }),
2832 ));
2833 }
2834 for i in 0..100 {
2835 let rkey: &'static str = Box::leak(format!("3l2bus{i:06}").into_boxed_str());
2836 let mut doc = shared_doc("p00", &format!("d{i}"), "2026-07-11T00:00:00Z");
2837 doc["textContent"] = json!(body);
2838 records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2839 }
2840 let (plc, _) = serve_repo(SHARED, records).await;
2841 let client = crate::feed::build_client().unwrap();
2842 let budget = 4 * 1024 * 1024;
2843 let alone = fetch_repo_capped(&client, &plc, SHARED, &rkeys[..1], 2_000, budget)
2844 .await
2845 .unwrap();
2846 let grouped = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, budget)
2847 .await
2848 .unwrap();
2849 let (a, g) = (alone[0].as_ref().unwrap(), grouped[0].as_ref().unwrap());
2850 assert_eq!((a.entries.len(), a.complete), (100, true));
2851 assert_eq!(
2852 (g.entries.len(), g.complete),
2853 (100, true),
2854 "grouping read it worse than alone"
2855 );
2856 }
2857
2858 #[tokio::test]
2861 async fn one_large_document_does_not_fail_its_publication_in_a_group() {
2862 let mut big = shared_doc("big", "huge", "2026-07-11T00:00:00Z");
2863 big["textContent"] = json!("w".repeat(500 * 1024));
2864 let records = vec![
2865 (
2866 nsid::STANDARD_PUBLICATION,
2867 "big",
2868 json!({ "name": "Big", "url": "https://b.example" }),
2869 ),
2870 (
2871 nsid::STANDARD_PUBLICATION,
2872 "other",
2873 json!({ "name": "Other", "url": "https://o.example" }),
2874 ),
2875 (nsid::STANDARD_DOCUMENT, "3l2hugeaaaa2a", big),
2876 ];
2877 let (plc, _) = serve_repo(SHARED, records).await;
2878 let client = crate::feed::build_client().unwrap();
2879 let rkeys = vec!["big".to_string(), "other".to_string()];
2880 let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 1024 * 1024)
2881 .await
2882 .unwrap();
2883 assert_eq!(
2884 reads[0].as_ref().unwrap().entries.len(),
2885 1,
2886 "a large document was dropped"
2887 );
2888 }
2889
2890 #[tokio::test]
2893 async fn every_requested_rkey_gets_a_result() {
2894 let (plc, _) = serve_repo(
2895 SHARED,
2896 vec![(
2897 nsid::STANDARD_PUBLICATION,
2898 "alpha",
2899 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2900 )],
2901 )
2902 .await;
2903 let client = crate::feed::build_client().unwrap();
2904 let reads = fetch_repo(
2905 &client,
2906 &plc,
2907 SHARED,
2908 &["nope".to_string(), "nope".to_string()],
2909 )
2910 .await
2911 .unwrap();
2912 assert_eq!(reads.len(), 2);
2913 assert!(reads.iter().all(|r| r.is_err()));
2914 }
2915
2916 #[tokio::test]
2919 async fn repeated_and_empty_requests_are_handled() {
2920 let (plc, hits) = serve_repo(
2921 SHARED,
2922 vec![
2923 (
2924 nsid::STANDARD_PUBLICATION,
2925 "alpha",
2926 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2927 ),
2928 (
2929 nsid::STANDARD_DOCUMENT,
2930 "3l2rpaaaaaa2a",
2931 shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
2932 ),
2933 ],
2934 )
2935 .await;
2936 let client = crate::feed::build_client().unwrap();
2937 let none = fetch_repo(&client, &plc, SHARED, &[]).await.unwrap();
2938 assert!(none.is_empty());
2939 assert_eq!(
2940 hits.load(std::sync::atomic::Ordering::SeqCst),
2941 0,
2942 "an empty request reached the network"
2943 );
2944 let twice = fetch_repo(
2945 &client,
2946 &plc,
2947 SHARED,
2948 &["alpha".to_string(), "alpha".to_string()],
2949 )
2950 .await
2951 .unwrap();
2952 for read in &twice {
2953 assert_eq!(
2954 read.as_ref().unwrap().entries.len(),
2955 1,
2956 "a repeated rkey read as empty"
2957 );
2958 }
2959 }
2960
2961 #[tokio::test]
2964 async fn a_group_stores_each_publications_documents_under_its_own_feed() {
2965 let (plc, _) = serve_repo(
2966 SHARED,
2967 vec![
2968 (
2969 nsid::STANDARD_PUBLICATION,
2970 "alpha",
2971 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2972 ),
2973 (
2974 nsid::STANDARD_PUBLICATION,
2975 "beta",
2976 json!({ "name": "Beta", "url": "https://beta.example" }),
2977 ),
2978 (
2979 nsid::STANDARD_DOCUMENT,
2980 "3l2grpaaaaa2a",
2981 shared_doc("alpha", "only-alpha", "2026-07-11T00:00:00Z"),
2982 ),
2983 (
2984 nsid::STANDARD_DOCUMENT,
2985 "3l2grpaaaaa2b",
2986 shared_doc("beta", "only-beta", "2026-07-10T00:00:00Z"),
2987 ),
2988 ],
2989 )
2990 .await;
2991 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2992 let mut feeds = Vec::new();
2993 for rkey in ["alpha", "beta"] {
2994 let url = format!("at://{SHARED}/{}/{rkey}", nsid::STANDARD_PUBLICATION);
2995 crate::store::upsert_feed(
2996 &pool,
2997 &crate::store::NewFeed {
2998 url: url.clone(),
2999 ..Default::default()
3000 },
3001 )
3002 .await
3003 .unwrap();
3004 feeds.push(
3005 crate::store::get_feed_by_url(&pool, &url)
3006 .await
3007 .unwrap()
3008 .unwrap(),
3009 );
3010 }
3011 let mut config = crate::config::Config::default();
3012 config.oauth.plc_directory = plc;
3013 let client = crate::feed::build_client().unwrap();
3014 let outcomes = crate::feed::poll_publication_group(&pool, &client, &config, &feeds).await;
3015 assert_eq!(outcomes.len(), 2);
3016 for (feed, want) in feeds.iter().zip(["only-alpha", "only-beta"]) {
3017 let titles: Vec<String> =
3018 sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ?")
3019 .bind(feed.id)
3020 .fetch_all(&pool)
3021 .await
3022 .unwrap();
3023 assert_eq!(
3024 titles,
3025 vec![want.to_string()],
3026 "{} got another publication's documents",
3027 feed.url
3028 );
3029 }
3030 }
3031
3032 #[tokio::test]
3038 async fn a_publication_subscription_delivers_entries_end_to_end() {
3039 const A0: &str = "did:plc:acceptanceaaaaaaaaaaaaaa";
3040 let site = format!("at://{A0}/{}/a0pub", nsid::STANDARD_PUBLICATION);
3041 let document = |title: &str, path: &str| {
3042 json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
3043 "path": path, "site": site, "textContent": "body" })
3044 };
3045 let (plc, _hits) = serve_repo(
3046 A0,
3047 vec![
3048 (
3049 nsid::STANDARD_PUBLICATION,
3050 "a0pub",
3051 json!({ "name": "A0 Journal", "url": "https://a0.example" }),
3052 ),
3053 (
3054 nsid::STANDARD_DOCUMENT,
3055 "3l2a0aaaaaa2a",
3056 document("First post", "/first"),
3057 ),
3058 (
3059 nsid::STANDARD_DOCUMENT,
3060 "3l2a0aaaaaa2b",
3061 document("Second post", "/second"),
3062 ),
3063 ],
3064 )
3065 .await;
3066
3067 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
3068 crate::store::upsert_feed(
3069 &pool,
3070 &crate::store::NewFeed {
3071 url: site.clone(),
3072 ..Default::default()
3073 },
3074 )
3075 .await
3076 .unwrap();
3077 let mut config = crate::config::Config::default();
3079 config.oauth.plc_directory = plc;
3080
3081 let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
3082 let feed =
3084 crate::store::due_feeds_of_kind(&pool, &now, crate::feed::FeedKind::Publication, 50)
3085 .await
3086 .unwrap()
3087 .into_iter()
3088 .find(|f| f.url == site)
3089 .expect("the publication is not handed to the publication poller");
3090
3091 let client = crate::feed::build_client().unwrap();
3092 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
3093 .await
3094 .unwrap();
3095 assert!(
3096 matches!(
3097 outcome,
3098 crate::feed::PollOutcome::Updated { new_entries: 2 }
3099 ),
3100 "expected two new entries, got {outcome:?}"
3101 );
3102 let titles: Vec<String> =
3103 sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ? ORDER BY title")
3104 .bind(feed.id)
3105 .fetch_all(&pool)
3106 .await
3107 .unwrap();
3108 assert_eq!(titles, vec!["First post", "Second post"]);
3109 }
3110
3111 fn canonical(rkey: &str) -> String {
3112 format!("at://{DID}/{}/{rkey}", nsid::STANDARD_PUBLICATION)
3113 }
3114
3115 #[test]
3124 fn the_site_filter_uses_the_uri_the_pds_minted() {
3125 let records = vec![publication("p", "https://scanash.com")];
3126 let (site, pubn) = publication_from_records("p", &records).expect("publication not found");
3127 assert_eq!(site, canonical("p"), "did not take the PDS's canonical URI");
3128
3129 let docs = vec![document("d1", &canonical("p"), "Hello", "/hello")];
3130 let entries = entries_from_records(&site, &pubn, &docs);
3131 assert_eq!(
3132 entries.len(),
3133 1,
3134 "a canonical-site document was not matched"
3135 );
3136 }
3137
3138 #[test]
3141 fn an_entry_maps_onto_the_stores_row() {
3142 let records = vec![publication("p", "https://example.com")];
3143 let (site, pubn) = publication_from_records("p", &records).unwrap();
3144 let docs = vec![document("rk1", &site, "Hello", "/hello")];
3145 let row: crate::store::NewEntry = entries_from_records(&site, &pubn, &docs)
3146 .pop()
3147 .unwrap()
3148 .into();
3149 assert_eq!(
3150 row.guid,
3151 format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
3152 );
3153 assert_eq!(row.url.as_deref(), Some("https://example.com/hello"));
3154 assert_eq!(row.title.as_deref(), Some("Hello"));
3155 assert_eq!(row.published.as_deref(), Some("2026-07-11T00:00:00Z"));
3156 assert_eq!(row.content_html.as_deref(), Some("body"));
3157 assert_eq!(row.author, None);
3158 assert_eq!(row.fetched_at, None);
3159 }
3160
3161 #[test]
3164 fn documents_are_filtered_by_their_site_field() {
3165 let records = vec![publication("mine", "https://example.com")];
3166 let (site, pubn) = publication_from_records("mine", &records).unwrap();
3167 let docs = vec![
3168 document("a", &site, "Mine", "/a"),
3169 document("b", &canonical("theirs"), "Theirs", "/b"),
3170 document("c", &site, "Mine again", "/c"),
3171 ];
3172 let titles: Vec<String> = entries_from_records(&site, &pubn, &docs)
3173 .into_iter()
3174 .map(|e| e.title)
3175 .collect();
3176 assert_eq!(titles, ["Mine", "Mine again"]);
3177 }
3178
3179 #[test]
3186 fn a_publication_with_a_hostile_url_is_refused() {
3187 for hostile in [
3188 "javascript:alert(1)",
3189 "data:text/html,<script>",
3190 "file:///etc/passwd",
3191 "",
3192 ] {
3193 let records = vec![publication("p", hostile)];
3194 assert!(
3195 publication_from_records("p", &records).is_none(),
3196 "accepted a publication whose url is {hostile:?}",
3197 );
3198 }
3199 }
3200
3201 #[test]
3207 fn entry_urls_are_joined_against_the_publication_base() {
3208 let records = vec![publication("p", "https://example.com/blog")];
3209 let (site, pubn) = publication_from_records("p", &records).unwrap();
3210 let docs = vec![
3211 document("a", &site, "Relative", "/a"),
3212 document("b", &site, "Absolute-looking", "https://evil.example/x"),
3213 ];
3214 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3215 .into_iter()
3216 .map(|e| e.url)
3217 .collect();
3218 assert_eq!(urls[0].as_deref(), Some("https://example.com/a"));
3219 assert_eq!(
3227 urls[1], None,
3228 "a document path that escapes its publication's origin must yield no URL"
3229 );
3230 }
3231
3232 #[test]
3235 fn a_documents_text_fields_are_bounded_before_they_are_stored() {
3236 let site = canonical("pub");
3237 let big = "x".repeat(8 * 1024 * 1024);
3238 let records = vec![
3239 publication("pub", "https://scanash.com"),
3240 rec(
3241 nsid::STANDARD_DOCUMENT,
3242 "3l2bigaaaaa2a",
3243 json!({ "title": big, "publishedAt": "2026-07-11T00:00:00Z",
3244 "path": format!("/{}", "p".repeat(20_000)), "site": site,
3245 "textContent": "<".repeat(3 * 1024 * 1024) }),
3246 ),
3247 ];
3248 let (_, publication) = publication_from_records("pub", &records).unwrap();
3249 let entries = entries_from_records(&site, &publication, &records);
3250 let e = &entries[0];
3251 assert!(
3252 e.title.len() <= crate::feed::MAX_TITLE_BYTES,
3253 "title: {}",
3254 e.title.len()
3255 );
3256 let url = e
3257 .url
3258 .as_ref()
3259 .expect("an overlong path is truncated, not dropped");
3260 assert!(
3261 url.len() <= crate::feed::MAX_URL_BYTES,
3262 "url: {}",
3263 url.len()
3264 );
3265 let summary = e.summary.as_ref().unwrap();
3266 assert!(
3267 summary.len() <= crate::feed::MAX_CONTENT_HTML_BYTES,
3268 "the ESCAPED summary is what is stored: {}",
3269 summary.len()
3270 );
3271 }
3272
3273 #[test]
3277 fn a_documents_uri_is_bounded_as_an_entry_id() {
3278 let site = canonical("pub");
3279 let records = vec![
3280 publication("pub", "https://scanash.com"),
3281 rec(
3282 nsid::STANDARD_DOCUMENT,
3283 &"k".repeat(100_000),
3284 json!({ "title": "t", "publishedAt": "2026-07-11T00:00:00Z",
3285 "path": "/p", "site": site }),
3286 ),
3287 ];
3288 let (_, publication) = publication_from_records("pub", &records).unwrap();
3289 let entries = entries_from_records(&site, &publication, &records);
3290 let stored: crate::store::NewEntry = entries[0].clone().into();
3291 assert!(
3292 stored.guid.len() <= crate::feed::MAX_GUID_BYTES,
3293 "guid: {}",
3294 stored.guid.len()
3295 );
3296 }
3297
3298 #[test]
3299 fn a_publications_own_name_is_bounded() {
3300 let records = vec![rec(
3301 nsid::STANDARD_PUBLICATION,
3302 "pub",
3303 json!({ "name": "n".repeat(100_000), "url": "https://scanash.com" }),
3304 )];
3305 let (_, publication) = publication_from_records("pub", &records).unwrap();
3306 assert!(publication.name.unwrap().len() <= crate::feed::MAX_TITLE_BYTES);
3307 }
3308
3309 #[test]
3311 fn a_document_with_neither_summary_field_still_yields_an_entry() {
3312 let records = vec![publication("p", "https://example.com/")];
3313 let (site, pubn) = publication_from_records("p", &records).unwrap();
3314 let bare = rec(
3315 nsid::STANDARD_DOCUMENT,
3316 "bare",
3317 json!({
3318 "title": "Bare",
3319 "publishedAt": "2026-07-11T00:00:00Z",
3320 "path": "/bare",
3321 "site": site,
3322 }),
3323 );
3324 let entries = entries_from_records(&site, &pubn, &[bare]);
3325 assert_eq!(entries.len(), 1);
3326 assert_eq!(entries[0].summary, None);
3327 assert_eq!(entries[0].url.as_deref(), Some("https://example.com/bare"));
3328 }
3329
3330 #[test]
3333 fn an_empty_description_does_not_shadow_the_body() {
3334 let records = vec![publication("p", "https://example.com")];
3335 let (site, pubn) = publication_from_records("p", &records).unwrap();
3336 let doc = rec(
3337 nsid::STANDARD_DOCUMENT,
3338 "d",
3339 json!({
3340 "title": "T",
3341 "publishedAt": "2026-07-11T00:00:00Z",
3342 "path": "/d",
3343 "site": site,
3344 "description": " ",
3345 "textContent": "the real body",
3346 }),
3347 );
3348 let entries = entries_from_records(&site, &pubn, &[doc]);
3349 assert_eq!(entries[0].summary.as_deref(), Some("the real body"));
3350 }
3351
3352 #[test]
3359 fn published_at_is_normalised_or_dropped() {
3360 let records = vec![publication("p", "https://example.com")];
3361 let (site, pubn) = publication_from_records("p", &records).unwrap();
3362 let with = |rkey: &str, published_at: serde_json::Value| {
3363 rec(
3364 nsid::STANDARD_DOCUMENT,
3365 rkey,
3366 json!({ "title": "T", "publishedAt": published_at, "path": "/x", "site": site }),
3367 )
3368 };
3369 let docs = vec![
3373 with("a", json!("2026-07-11T09:30:00.123+02:00")),
3374 with("b", json!("yesterday-ish")),
3375 with("c", json!("2026-07-11T00:00:00Z")),
3376 ];
3377 let published: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3378 .into_iter()
3379 .map(|e| e.published)
3380 .collect();
3381 assert_eq!(
3382 published,
3383 vec![
3384 Some("2026-07-11T07:30:00Z".to_string()),
3385 None,
3386 Some("2026-07-11T00:00:00Z".to_string()),
3387 ],
3388 "publishedAt was not normalised to the store's spelling"
3389 );
3390 }
3391
3392 #[test]
3404 fn an_undated_document_is_dated_from_its_tid_rkey() {
3405 let records = vec![publication("p", "https://example.com")];
3406 let (site, pubn) = publication_from_records("p", &records).unwrap();
3407 let docs = vec![rec(
3408 nsid::STANDARD_DOCUMENT,
3409 PAST_TID,
3410 json!({ "title": "T", "path": "/x", "site": site }),
3411 )];
3412 assert_eq!(
3413 entries_from_records(&site, &pubn, &docs)
3414 .into_iter()
3415 .next()
3416 .expect("the document is an entry")
3417 .published,
3418 Some(PAST_TID_WRITTEN_AT.to_string()),
3419 "the date must come from the record key, and be spelled the way the store spells dates"
3420 );
3421 }
3422
3423 #[test]
3424 fn an_unparseable_published_at_falls_back_to_the_tid_rkey() {
3425 let records = vec![publication("p", "https://example.com")];
3426 let (site, pubn) = publication_from_records("p", &records).unwrap();
3427 let docs = vec![rec(
3428 nsid::STANDARD_DOCUMENT,
3429 PAST_TID,
3430 json!({ "title": "T", "publishedAt": "yesterday-ish", "path": "/x", "site": site }),
3431 )];
3432 assert_eq!(
3433 entries_from_records(&site, &pubn, &docs)
3434 .into_iter()
3435 .next()
3436 .expect("the document is an entry")
3437 .published,
3438 Some(PAST_TID_WRITTEN_AT.to_string()),
3439 "a date the parser cannot read is no date at all, so the rkey must stand in"
3440 );
3441 }
3442
3443 #[test]
3444 fn a_stated_date_outranks_the_rkey() {
3445 let records = vec![publication("p", "https://example.com")];
3446 let (site, pubn) = publication_from_records("p", &records).unwrap();
3447 let docs = vec![rec(
3450 nsid::STANDARD_DOCUMENT,
3451 PAST_TID,
3452 json!({ "title": "T", "publishedAt": "2020-01-02T00:00:00Z", "path": "/x", "site": site }),
3453 )];
3454 assert_eq!(
3455 entries_from_records(&site, &pubn, &docs)
3456 .into_iter()
3457 .next()
3458 .expect("the document is an entry")
3459 .published,
3460 Some("2020-01-02T00:00:00Z".to_string()),
3461 "the rkey records when the file was written, which is not when the post was published"
3462 );
3463 }
3464
3465 #[test]
3474 fn a_future_dated_document_falls_back_to_its_rkey() {
3475 let records = vec![publication("p", "https://example.com")];
3476 let (site, pubn) = publication_from_records("p", &records).unwrap();
3477 let docs = vec![rec(
3478 nsid::STANDARD_DOCUMENT,
3479 PAST_TID,
3480 json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
3481 )];
3482 assert_eq!(
3483 entries_from_records(&site, &pubn, &docs)
3484 .into_iter()
3485 .next()
3486 .expect("the document is an entry")
3487 .published,
3488 Some(PAST_TID_WRITTEN_AT.to_string()),
3489 "the date must be the record's write time, not the hour the poll happened to run"
3490 );
3491 }
3492
3493 #[test]
3494 fn a_future_dated_document_without_a_tid_rkey_is_undated() {
3495 let records = vec![publication("p", "https://example.com")];
3496 let (site, pubn) = publication_from_records("p", &records).unwrap();
3497 let docs = vec![rec(
3498 nsid::STANDARD_DOCUMENT,
3499 "self",
3500 json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
3501 )];
3502 assert_eq!(
3503 entries_from_records(&site, &pubn, &docs)
3504 .into_iter()
3505 .next()
3506 .expect("the document is an entry")
3507 .published,
3508 None,
3509 "with nothing credible to date it by, the row falls to fetched_at, which holds still"
3510 );
3511 }
3512
3513 #[test]
3521 fn a_stated_date_a_little_ahead_of_our_clock_is_still_believed() {
3522 let records = vec![publication("p", "https://example.com")];
3523 let (site, pubn) = publication_from_records("p", &records).unwrap();
3524 let slightly_ahead =
3525 crate::feed::fmt_time(chrono::Utc::now() + chrono::Duration::seconds(10));
3526 let docs = vec![rec(
3527 nsid::STANDARD_DOCUMENT,
3528 "self",
3529 json!({ "title": "T", "publishedAt": slightly_ahead, "path": "/x", "site": site }),
3530 )];
3531 assert_eq!(
3532 entries_from_records(&site, &pubn, &docs)
3533 .into_iter()
3534 .next()
3535 .expect("the document is an entry")
3536 .published,
3537 Some(slightly_ahead),
3538 "a few seconds of clock skew must not cost the entry its date"
3539 );
3540 }
3541
3542 #[test]
3543 fn a_document_with_neither_a_date_nor_a_tid_rkey_stays_undated() {
3544 let records = vec![publication("p", "https://example.com")];
3545 let (site, pubn) = publication_from_records("p", &records).unwrap();
3546 let docs = vec![
3553 rec(
3554 nsid::STANDARD_DOCUMENT,
3555 "my-first-post",
3556 json!({ "title": "T", "path": "/x", "site": site }),
3557 ),
3558 rec(
3559 nsid::STANDARD_DOCUMENT,
3560 "abcdefghijklm",
3561 json!({ "title": "T", "path": "/y", "site": site }),
3562 ),
3563 ];
3564 assert_eq!(
3565 entries_from_records(&site, &pubn, &docs)
3566 .into_iter()
3567 .map(|e| e.published)
3568 .collect::<Vec<_>>(),
3569 vec![None, None],
3570 "an invented date is worse than no date; the store decides what to do with undated rows"
3571 );
3572 }
3573
3574 #[test]
3584 fn summaries_are_escaped_as_plain_text_not_sanitised_as_markup() {
3585 let records = vec![publication("p", "https://example.com")];
3586 let (site, pubn) = publication_from_records("p", &records).unwrap();
3587 let doc = |rkey: &str, body: &str| {
3588 rec(
3589 nsid::STANDARD_DOCUMENT,
3590 rkey,
3591 json!({
3592 "title": "T",
3593 "publishedAt": "2026-07-11T00:00:00Z",
3594 "path": "/d",
3595 "site": site,
3596 "textContent": body,
3597 }),
3598 )
3599 };
3600 let summaries: Vec<String> = entries_from_records(
3601 &site,
3602 &pubn,
3603 &[
3604 doc("a", "Vec<String> is a type"),
3605 doc("b", "<script>alert(1)</script>"),
3606 ],
3607 )
3608 .into_iter()
3609 .filter_map(|e| e.summary)
3610 .collect();
3611 assert_eq!(
3612 summaries[0], "Vec<String> is a type",
3613 "prose was eaten by an HTML parser"
3614 );
3615 assert!(
3616 !summaries[1].contains("<script"),
3617 "escaping failed: {}",
3618 summaries[1]
3619 );
3620 }
3621
3622 #[test]
3630 fn an_entry_url_is_never_an_unvetted_path() {
3631 let pubn = Publication {
3632 name: None,
3633 url: "not a url".to_string(),
3636 };
3637 let site = canonical("p");
3638 let docs = vec![
3639 document("a", &site, "Hostile", "javascript:alert(1)"),
3640 document("b", &site, "Fine", "https://example.com/ok"),
3641 ];
3642 let entries = entries_from_records(&site, &pubn, &docs);
3643 assert_eq!(
3644 entries[0].url, None,
3645 "an unvetted path became an entry link"
3646 );
3647 assert_eq!(
3652 entries[1].url, None,
3653 "an off-origin absolute URL was published under the publication's name"
3654 );
3655 }
3656
3657 #[test]
3661 fn a_document_without_published_at_is_still_an_entry() {
3662 let records = vec![publication("p", "https://example.com")];
3663 let (site, pubn) = publication_from_records("p", &records).unwrap();
3664 let doc = rec(
3665 nsid::STANDARD_DOCUMENT,
3666 "d",
3667 json!({ "title": "T", "path": "/d", "site": site }),
3668 );
3669 let entries = entries_from_records(&site, &pubn, &[doc]);
3670 assert_eq!(
3671 entries.len(),
3672 1,
3673 "a missing publishedAt dropped the document"
3674 );
3675 assert_eq!(entries[0].published, None);
3676 }
3677
3678 #[tokio::test]
3684 async fn fetch_refuses_a_uri_for_another_collection() {
3685 let uri = AtUri::parse(&format!("at://{DID}/app.bsky.feed.post/3lab")).unwrap();
3686 let err = fetch(&reqwest::Client::new(), "https://plc.example", &uri)
3687 .await
3688 .expect_err("read a feed post as a publication");
3689 assert!(
3690 format!("{err:#}").contains(nsid::STANDARD_PUBLICATION),
3691 "failed for the wrong reason: {err:#}"
3692 );
3693 }
3694
3695 #[test]
3700 fn a_subpath_publication_keeps_its_base_path() {
3701 let records = vec![publication("p", "https://example.com/blog")];
3702 let (site, pubn) = publication_from_records("p", &records).unwrap();
3703 let docs = vec![document("a", &site, "Relative", "posts/a")];
3704 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3705 .into_iter()
3706 .map(|e| e.url)
3707 .collect();
3708 assert_eq!(urls[0].as_deref(), Some("https://example.com/blog/posts/a"));
3709 }
3710
3711 #[test]
3715 fn no_parseable_base_means_no_url_not_any_url() {
3716 let pubn = Publication {
3717 name: None,
3718 url: "not a url".to_string(),
3719 };
3720 let site = canonical("p");
3721 let docs = vec![document("a", &site, "Absolute", "https://evil.example/x")];
3722 let entries = entries_from_records(&site, &pubn, &docs);
3723 assert_eq!(
3724 entries[0].url, None,
3725 "an off-origin absolute URL was published"
3726 );
3727 }
3728
3729 #[test]
3735 fn an_off_origin_path_yields_no_url_rather_than_the_homepage() {
3736 let records = vec![publication("p", "https://example.com/blog")];
3737 let (site, pubn) = publication_from_records("p", &records).unwrap();
3738 let docs = vec![
3739 document("a", &site, "Elsewhere", "https://www.example.com/post"),
3740 document("b", &site, "Home", "/ok"),
3741 ];
3742 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3743 .into_iter()
3744 .map(|e| e.url)
3745 .collect();
3746 assert_eq!(
3747 urls[0], None,
3748 "an off-origin path was rewritten to the base"
3749 );
3750 assert_eq!(urls[1].as_deref(), Some("https://example.com/ok"));
3751 }
3752
3753 #[test]
3758 fn a_blank_path_yields_no_url() {
3759 let records = vec![publication("p", "https://example.com/blog")];
3760 let (site, pubn) = publication_from_records("p", &records).unwrap();
3761 let docs = vec![
3762 document("a", &site, "Blank", ""),
3763 document("b", &site, "Spaces", " "),
3764 document("c", &site, "Real", "/real"),
3765 ];
3766 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3767 .into_iter()
3768 .map(|e| e.url)
3769 .collect();
3770 assert_eq!(urls[0], None, "a blank path became the homepage");
3771 assert_eq!(urls[1], None, "a whitespace path became the homepage");
3772 assert_eq!(urls[2].as_deref(), Some("https://example.com/real"));
3773 }
3774
3775 #[test]
3778 fn a_document_without_a_path_is_still_an_entry() {
3779 let records = vec![publication("p", "https://example.com")];
3780 let (site, pubn) = publication_from_records("p", &records).unwrap();
3781 let doc = rec(
3782 nsid::STANDARD_DOCUMENT,
3783 "d",
3784 json!({ "title": "T", "publishedAt": "2026-07-11T00:00:00Z", "site": site }),
3785 );
3786 let entries = entries_from_records(&site, &pubn, &[doc]);
3787 assert_eq!(
3788 entries.len(),
3789 1,
3790 "a missing path dropped the whole document"
3791 );
3792 assert_eq!(entries[0].title, "T");
3793 assert_eq!(entries[0].url, None);
3794 }
3795
3796 #[test]
3802 fn documents_are_classified_keep_sibling_or_orphan() {
3803 let pubs = [
3804 publication("a", "https://example.com"),
3805 publication("b", "https://b.example"),
3806 ];
3807 let known: std::collections::HashSet<&str> = pubs.iter().map(|p| p.uri.as_str()).collect();
3808 let mine = canonical("a");
3809 let wanted: std::collections::HashMap<String, usize> = [(mine.clone(), 0)].into();
3810 let fate = |d: &crate::atproto::RecordEntry| classify_document(d, &wanted, &known);
3811
3812 assert_eq!(
3813 fate(&document("d1", &mine, "Mine", "/1")),
3814 DocumentFate::Keep(0)
3815 );
3816 assert_eq!(
3817 fate(&document("d2", &canonical("b"), "B's", "/2")),
3818 DocumentFate::Sibling,
3819 "a sibling publication's document is not an orphan"
3820 );
3821 assert_eq!(
3822 fate(&document(
3823 "d3",
3824 "at://did:plc:other/site.standard.publication/x",
3825 "?",
3826 "/3"
3827 )),
3828 DocumentFate::Orphan
3829 );
3830 assert_eq!(
3831 fate(&rec(
3832 nsid::STANDARD_DOCUMENT,
3833 "d4",
3834 json!({"title": "no rest"})
3835 )),
3836 DocumentFate::Malformed
3837 );
3838 }
3839
3840 #[test]
3842 fn at_uri_parsing_uses_the_shared_prefix() {
3843 let uri = format!(
3844 "{}{DID}/{}/abc",
3845 crate::atproto::AT_URI_PREFIX,
3846 nsid::STANDARD_PUBLICATION
3847 );
3848 assert!(AtUri::parse(&uri).is_some());
3849 }
3850
3851 #[test]
3854 fn the_guid_is_the_record_uri_not_the_path() {
3855 let records = vec![publication("p", "https://example.com")];
3856 let (site, pubn) = publication_from_records("p", &records).unwrap();
3857 let docs = vec![document("rk1", &site, "T", "/moved")];
3858 let entries = entries_from_records(&site, &pubn, &docs);
3859 assert_eq!(
3860 entries[0].guid,
3861 format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
3862 );
3863 }
3864
3865 #[test]
3867 fn a_malformed_document_is_skipped_rather_than_fatal() {
3868 let records = vec![publication("p", "https://example.com")];
3869 let (site, pubn) = publication_from_records("p", &records).unwrap();
3870 let docs = vec![
3871 rec(
3872 nsid::STANDARD_DOCUMENT,
3873 "bad",
3874 json!({ "title": "no rest" }),
3875 ),
3876 document("ok", &site, "Good", "/good"),
3877 ];
3878 let entries = entries_from_records(&site, &pubn, &docs);
3879 assert_eq!(entries.len(), 1);
3880 assert_eq!(entries[0].title, "Good");
3881 }
3882
3883 #[test]
3885 fn a_missing_publication_is_none() {
3886 let records = vec![publication("other", "https://example.com")];
3887 assert!(publication_from_records("p", &records).is_none());
3888 }
3889
3890 #[test]
3891 fn at_uris_parse_in_both_forms_and_reject_malformed_ones() {
3892 let did = AtUri::parse(&format!("at://{DID}/site.standard.publication/abc")).unwrap();
3893 assert_eq!(did.authority, DID);
3894 assert_eq!(did.rkey, "abc");
3895 assert_eq!(
3896 did.to_string(),
3897 format!("at://{DID}/site.standard.publication/abc")
3898 );
3899 assert!(AtUri::parse("at://alice.example.com/site.standard.publication/abc").is_some());
3900 for bad in [
3901 "at://",
3902 "at://only-authority",
3903 "at://authority/collection",
3904 "at://authority/collection/",
3905 "at:///collection/rkey",
3906 "at://authority/collection/rkey/extra",
3907 "https://example.com/feed.xml",
3908 "at:authority/collection/rkey",
3909 ] {
3910 assert!(AtUri::parse(bad).is_none(), "parsed {bad:?}");
3911 }
3912 }
3913}