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}
99
100#[derive(Debug, Clone, PartialEq, Eq)]
103pub struct Publication {
104 pub name: Option<String>,
105 pub url: String,
106}
107
108#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct Entry {
111 pub guid: String,
117 pub title: String,
118 pub published: Option<String>,
127 pub url: Option<String>,
132 pub summary: Option<String>,
136}
137
138impl From<Entry> for crate::store::NewEntry {
139 fn from(e: Entry) -> Self {
142 crate::store::NewEntry {
143 guid: crate::feed::bound_guid(e.guid),
146 url: e.url,
147 title: Some(e.title),
148 author: None,
149 published: e.published,
150 content_html: e.summary,
151 fetched_at: None,
152 keep_stored_content: false,
153 }
154 }
155}
156
157#[derive(Debug, Deserialize)]
158struct PublicationValue {
159 name: Option<String>,
160 url: String,
161}
162
163#[derive(Debug, Deserialize)]
164struct DocumentValue {
165 title: String,
166 #[serde(rename = "publishedAt")]
171 published_at: Option<String>,
172 path: Option<String>,
176 site: String,
182 #[serde(rename = "textContent")]
183 text_content: Option<String>,
184 description: Option<String>,
185}
186
187pub fn publication_from_records(
195 rkey: &str,
196 records: &[crate::atproto::RecordEntry],
197) -> Option<(String, Publication)> {
198 let entry = records
199 .iter()
200 .find(|r| AtUri::parse(&r.uri).is_some_and(|u| u.rkey == rkey))?;
201 let value: PublicationValue = serde_json::from_value(entry.value.clone()).ok()?;
202 let url = crate::net::safe_link(&value.url)?;
207 Some((
208 entry.uri.clone(),
209 Publication {
210 name: value
211 .name
212 .map(|n| crate::feed::bound_text(n, crate::feed::MAX_TITLE_BYTES)),
213 url: crate::feed::bound_text(url, crate::feed::MAX_URL_BYTES),
214 },
215 ))
216}
217
218pub fn entries_from_records(
221 canonical_site: &str,
222 publication: &Publication,
223 records: &[crate::atproto::RecordEntry],
224) -> Vec<Entry> {
225 let base = url::Url::parse(&publication.url).ok().map(|mut u| {
231 if !u.path().ends_with('/') {
232 u.set_path(&format!("{}/", u.path()));
233 }
234 u
235 });
236 let now = chrono::Utc::now();
242 let ceiling = now + chrono::Duration::seconds(crate::atproto::CLOCK_SKEW_GRACE_SECS);
248 records
249 .iter()
250 .filter_map(|record| {
251 let doc: DocumentValue = serde_json::from_value(record.value.clone()).ok()?;
254 if doc.site != canonical_site {
255 return None;
256 }
257 Some(Entry {
258 guid: record.uri.clone(),
259 title: crate::feed::bound_text(doc.title, crate::feed::MAX_TITLE_BYTES),
260 published: doc
276 .published_at
277 .as_deref()
278 .and_then(|raw| chrono::DateTime::parse_from_rfc3339(raw).ok())
279 .map(|d| d.with_timezone(&chrono::Utc))
280 .filter(|d| *d <= ceiling)
281 .or_else(|| {
282 AtUri::parse(&record.uri)
283 .and_then(|uri| crate::atproto::tid_timestamp(&uri.rkey))
284 })
285 .map(crate::feed::fmt_time),
286 url: non_blank(doc.path)
291 .as_deref()
292 .and_then(|path| join_path(base.as_ref(), path))
293 .map(|u| crate::feed::bound_text(u, crate::feed::MAX_URL_BYTES)),
294 summary: non_blank(doc.description)
298 .or_else(|| non_blank(doc.text_content))
299 .map(|raw| {
300 crate::feed::plain_text_to_html_bounded(
301 &raw,
302 crate::feed::MAX_CONTENT_HTML_BYTES,
303 )
304 }),
305 })
306 })
307 .collect()
308}
309
310fn ingest_floor(
358 retention_days: u32,
359 retention_hard_days: u32,
360 now: chrono::DateTime<chrono::Utc>,
361) -> Option<String> {
362 let days = if retention_days > 0 {
363 retention_days
364 } else if retention_hard_days > 0 {
365 retention_hard_days
366 } else {
367 return None;
368 };
369 let window = chrono::Duration::try_days(days.into())?;
370 now.checked_sub_signed(window).map(crate::feed::fmt_time)
371}
372
373pub async fn store_publication(
408 pool: &sqlx::SqlitePool,
409 url: &str,
410 read: PublicationRead,
411 max_entries_per_feed: i64,
412 retention_days: u32,
413 retention_hard_days: u32,
414) -> anyhow::Result<crate::feed::PollOutcome> {
415 let offered = read.entries.len();
416
417 if !read.complete && offered == 0 {
428 return Ok(crate::feed::PollOutcome::Failed {
429 backoff: crate::feed::backoff_for(1),
430 kind: crate::feed::FailureKind::Body,
431 detail: crate::feed::failure_detail(
432 "the publication read stopped before its first document",
433 ),
434 });
435 }
436
437 let floor = ingest_floor(retention_days, retention_hard_days, chrono::Utc::now());
439 let rows: Vec<crate::store::NewEntry> = read
440 .entries
441 .into_iter()
442 .filter(|e| match (&floor, &e.published) {
443 (Some(floor), Some(published)) => published.as_str() >= floor.as_str(),
446 _ => true,
450 })
451 .map(Into::into)
452 .collect();
453
454 if offered > rows.len() {
460 tracing::info!(
461 feed = %url,
462 offered,
463 stored = rows.len(),
464 "the retention floor dropped documents older than the window"
465 );
466 }
467
468 let feed_id = crate::store::upsert_feed(
469 pool,
470 &crate::store::NewFeed {
471 url: url.to_string(),
472 title: read.publication.name.clone(),
473 site_url: Some(read.publication.url.clone()),
474 last_polled: Some(crate::feed::fmt_time(chrono::Utc::now())),
475 ..Default::default()
476 },
477 )
478 .await?;
479
480 let new_entries =
481 crate::store::insert_entries(pool, feed_id, &rows, max_entries_per_feed).await?;
482 Ok(crate::feed::PollOutcome::Updated { new_entries })
483}
484
485fn non_blank(s: Option<String>) -> Option<String> {
486 s.filter(|v| !v.trim().is_empty())
487}
488
489fn join_path(base: Option<&url::Url>, path: &str) -> Option<String> {
497 let base = base?;
511 match base.join(path) {
512 Ok(joined) if joined.origin() == base.origin() => crate::net::safe_link(joined.as_str()),
515 _ => None,
521 }
522}
523
524#[derive(Debug, Clone, Copy, PartialEq, Eq)]
529enum DocumentFate {
530 Keep(usize),
532 Sibling,
535 Orphan,
539 Malformed,
541}
542
543fn classify_document(
544 record: &crate::atproto::RecordEntry,
545 wanted: &std::collections::HashMap<String, usize>,
546 known: &std::collections::HashSet<&str>,
547) -> DocumentFate {
548 match serde_json::from_value::<DocumentValue>(record.value.clone()) {
549 Ok(doc) if wanted.contains_key(&doc.site) => DocumentFate::Keep(wanted[&doc.site]),
550 Ok(doc) if known.contains(doc.site.as_str()) => DocumentFate::Sibling,
551 Ok(_) => DocumentFate::Orphan,
552 Err(_) => DocumentFate::Malformed,
553 }
554}
555
556pub async fn fetch(
577 http: &reqwest::Client,
578 plc_directory: &str,
579 uri: &AtUri,
580) -> anyhow::Result<PublicationRead> {
581 if uri.collection != nsid::STANDARD_PUBLICATION {
586 return Err(
587 NotAPublication(format!("{uri} is not a {} URI", nsid::STANDARD_PUBLICATION)).into(),
588 );
589 }
590 fetch_repo(
591 http,
592 plc_directory,
593 &uri.authority,
594 std::slice::from_ref(&uri.rkey),
595 )
596 .await?
597 .pop()
598 .unwrap_or_else(|| Err(NotAPublication(format!("{uri} was not read")).into()))
599}
600
601#[derive(Debug)]
609pub struct NotAPublication(pub String);
610
611impl std::fmt::Display for NotAPublication {
612 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
613 f.write_str(&self.0)
614 }
615}
616
617impl std::error::Error for NotAPublication {}
618
619pub async fn fetch_repo(
621 http: &reqwest::Client,
622 plc_directory: &str,
623 did: &str,
624 rkeys: &[String],
625) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
626 fetch_repo_capped(
627 http,
628 plc_directory,
629 did,
630 rkeys,
631 crate::atproto::MAX_LARGE_RECORDS,
632 crate::atproto::MAX_LIST_BYTES,
633 )
634 .await
635}
636
637pub(crate) async fn fetch_repo_capped(
640 http: &reqwest::Client,
641 plc_directory: &str,
642 did: &str,
643 requested: &[String],
644 per_site_cap: usize,
645 budget_bytes: usize,
646) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
647 use anyhow::Context;
648
649 if requested.is_empty() {
652 return Ok(Vec::new());
653 }
654 let mut rkeys: Vec<String> = Vec::new();
657 for r in requested {
658 if !rkeys.contains(r) {
659 rkeys.push(r.clone());
660 }
661 }
662 let rkeys = rkeys.as_slice();
663
664 let pds = crate::atproto::resolve_did_to_pds(http, plc_directory, did)
665 .await
666 .with_context(|| format!("resolving the PDS for {did}"))?;
667 let client = crate::atproto::PdsClient::anonymous(http.clone(), pds, did.to_string());
668
669 let mut budget = crate::atproto::ByteBudget::new(budget_bytes);
673 let (publications, skipped_publications) = client
674 .list_all_records_skipping_within(nsid::STANDARD_PUBLICATION, &mut budget)
675 .await
676 .with_context(|| format!("listing publications for {did}"))?;
677 if skipped_publications > 0 {
678 tracing::warn!(
679 repo = %did,
680 skipped = skipped_publications,
681 "skipped malformed publication records in this repo"
682 );
683 }
684
685 let wanted: Vec<Option<(String, Publication)>> = rkeys
687 .iter()
688 .map(|rkey| publication_from_records(rkey, &publications))
689 .collect();
690 let index_of: std::collections::HashMap<String, usize> = wanted
691 .iter()
692 .enumerate()
693 .filter_map(|(i, w)| w.as_ref().map(|(site, _)| (site.clone(), i)))
694 .collect();
695 let not_found = |rkey: &str| {
696 anyhow::Error::new(NotAPublication(format!(
697 "at://{did}/{}/{rkey} is not a readable site.standard.publication",
698 nsid::STANDARD_PUBLICATION
699 )))
700 };
701 if index_of.is_empty() {
702 return Ok(requested.iter().map(|r| Err(not_found(r))).collect());
703 }
704
705 let known: std::collections::HashSet<&str> =
711 publications.iter().map(|p| p.uri.as_str()).collect();
712 let mut kept_per = vec![0usize; rkeys.len()];
713 let mut capped = vec![false; rkeys.len()];
714 let mut orphaned = 0usize;
715 let documents = client
716 .list_recent_matching_within(
717 nsid::STANDARD_DOCUMENT,
718 per_site_cap.saturating_mul(index_of.len()),
719 &mut budget,
720 DOCUMENT_PAGE_SIZE,
721 |record| match classify_document(record, &index_of, &known) {
722 DocumentFate::Keep(i) if kept_per[i] < per_site_cap => {
723 kept_per[i] += 1;
724 true
725 }
726 DocumentFate::Keep(i) => {
727 capped[i] = true;
728 false
729 }
730 DocumentFate::Orphan => {
734 orphaned += 1;
735 false
736 }
737 DocumentFate::Sibling | DocumentFate::Malformed => false,
738 },
739 )
740 .await
741 .with_context(|| format!("listing documents for {did}"))?;
742
743 let mut per_site: Vec<Vec<crate::atproto::RecordEntry>> = vec![Vec::new(); rkeys.len()];
745 for record in documents.records {
746 let site = serde_json::from_value::<DocumentValue>(record.value.clone()).map(|d| d.site);
747 if let Some(&i) = site.ok().as_deref().and_then(|s| index_of.get(s)) {
748 per_site[i].push(record);
749 }
750 }
751 let reads: Vec<anyhow::Result<PublicationRead>> = wanted
752 .into_iter()
753 .zip(per_site)
754 .enumerate()
755 .map(|(i, (w, records))| match w {
756 None => Err(not_found(&rkeys[i])),
757 Some((site, publication)) => {
758 let walk = crate::atproto::RecordWalk {
759 records,
760 complete: documents.complete && !capped[i],
761 malformed: documents.malformed,
762 };
763 Ok(read_from(publication, &site, walk, orphaned))
764 }
765 })
766 .collect();
767 let mut reads: Vec<Option<anyhow::Result<PublicationRead>>> =
770 reads.into_iter().map(Some).collect();
771 let positions: Vec<usize> = requested
772 .iter()
773 .map(|r| {
774 rkeys
775 .iter()
776 .position(|k| k == r)
777 .expect("every requested rkey is in rkeys")
778 })
779 .collect();
780 Ok(positions
781 .iter()
782 .enumerate()
783 .map(|(n, &u)| {
784 let repeated_later = positions[n + 1..].contains(&u);
785 match (&reads[u], repeated_later) {
786 (Some(Ok(read)), true) => Ok(read.clone()),
787 (Some(Err(_)), _) | (None, _) => Err(not_found(&requested[n])),
788 (Some(Ok(_)), false) => match reads[u].take() {
789 Some(Ok(read)) => Ok(read),
790 _ => Err(not_found(&requested[n])),
791 },
792 }
793 })
794 .collect())
795}
796
797fn read_from(
804 publication: Publication,
805 canonical_site: &str,
806 documents: crate::atproto::RecordWalk,
807 orphaned: usize,
808) -> PublicationRead {
809 let entries = entries_from_records(canonical_site, &publication, &documents.records);
810
811 if !documents.complete {
821 tracing::warn!(
822 site = %canonical_site,
823 kept = entries.len(),
824 "stopped reading this publication before its documents ran out"
825 );
826 }
827 if documents.malformed > 0 {
836 tracing::warn!(
837 site = %canonical_site,
838 skipped = documents.malformed,
839 "skipped malformed document records for this publication"
840 );
841 }
842 if orphaned > 0 {
843 tracing::warn!(
844 site = %canonical_site,
845 orphaned,
846 "documents in this repo reference no publication in it — a `site` spelling nothing matches"
847 );
848 }
849 PublicationRead {
850 publication,
851 entries,
852 complete: documents.complete,
853 }
854}
855
856#[cfg(test)]
857pub(crate) mod tests {
858 use super::*;
859 use crate::atproto::RecordEntry;
860 use serde_json::json;
861
862 const DID: &str = "did:plc:ohutz6x5acjmpuulp3x7wxxc";
863
864 #[test]
870 fn a_read_carries_whether_the_walk_finished() {
871 let publication = Publication {
872 name: Some("Scan's Lab".to_string()),
873 url: "https://example.com/blog/".to_string(),
874 };
875 for complete in [true, false] {
876 let walk = crate::atproto::RecordWalk {
877 records: Vec::new(),
878 complete,
879 malformed: 0,
880 };
881 let read = read_from(publication.clone(), "at://d/c/r", walk, 0);
882 assert_eq!(
883 read.complete, complete,
884 "the walk said complete={complete} and the read said {}",
885 read.complete
886 );
887 }
888 }
889
890 fn read_of(entries: Vec<Entry>, complete: bool) -> PublicationRead {
893 PublicationRead {
894 publication: Publication {
895 name: Some("Scan's Lab".to_string()),
896 url: "https://example.com/blog/".to_string(),
897 },
898 entries,
899 complete,
900 }
901 }
902
903 fn entry_dated(guid: &str, published: Option<&str>) -> Entry {
904 Entry {
905 guid: guid.to_string(),
906 title: "T".to_string(),
907 published: published.map(str::to_string),
908 url: Some("https://example.com/blog/a".to_string()),
909 summary: None,
910 }
911 }
912
913 fn days_ago(n: i64) -> String {
914 crate::feed::fmt_time(chrono::Utc::now() - chrono::Duration::days(n))
915 }
916
917 const PUB_URL: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
918
919 async fn pool() -> sqlx::SqlitePool {
920 crate::store::init_url("sqlite::memory:").await.unwrap()
921 }
922
923 #[tokio::test]
924 async fn a_complete_read_stores_its_entries_and_reports_them() {
925 let pool = pool().await;
926 let read = read_of(
927 vec![
928 entry_dated("at://d/c/1", Some(&days_ago(1))),
929 entry_dated("at://d/c/2", Some(&days_ago(2))),
930 ],
931 true,
932 );
933 let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
934 .await
935 .unwrap();
936 assert!(
937 matches!(
938 outcome,
939 crate::feed::PollOutcome::Updated { new_entries: 2 }
940 ),
941 "expected two new entries, got {outcome:?}"
942 );
943 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
944 .fetch_one(&pool)
945 .await
946 .unwrap();
947 assert_eq!(n, 2, "the entries were not stored");
948 }
949
950 #[tokio::test]
958 async fn an_incomplete_read_that_offered_nothing_is_a_failure() {
959 let pool = pool().await;
960 let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
961 .await
962 .unwrap();
963 assert!(
964 matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
965 "a truncated read that produced nothing is not a healthy poll: {outcome:?}"
966 );
967 }
968
969 #[tokio::test]
970 async fn an_incomplete_read_whose_entries_the_floor_dropped_is_not_a_failure() {
971 let pool = pool().await;
972 let read = read_of(
973 vec![entry_dated("at://d/c/old", Some(&days_ago(900)))],
974 false,
975 );
976 let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
977 .await
978 .unwrap();
979 assert!(
980 !matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
981 "the read offered an entry; the floor dropping it is not a failed poll: {outcome:?}"
982 );
983 }
984
985 #[tokio::test]
987 async fn a_failed_read_does_not_stamp_last_polled() {
988 let pool = pool().await;
989 let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
990 .await
991 .unwrap();
992 assert!(matches!(outcome, crate::feed::PollOutcome::Failed { .. }));
993 let stamped: Option<String> =
994 sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
995 .bind(PUB_URL)
996 .fetch_optional(&pool)
997 .await
998 .unwrap()
999 .flatten();
1000 assert_eq!(
1001 stamped, None,
1002 "a failed poll stamped last_polled, so the feed reads as freshly polled"
1003 );
1004 }
1005
1006 #[tokio::test]
1007 async fn an_entry_already_older_than_the_window_is_not_stored() {
1008 let pool = pool().await;
1009 let read = read_of(
1010 vec![
1011 entry_dated("at://d/c/fresh", Some(&days_ago(1))),
1012 entry_dated("at://d/c/ancient", Some(&days_ago(900))),
1013 ],
1014 true,
1015 );
1016 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1017 .await
1018 .unwrap();
1019 let guids: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
1020 .fetch_all(&pool)
1021 .await
1022 .unwrap();
1023 assert_eq!(
1024 guids,
1025 vec!["at://d/c/fresh".to_string()],
1026 "an entry the next sweep would delete was stored anyway"
1027 );
1028 }
1029
1030 #[tokio::test]
1038 async fn the_floor_follows_the_hard_ceiling_when_the_window_is_disabled() {
1039 let pool = pool().await;
1040 let read = read_of(
1041 vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
1042 true,
1043 );
1044 store_publication(&pool, PUB_URL, read, 0, 0, 180)
1045 .await
1046 .unwrap();
1047 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1048 .fetch_one(&pool)
1049 .await
1050 .unwrap();
1051 assert_eq!(
1052 n, 0,
1053 "an entry the hard ceiling will delete was stored, so it will resurrect"
1054 );
1055 }
1056
1057 #[tokio::test]
1064 async fn the_floor_follows_the_window_when_both_are_set() {
1065 let pool = pool().await;
1066 let read = read_of(
1067 vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
1068 true,
1069 );
1070 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1071 .await
1072 .unwrap();
1073 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1074 .fetch_one(&pool)
1075 .await
1076 .unwrap();
1077 assert_eq!(
1078 n, 0,
1079 "an entry inside the ceiling but outside the window was stored, so it will cycle"
1080 );
1081 }
1082
1083 #[test]
1097 fn a_publications_floor_is_its_archive_ceiling_not_the_rss_window() {
1098 let config = crate::config::Config::default();
1099 let now = chrono::DateTime::parse_from_rfc3339("2026-09-29T12:00:00Z")
1100 .unwrap()
1101 .with_timezone(&chrono::Utc);
1102
1103 let (days, hard) = config.retention_for(crate::feed::FeedKind::Publication);
1104 let floor = ingest_floor(days, hard, now).expect("a publication has an ingest floor");
1105 let expected = crate::feed::fmt_time(
1106 now - chrono::Duration::days(config.publication_retention_days.into()),
1107 );
1108 assert_eq!(
1109 floor, expected,
1110 "a publication's floor must be its archive ceiling ({} days), because \
1111 that is the only sweep pass that can delete its rows",
1112 config.publication_retention_days,
1113 );
1114
1115 let rss_floor = ingest_floor(config.retention_days, config.retention_hard_days, now)
1118 .expect("an RSS feed has an ingest floor");
1119 assert!(
1120 floor < rss_floor,
1121 "the publication floor ({floor}) is no older than the RSS one \
1122 ({rss_floor}), so a months-old document would still be dropped",
1123 );
1124 let a_real_publications_newest_document =
1125 crate::feed::fmt_time(now - chrono::Duration::days(109));
1126 assert!(
1127 a_real_publications_newest_document.as_str() >= floor.as_str(),
1128 "the newest document a real publication offered would be refused at \
1129 ingest: {a_real_publications_newest_document} against a floor of {floor}",
1130 );
1131 assert!(
1132 a_real_publications_newest_document.as_str() < rss_floor.as_str(),
1133 "this assertion is only meaningful while the RSS window WOULD have \
1134 dropped it, and it no longer does",
1135 );
1136 }
1137
1138 #[test]
1145 fn the_ingest_floor_is_the_window_the_sweep_would_use() {
1146 let now = chrono::DateTime::parse_from_rfc3339("2026-09-26T12:00:00Z")
1147 .unwrap()
1148 .with_timezone(&chrono::Utc);
1149
1150 assert_eq!(
1151 ingest_floor(14, 180, now).as_deref(),
1152 Some("2026-09-12T12:00:00Z"),
1153 "with both set, the floor is the WINDOW — the thing that deletes first",
1154 );
1155 assert_eq!(
1156 ingest_floor(180, 30, now).as_deref(),
1157 Some("2026-03-30T12:00:00Z"),
1158 "a ceiling INSIDE the window is one `prune_old_entries` ignores, so it \
1159 must not lower the floor — `min` here discarded five months of archive \
1160 that nothing would have deleted",
1161 );
1162 assert_eq!(
1163 ingest_floor(0, 30, now).as_deref(),
1164 Some("2026-08-27T12:00:00Z"),
1165 "with no window the ceiling stands alone, and it still deletes",
1166 );
1167 assert_eq!(
1168 ingest_floor(0, 0, now),
1169 None,
1170 "with no retention at all there is nothing to floor against",
1171 );
1172 assert_eq!(
1177 ingest_floor(u32::MAX, u32::MAX, now),
1178 None,
1179 "an unrepresentable window produced a floor, so the comparison is \
1180 resting on how `fmt_time` renders an out-of-range year",
1181 );
1182 }
1183
1184 #[tokio::test]
1199 async fn a_ceiling_inside_the_window_does_not_lower_the_floor() {
1200 let pool = pool().await;
1201 let read = read_of(
1202 vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
1203 true,
1204 );
1205 store_publication(&pool, PUB_URL, read, 0, 180, 30)
1206 .await
1207 .unwrap();
1208 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1209 .fetch_one(&pool)
1210 .await
1211 .unwrap();
1212 assert_eq!(
1213 n, 1,
1214 "an entry inside the 180-day window was dropped because of a 30-day \
1215 ceiling the sweep ignores — five months of archive discarded at ingest \
1216 that nothing would have deleted",
1217 );
1218 }
1219
1220 #[tokio::test]
1225 async fn a_complete_read_of_nothing_is_not_a_failure() {
1226 let pool = pool().await;
1227 let outcome = store_publication(&pool, PUB_URL, read_of(vec![], true), 0, 14, 180)
1228 .await
1229 .unwrap();
1230 assert!(
1231 matches!(
1232 outcome,
1233 crate::feed::PollOutcome::Updated { new_entries: 0 }
1234 ),
1235 "a complete read of an empty publication was not a healthy poll: {outcome:?}"
1236 );
1237 }
1238
1239 #[tokio::test]
1243 async fn a_successful_read_stamps_last_polled() {
1244 let pool = pool().await;
1245 let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
1246 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1247 .await
1248 .unwrap();
1249 let stamped: Option<String> =
1250 sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
1251 .bind(PUB_URL)
1252 .fetch_one(&pool)
1253 .await
1254 .unwrap();
1255 assert!(
1256 stamped.is_some(),
1257 "a successful read left `last_polled` NULL, so the publication stays \
1258 due forever and `/stats` never shows it as polled",
1259 );
1260 }
1261
1262 #[tokio::test]
1269 async fn an_absurd_retention_window_keeps_everything_rather_than_nothing() {
1270 let pool = pool().await;
1271 let read = read_of(
1272 vec![entry_dated("at://d/c/ancient", Some(&days_ago(10_000)))],
1273 true,
1274 );
1275 store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
1276 .await
1277 .expect("a huge window is a wide floor, not a crash");
1278 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1279 .fetch_one(&pool)
1280 .await
1281 .unwrap();
1282 assert_eq!(
1283 n, 1,
1284 "a window of u32::MAX days dropped a 27-year-old entry, so the \
1285 saturation went the wrong way",
1286 );
1287 }
1288
1289 #[tokio::test]
1294 async fn the_per_feed_cap_is_the_one_the_caller_passed() {
1295 let pool = pool().await;
1296 let read = read_of(
1297 (0..5)
1298 .map(|i| entry_dated(&format!("at://d/c/{i}"), Some(&days_ago(i + 1))))
1299 .collect(),
1300 true,
1301 );
1302 store_publication(&pool, PUB_URL, read, 2, 0, 0)
1303 .await
1304 .unwrap();
1305 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1306 .fetch_one(&pool)
1307 .await
1308 .unwrap();
1309 assert_eq!(
1310 n, 2,
1311 "five entries under a cap of two left {n} rows, so the caller's cap is \
1312 not the one being applied",
1313 );
1314 }
1315
1316 #[tokio::test]
1317 async fn no_retention_at_all_means_no_ingest_floor() {
1318 let pool = pool().await;
1319 let read = read_of(
1320 vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
1321 true,
1322 );
1323 store_publication(&pool, PUB_URL, read, 0, 0, 0)
1324 .await
1325 .unwrap();
1326 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1327 .fetch_one(&pool)
1328 .await
1329 .unwrap();
1330 assert_eq!(
1331 n, 1,
1332 "with nothing deleting it, an old entry is worth keeping"
1333 );
1334 }
1335
1336 #[tokio::test]
1338 async fn an_absurd_retention_window_does_not_panic() {
1339 let pool = pool().await;
1340 let read = read_of(vec![entry_dated("at://d/c/x", Some(&days_ago(1)))], true);
1341 store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
1342 .await
1343 .expect("a huge window is a wide floor, not a crash");
1344 }
1345
1346 #[tokio::test]
1350 async fn the_feed_row_learns_the_publications_name_and_site() {
1351 let pool = pool().await;
1352 let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
1353 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1354 .await
1355 .unwrap();
1356 let (title, site): (Option<String>, Option<String>) =
1357 sqlx::query_as("SELECT title, site_url FROM feeds WHERE url = ?1")
1358 .bind(PUB_URL)
1359 .fetch_one(&pool)
1360 .await
1361 .unwrap();
1362 assert_eq!(
1363 title.as_deref(),
1364 Some("Scan's Lab"),
1365 "the name never reached the row"
1366 );
1367 assert_eq!(
1368 site.as_deref(),
1369 Some("https://example.com/blog/"),
1370 "the homepage never reached the row"
1371 );
1372 }
1373
1374 #[tokio::test]
1378 async fn an_undated_entry_is_stored_rather_than_dropped() {
1379 let pool = pool().await;
1380 let read = read_of(vec![entry_dated("at://d/c/undated", None)], true);
1381 store_publication(&pool, PUB_URL, read, 0, 14, 180)
1382 .await
1383 .unwrap();
1384 let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1385 .fetch_one(&pool)
1386 .await
1387 .unwrap();
1388 assert_eq!(n, 1, "an entry with no date was discarded");
1389 }
1390
1391 const PAST_TID: &str = "3jzfcijpj2z2a";
1399 const PAST_TID_WRITTEN_AT: &str = "2023-06-30T15:03:01Z";
1400
1401 fn rec(collection: &str, rkey: &str, value: serde_json::Value) -> RecordEntry {
1402 RecordEntry {
1403 uri: format!("at://{DID}/{collection}/{rkey}"),
1404 cid: None,
1405 value,
1406 }
1407 }
1408
1409 fn publication(rkey: &str, url: &str) -> RecordEntry {
1410 rec(
1411 nsid::STANDARD_PUBLICATION,
1412 rkey,
1413 json!({ "name": "Scan's Lab", "url": url }),
1414 )
1415 }
1416
1417 fn document(rkey: &str, site: &str, title: &str, path: &str) -> RecordEntry {
1418 rec(
1419 nsid::STANDARD_DOCUMENT,
1420 rkey,
1421 json!({
1422 "title": title,
1423 "publishedAt": "2026-07-11T00:00:00Z",
1424 "path": path,
1425 "site": site,
1426 "textContent": "body",
1427 }),
1428 )
1429 }
1430
1431 pub(crate) async fn serve_repo(
1438 did: &'static str,
1439 records: Vec<(&'static str, &'static str, serde_json::Value)>,
1440 ) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
1441 use axum::extract::{Query, Request};
1442 use std::collections::HashMap;
1443 use std::sync::atomic::{AtomicUsize, Ordering};
1444 use std::sync::Arc;
1445
1446 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1447 let addr = listener.local_addr().unwrap();
1448 let port = addr.port();
1449 let (plc_host, pds_host) = (
1452 format!("plc-{port}.repo.test"),
1453 format!("pds-{port}.repo.test"),
1454 );
1455 for host in [&plc_host, &pds_host] {
1456 crate::net::test_host_override(host, addr);
1457 }
1458 let hits = Arc::new(AtomicUsize::new(0));
1459 let counter = Arc::clone(&hits);
1460 let records = Arc::new(records);
1461 let app = axum::Router::new().fallback(
1462 move |Query(q): Query<HashMap<String, String>>, req: Request| {
1463 let records = Arc::clone(&records);
1464 let counter = Arc::clone(&counter);
1465 async move {
1466 let path = req.uri().path().to_string();
1467 if path == format!("/{did}") {
1468 return axum::Json(json!({
1469 "id": did,
1470 "service": [{
1471 "id": "#atproto_pds",
1472 "type": "AtprotoPersonalDataServer",
1473 "serviceEndpoint": format!("http://{pds_host}:{port}"),
1474 }],
1475 }));
1476 }
1477 assert_eq!(path, "/xrpc/com.atproto.repo.listRecords", "unexpected request");
1478 counter.fetch_add(1, Ordering::SeqCst);
1479 let collection = q.get("collection").cloned().unwrap_or_default();
1480 let limit: usize = q.get("limit").and_then(|l| l.parse().ok()).unwrap_or(50);
1483 let offset: usize = q.get("cursor").and_then(|c| c.parse().ok()).unwrap_or(0);
1484 let page: Vec<_> = records
1485 .iter()
1486 .filter(|(c, _, _)| *c == collection)
1487 .skip(offset)
1488 .take(limit)
1489 .map(|(c, rkey, value)| {
1490 if rkey.is_empty() {
1493 return json!({ "cid": "bafy", "value": value });
1494 }
1495 json!({ "uri": format!("at://{did}/{c}/{rkey}"), "cid": "bafy", "value": value })
1496 })
1497 .collect();
1498 let next = (offset + page.len()).to_string();
1499 axum::Json(json!({ "records": page, "cursor": next }))
1500 }
1501 },
1502 );
1503 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1504 (format!("http://{plc_host}:{port}"), hits)
1505 }
1506
1507 #[tokio::test]
1512 async fn fetch_reads_a_publication_from_a_mocked_repo() {
1513 const OWN: &str = "did:plc:fetchmock";
1514 let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
1515 let sibling = format!("at://{OWN}/{}/other", nsid::STANDARD_PUBLICATION);
1516 let doc = |site: &str, title: &str| {
1517 json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
1518 "path": "/p", "site": site })
1519 };
1520 let (plc, hits) = serve_repo(
1521 OWN,
1522 vec![
1523 (
1524 nsid::STANDARD_PUBLICATION,
1525 "mine",
1526 json!({ "name": "Mine", "url": "https://mine.example" }),
1527 ),
1528 (
1529 nsid::STANDARD_PUBLICATION,
1530 "other",
1531 json!({ "name": "Other", "url": "https://other.example" }),
1532 ),
1533 (nsid::STANDARD_DOCUMENT, "3l2fmaaaaaa2a", doc(&site, "kept")),
1534 (
1535 nsid::STANDARD_DOCUMENT,
1536 "3l2fmaaaaaa2b",
1537 doc(&sibling, "sibling's"),
1538 ),
1539 ],
1540 )
1541 .await;
1542 let client = crate::feed::build_client().unwrap();
1543 let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
1544 .await
1545 .unwrap();
1546 assert!(read.complete, "the walk did not finish");
1547 assert_eq!(read.publication.name.as_deref(), Some("Mine"));
1548 let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
1549 assert_eq!(
1550 titles,
1551 vec!["kept"],
1552 "the site filter let a sibling through"
1553 );
1554 assert_eq!(
1555 hits.load(std::sync::atomic::Ordering::SeqCst),
1556 4,
1557 "two collections, each a full page then an empty one"
1558 );
1559 }
1560
1561 #[tokio::test]
1565 async fn a_failed_publication_read_is_a_poll_failure() {
1566 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1567 let url = format!(
1570 "at://did:plc:unreachableaaaaaaaaaaaaa/{}/x",
1571 nsid::STANDARD_PUBLICATION
1572 );
1573 crate::store::upsert_feed(
1574 &pool,
1575 &crate::store::NewFeed {
1576 url: url.clone(),
1577 ..Default::default()
1578 },
1579 )
1580 .await
1581 .unwrap();
1582 let feed = crate::store::get_feed_by_url(&pool, &url)
1583 .await
1584 .unwrap()
1585 .unwrap();
1586 let mut config = crate::config::Config::default();
1587 config.oauth.plc_directory = "http://plc.nowhere.invalid".into();
1588 let client = crate::feed::build_client().unwrap();
1589 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1590 .await
1591 .expect("a source failure surfaced as a store error");
1592 assert!(
1593 matches!(
1594 outcome,
1595 crate::feed::PollOutcome::Failed {
1596 kind: crate::feed::FailureKind::Fetch,
1597 ..
1598 }
1599 ),
1600 "an unreachable publication was not a fetch failure: {outcome:?}"
1601 );
1602 }
1603
1604 #[tokio::test]
1608 async fn a_deleted_publication_record_is_not_an_unreachable_publisher() {
1609 const GONE: &str = "did:plc:goneaaaaaaaaaaaaaaaaaaaa";
1610 let site = format!("at://{GONE}/{}/deleted", nsid::STANDARD_PUBLICATION);
1611 let (plc, _) = serve_repo(
1613 GONE,
1614 vec![(
1615 nsid::STANDARD_PUBLICATION,
1616 "another",
1617 json!({ "name": "Other", "url": "https://o.example" }),
1618 )],
1619 )
1620 .await;
1621 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1622 crate::store::upsert_feed(
1623 &pool,
1624 &crate::store::NewFeed {
1625 url: site.clone(),
1626 ..Default::default()
1627 },
1628 )
1629 .await
1630 .unwrap();
1631 let feed = crate::store::get_feed_by_url(&pool, &site)
1632 .await
1633 .unwrap()
1634 .unwrap();
1635 let mut config = crate::config::Config::default();
1636 config.oauth.plc_directory = plc;
1637 let client = crate::feed::build_client().unwrap();
1638 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1639 .await
1640 .unwrap();
1641 assert!(
1642 matches!(
1643 outcome,
1644 crate::feed::PollOutcome::Failed {
1645 kind: crate::feed::FailureKind::Parse,
1646 ..
1647 }
1648 ),
1649 "a deleted publication was filed as a network failure: {outcome:?}"
1650 );
1651 }
1652
1653 async fn serve_answering(
1658 did: &'static str,
1659 plc_status: u16,
1660 pds: (u16, &'static str),
1661 delay: std::time::Duration,
1662 endless: bool,
1663 ) -> String {
1664 use std::sync::atomic::{AtomicUsize, Ordering};
1665 use std::sync::Arc;
1666 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1667 let addr = listener.local_addr().unwrap();
1668 let port = addr.port();
1669 let (plc_host, pds_host) = (
1670 format!("plc-{port}.answer.test"),
1671 format!("pds-{port}.answer.test"),
1672 );
1673 crate::net::test_host_override(&plc_host, addr);
1674 crate::net::test_host_override(&pds_host, addr);
1675 let pages = Arc::new(AtomicUsize::new(0));
1676 let endpoint = format!("http://{pds_host}:{port}");
1677 let app = axum::Router::new().fallback(move |req: axum::extract::Request| {
1678 let pages = Arc::clone(&pages);
1679 let endpoint = endpoint.clone();
1680 async move {
1681 use axum::response::IntoResponse;
1682 if req.uri().path() == format!("/{did}") {
1683 let status = axum::http::StatusCode::from_u16(plc_status).unwrap();
1684 let doc = json!({ "id": did, "service": [{ "id": "#atproto_pds",
1685 "type": "AtprotoPersonalDataServer", "serviceEndpoint": endpoint }] });
1686 return (status, axum::Json(doc)).into_response();
1687 }
1688 tokio::time::sleep(delay).await;
1689 if endless {
1690 let n = pages.fetch_add(1, Ordering::SeqCst);
1691 let body = json!({ "records": [{
1692 "uri": format!("at://{did}/{}/3lend{n:08}", nsid::STANDARD_DOCUMENT),
1693 "cid": "b",
1694 "value": { "title": "x", "path": "/x", "publishedAt": "2026-07-11T00:00:00Z",
1695 "site": format!("at://{did}/{}/other", nsid::STANDARD_PUBLICATION) } }],
1696 "cursor": format!("c{n}") });
1697 return axum::Json(body).into_response();
1698 }
1699 let status = axum::http::StatusCode::from_u16(pds.0).unwrap();
1700 (status, [("content-type", "application/json")], pds.1).into_response()
1701 }
1702 });
1703 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1704 format!("http://{plc_host}:{port}")
1705 }
1706
1707 async fn poll_publication_at(
1708 did: &str,
1709 plc: String,
1710 deadline: Option<std::time::Duration>,
1711 ) -> crate::feed::PollOutcome {
1712 let site = format!("at://{did}/{}/mine", nsid::STANDARD_PUBLICATION);
1713 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1714 crate::store::upsert_feed(
1715 &pool,
1716 &crate::store::NewFeed {
1717 url: site.clone(),
1718 ..Default::default()
1719 },
1720 )
1721 .await
1722 .unwrap();
1723 let feed = crate::store::get_feed_by_url(&pool, &site)
1724 .await
1725 .unwrap()
1726 .unwrap();
1727 let mut config = crate::config::Config::default();
1728 config.oauth.plc_directory = plc;
1729 if let Some(d) = deadline {
1730 config.publication_read_deadline = d;
1731 }
1732 let client = crate::feed::build_client().unwrap();
1733 crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1734 .await
1735 .unwrap()
1736 }
1737
1738 fn kind_of(outcome: &crate::feed::PollOutcome) -> Option<crate::feed::FailureKind> {
1739 match outcome {
1740 crate::feed::PollOutcome::Failed { kind, .. } => Some(*kind),
1741 _ => None,
1742 }
1743 }
1744
1745 #[tokio::test]
1748 async fn a_publication_read_has_an_overall_deadline() {
1749 const DID: &str = "did:plc:slowrepoaaaaaaaaaaaaaaaa";
1750 let plc = serve_answering(
1751 DID,
1752 200,
1753 (200, ""),
1754 std::time::Duration::from_millis(50),
1755 true,
1756 )
1757 .await;
1758 let started = std::time::Instant::now();
1759 let outcome =
1760 poll_publication_at(DID, plc, Some(std::time::Duration::from_millis(300))).await;
1761 assert!(
1762 started.elapsed() < std::time::Duration::from_secs(3),
1763 "the read ran {:?}",
1764 started.elapsed()
1765 );
1766 assert_eq!(
1767 kind_of(&outcome),
1768 Some(crate::feed::FailureKind::Fetch),
1769 "{outcome:?}"
1770 );
1771 }
1772
1773 #[tokio::test]
1777 async fn a_publication_failure_is_filed_under_what_happened() {
1778 let zero = std::time::Duration::ZERO;
1779 const GONE: &str = "did:plc:tombstonedaaaaaaaaaaaaaa";
1780 let plc = serve_answering(GONE, 404, (200, ""), zero, false).await;
1781 let outcome = poll_publication_at(GONE, plc, None).await;
1782 assert_eq!(
1783 kind_of(&outcome),
1784 Some(crate::feed::FailureKind::Status),
1785 "PLC 404: {outcome:?}"
1786 );
1787
1788 const NOREPO: &str = "did:plc:norepoaaaaaaaaaaaaaaaaaa";
1789 let plc = serve_answering(
1790 NOREPO,
1791 200,
1792 (400, r#"{"error":"RepoNotFound"}"#),
1793 zero,
1794 false,
1795 )
1796 .await;
1797 let outcome = poll_publication_at(NOREPO, plc, None).await;
1798 assert_eq!(
1799 kind_of(&outcome),
1800 Some(crate::feed::FailureKind::Status),
1801 "RepoNotFound: {outcome:?}"
1802 );
1803
1804 const GARBLED: &str = "did:plc:garbledaaaaaaaaaaaaaaaaa";
1805 let plc = serve_answering(GARBLED, 200, (200, r#"{"records":"x"}"#), zero, false).await;
1806 let outcome = poll_publication_at(GARBLED, plc, None).await;
1807 assert_eq!(
1808 kind_of(&outcome),
1809 Some(crate::feed::FailureKind::Parse),
1810 "garbled body: {outcome:?}"
1811 );
1812 }
1813
1814 #[tokio::test]
1817 async fn an_empty_publication_is_a_healthy_poll() {
1818 const EMPTY: &str = "did:plc:emptypubaaaaaaaaaaaaaaaa";
1819 let site = format!("at://{EMPTY}/{}/quiet", nsid::STANDARD_PUBLICATION);
1820 let (plc, _) = serve_repo(
1821 EMPTY,
1822 vec![(
1823 nsid::STANDARD_PUBLICATION,
1824 "quiet",
1825 json!({ "name": "Quiet", "url": "https://quiet.example" }),
1826 )],
1827 )
1828 .await;
1829 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1830 crate::store::upsert_feed(
1831 &pool,
1832 &crate::store::NewFeed {
1833 url: site.clone(),
1834 ..Default::default()
1835 },
1836 )
1837 .await
1838 .unwrap();
1839 let feed = crate::store::get_feed_by_url(&pool, &site)
1840 .await
1841 .unwrap()
1842 .unwrap();
1843 let mut config = crate::config::Config::default();
1844 config.oauth.plc_directory = plc;
1845 let client = crate::feed::build_client().unwrap();
1846 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1847 .await
1848 .unwrap();
1849 assert!(
1850 matches!(
1851 outcome,
1852 crate::feed::PollOutcome::Updated { new_entries: 0 }
1853 ),
1854 "an empty publication was not a healthy poll: {outcome:?}"
1855 );
1856 }
1857
1858 #[tokio::test]
1862 async fn fetch_skips_malformed_records_beside_good_ones() {
1863 const OWN: &str = "did:plc:malformedrepo";
1864 let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
1865 let doc = |title: &str| {
1866 json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
1867 "path": "/p", "site": site })
1868 };
1869 let (plc, _) = serve_repo(
1870 OWN,
1871 vec![
1872 (
1873 nsid::STANDARD_PUBLICATION,
1874 "",
1875 json!({ "name": "Broken", "url": "https://x.example" }),
1876 ),
1877 (
1878 nsid::STANDARD_PUBLICATION,
1879 "mine",
1880 json!({ "name": "Mine", "url": "https://mine.example" }),
1881 ),
1882 (nsid::STANDARD_DOCUMENT, "", doc("unreadable")),
1883 (nsid::STANDARD_DOCUMENT, "3l2mfaaaaaa2a", doc("kept")),
1884 ],
1885 )
1886 .await;
1887 let client = crate::feed::build_client().unwrap();
1888 let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
1889 .await
1890 .expect("a malformed record stalled a stranger's publication");
1891 assert!(read.complete);
1892 let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
1893 assert_eq!(titles, vec!["kept"]);
1894 }
1895
1896 const SHARED: &str = "did:plc:sharedrepoaaaaaaaaaaaaaa";
1899
1900 fn shared_doc(site_rkey: &str, title: &str, at: &str) -> serde_json::Value {
1901 json!({ "title": title, "publishedAt": at, "path": format!("/{title}"),
1902 "site": format!("at://{SHARED}/{}/{site_rkey}", nsid::STANDARD_PUBLICATION) })
1903 }
1904
1905 #[tokio::test]
1909 async fn two_publications_in_one_repo_cost_one_walk() {
1910 let (plc, hits) = serve_repo(
1911 SHARED,
1912 vec![
1913 (
1914 nsid::STANDARD_PUBLICATION,
1915 "alpha",
1916 json!({ "name": "Alpha", "url": "https://alpha.example" }),
1917 ),
1918 (
1919 nsid::STANDARD_PUBLICATION,
1920 "beta",
1921 json!({ "name": "Beta", "url": "https://beta.example" }),
1922 ),
1923 (
1924 nsid::STANDARD_DOCUMENT,
1925 "3l2shaaaaaa2a",
1926 shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
1927 ),
1928 (
1929 nsid::STANDARD_DOCUMENT,
1930 "3l2shaaaaaa2b",
1931 shared_doc("beta", "b1", "2026-07-10T00:00:00Z"),
1932 ),
1933 (
1934 nsid::STANDARD_DOCUMENT,
1935 "3l2shaaaaaa2c",
1936 shared_doc("alpha", "a2", "2026-07-09T00:00:00Z"),
1937 ),
1938 ],
1939 )
1940 .await;
1941 let client = crate::feed::build_client().unwrap();
1942 let reads = fetch_repo(
1943 &client,
1944 &plc,
1945 SHARED,
1946 &["alpha".to_string(), "beta".to_string()],
1947 )
1948 .await
1949 .unwrap();
1950 assert_eq!(
1951 hits.load(std::sync::atomic::Ordering::SeqCst),
1952 4,
1953 "two collections walked once each (a page then an empty page), not once per publication"
1954 );
1955 let titles = |r: &anyhow::Result<PublicationRead>| {
1956 let mut t: Vec<String> = r
1957 .as_ref()
1958 .unwrap()
1959 .entries
1960 .iter()
1961 .map(|e| e.title.clone())
1962 .collect();
1963 t.sort();
1964 t
1965 };
1966 assert_eq!(titles(&reads[0]), vec!["a1", "a2"]);
1967 assert_eq!(titles(&reads[1]), vec!["b1"]);
1968 }
1969
1970 #[tokio::test]
1973 async fn a_busy_publication_does_not_starve_its_quiet_sibling() {
1974 let mut records = vec![
1975 (
1976 nsid::STANDARD_PUBLICATION,
1977 "busy",
1978 json!({ "name": "Busy", "url": "https://busy.example" }),
1979 ),
1980 (
1981 nsid::STANDARD_PUBLICATION,
1982 "quiet",
1983 json!({ "name": "Quiet", "url": "https://quiet.example" }),
1984 ),
1985 ];
1986 for (rkey, title) in [
1987 ("3l2bsaaaaaa2a", "b1"),
1988 ("3l2bsaaaaaa2b", "b2"),
1989 ("3l2bsaaaaaa2c", "b3"),
1990 ("3l2bsaaaaaa2d", "b4"),
1991 ("3l2bsaaaaaa2e", "b5"),
1992 ] {
1993 records.push((
1994 nsid::STANDARD_DOCUMENT,
1995 rkey,
1996 shared_doc("busy", title, "2026-07-11T00:00:00Z"),
1997 ));
1998 }
1999 records.push((
2000 nsid::STANDARD_DOCUMENT,
2001 "3l2bsaaaaaa2f",
2002 shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
2003 ));
2004 let (plc, _) = serve_repo(SHARED, records).await;
2005 let client = crate::feed::build_client().unwrap();
2006 let reads = fetch_repo_capped(
2007 &client,
2008 &plc,
2009 SHARED,
2010 &["busy".to_string(), "quiet".to_string()],
2011 2,
2012 crate::atproto::MAX_LIST_BYTES,
2013 )
2014 .await
2015 .unwrap();
2016 let busy = reads[0].as_ref().unwrap();
2017 let quiet = reads[1].as_ref().unwrap();
2018 assert_eq!(
2019 busy.entries.len(),
2020 2,
2021 "the busy publication was not capped at its own cap"
2022 );
2023 assert!(!busy.complete, "a capped publication was reported complete");
2024 assert_eq!(quiet.entries.len(), 1, "the quiet sibling was starved");
2025 assert!(quiet.complete);
2026 }
2027
2028 #[tokio::test]
2035 #[ignore = "known limitation: a one-repo group shares one byte budget (#229)"]
2036 async fn big_siblings_do_not_spend_a_quiet_publications_share() {
2037 let body = "w".repeat(20 * 1024);
2038 let mut records = vec![
2039 (
2040 nsid::STANDARD_PUBLICATION,
2041 "big1",
2042 json!({ "name": "Big 1", "url": "https://b1.example" }),
2043 ),
2044 (
2045 nsid::STANDARD_PUBLICATION,
2046 "big2",
2047 json!({ "name": "Big 2", "url": "https://b2.example" }),
2048 ),
2049 (
2050 nsid::STANDARD_PUBLICATION,
2051 "quiet",
2052 json!({ "name": "Quiet", "url": "https://q.example" }),
2053 ),
2054 ];
2055 for i in 0..200 {
2056 let rkey: &'static str = Box::leak(format!("3l2big{i:06}").into_boxed_str());
2057 let site = if i % 2 == 0 { "big1" } else { "big2" };
2058 let mut doc = shared_doc(site, &format!("d{i}"), "2026-07-11T00:00:00Z");
2059 doc["textContent"] = json!(body);
2060 records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2061 }
2062 records.push((
2063 nsid::STANDARD_DOCUMENT,
2064 "3l2zzzzzzzzzz",
2065 shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
2066 ));
2067 let (plc, _) = serve_repo(SHARED, records).await;
2068 let client = crate::feed::build_client().unwrap();
2069 let rkeys: Vec<String> = ["big1", "big2", "quiet"]
2070 .iter()
2071 .map(|s| s.to_string())
2072 .collect();
2073 let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 4 * 1024 * 1024)
2074 .await
2075 .unwrap();
2076 let quiet = reads[2].as_ref().unwrap();
2077 assert_eq!(
2078 quiet.entries.len(),
2079 1,
2080 "the quiet publication was starved of bytes by its siblings"
2081 );
2082 assert!(
2083 !reads[0].as_ref().unwrap().complete,
2084 "a publication over its share was reported complete"
2085 );
2086 }
2087
2088 #[tokio::test]
2093 async fn a_busy_publication_beside_idle_siblings_reads_as_it_would_alone() {
2094 let body = "w".repeat(20 * 1024);
2095 let mut records: Vec<(&'static str, &'static str, serde_json::Value)> = Vec::new();
2096 let rkeys: Vec<String> = (0..16).map(|i| format!("p{i:02}")).collect();
2097 for r in &rkeys {
2098 let r: &'static str = Box::leak(r.clone().into_boxed_str());
2099 records.push((
2100 nsid::STANDARD_PUBLICATION,
2101 r,
2102 json!({ "name": r, "url": "https://p.example" }),
2103 ));
2104 }
2105 for i in 0..100 {
2106 let rkey: &'static str = Box::leak(format!("3l2bus{i:06}").into_boxed_str());
2107 let mut doc = shared_doc("p00", &format!("d{i}"), "2026-07-11T00:00:00Z");
2108 doc["textContent"] = json!(body);
2109 records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2110 }
2111 let (plc, _) = serve_repo(SHARED, records).await;
2112 let client = crate::feed::build_client().unwrap();
2113 let budget = 4 * 1024 * 1024;
2114 let alone = fetch_repo_capped(&client, &plc, SHARED, &rkeys[..1], 2_000, budget)
2115 .await
2116 .unwrap();
2117 let grouped = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, budget)
2118 .await
2119 .unwrap();
2120 let (a, g) = (alone[0].as_ref().unwrap(), grouped[0].as_ref().unwrap());
2121 assert_eq!((a.entries.len(), a.complete), (100, true));
2122 assert_eq!(
2123 (g.entries.len(), g.complete),
2124 (100, true),
2125 "grouping read it worse than alone"
2126 );
2127 }
2128
2129 #[tokio::test]
2132 async fn one_large_document_does_not_fail_its_publication_in_a_group() {
2133 let mut big = shared_doc("big", "huge", "2026-07-11T00:00:00Z");
2134 big["textContent"] = json!("w".repeat(500 * 1024));
2135 let records = vec![
2136 (
2137 nsid::STANDARD_PUBLICATION,
2138 "big",
2139 json!({ "name": "Big", "url": "https://b.example" }),
2140 ),
2141 (
2142 nsid::STANDARD_PUBLICATION,
2143 "other",
2144 json!({ "name": "Other", "url": "https://o.example" }),
2145 ),
2146 (nsid::STANDARD_DOCUMENT, "3l2hugeaaaa2a", big),
2147 ];
2148 let (plc, _) = serve_repo(SHARED, records).await;
2149 let client = crate::feed::build_client().unwrap();
2150 let rkeys = vec!["big".to_string(), "other".to_string()];
2151 let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 1024 * 1024)
2152 .await
2153 .unwrap();
2154 assert_eq!(
2155 reads[0].as_ref().unwrap().entries.len(),
2156 1,
2157 "a large document was dropped"
2158 );
2159 }
2160
2161 #[tokio::test]
2164 async fn every_requested_rkey_gets_a_result() {
2165 let (plc, _) = serve_repo(
2166 SHARED,
2167 vec![(
2168 nsid::STANDARD_PUBLICATION,
2169 "alpha",
2170 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2171 )],
2172 )
2173 .await;
2174 let client = crate::feed::build_client().unwrap();
2175 let reads = fetch_repo(
2176 &client,
2177 &plc,
2178 SHARED,
2179 &["nope".to_string(), "nope".to_string()],
2180 )
2181 .await
2182 .unwrap();
2183 assert_eq!(reads.len(), 2);
2184 assert!(reads.iter().all(|r| r.is_err()));
2185 }
2186
2187 #[tokio::test]
2190 async fn repeated_and_empty_requests_are_handled() {
2191 let (plc, hits) = serve_repo(
2192 SHARED,
2193 vec![
2194 (
2195 nsid::STANDARD_PUBLICATION,
2196 "alpha",
2197 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2198 ),
2199 (
2200 nsid::STANDARD_DOCUMENT,
2201 "3l2rpaaaaaa2a",
2202 shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
2203 ),
2204 ],
2205 )
2206 .await;
2207 let client = crate::feed::build_client().unwrap();
2208 let none = fetch_repo(&client, &plc, SHARED, &[]).await.unwrap();
2209 assert!(none.is_empty());
2210 assert_eq!(
2211 hits.load(std::sync::atomic::Ordering::SeqCst),
2212 0,
2213 "an empty request reached the network"
2214 );
2215 let twice = fetch_repo(
2216 &client,
2217 &plc,
2218 SHARED,
2219 &["alpha".to_string(), "alpha".to_string()],
2220 )
2221 .await
2222 .unwrap();
2223 for read in &twice {
2224 assert_eq!(
2225 read.as_ref().unwrap().entries.len(),
2226 1,
2227 "a repeated rkey read as empty"
2228 );
2229 }
2230 }
2231
2232 #[tokio::test]
2235 async fn a_group_stores_each_publications_documents_under_its_own_feed() {
2236 let (plc, _) = serve_repo(
2237 SHARED,
2238 vec![
2239 (
2240 nsid::STANDARD_PUBLICATION,
2241 "alpha",
2242 json!({ "name": "Alpha", "url": "https://alpha.example" }),
2243 ),
2244 (
2245 nsid::STANDARD_PUBLICATION,
2246 "beta",
2247 json!({ "name": "Beta", "url": "https://beta.example" }),
2248 ),
2249 (
2250 nsid::STANDARD_DOCUMENT,
2251 "3l2grpaaaaa2a",
2252 shared_doc("alpha", "only-alpha", "2026-07-11T00:00:00Z"),
2253 ),
2254 (
2255 nsid::STANDARD_DOCUMENT,
2256 "3l2grpaaaaa2b",
2257 shared_doc("beta", "only-beta", "2026-07-10T00:00:00Z"),
2258 ),
2259 ],
2260 )
2261 .await;
2262 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2263 let mut feeds = Vec::new();
2264 for rkey in ["alpha", "beta"] {
2265 let url = format!("at://{SHARED}/{}/{rkey}", nsid::STANDARD_PUBLICATION);
2266 crate::store::upsert_feed(
2267 &pool,
2268 &crate::store::NewFeed {
2269 url: url.clone(),
2270 ..Default::default()
2271 },
2272 )
2273 .await
2274 .unwrap();
2275 feeds.push(
2276 crate::store::get_feed_by_url(&pool, &url)
2277 .await
2278 .unwrap()
2279 .unwrap(),
2280 );
2281 }
2282 let mut config = crate::config::Config::default();
2283 config.oauth.plc_directory = plc;
2284 let client = crate::feed::build_client().unwrap();
2285 let outcomes = crate::feed::poll_publication_group(&pool, &client, &config, &feeds).await;
2286 assert_eq!(outcomes.len(), 2);
2287 for (feed, want) in feeds.iter().zip(["only-alpha", "only-beta"]) {
2288 let titles: Vec<String> =
2289 sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ?")
2290 .bind(feed.id)
2291 .fetch_all(&pool)
2292 .await
2293 .unwrap();
2294 assert_eq!(
2295 titles,
2296 vec![want.to_string()],
2297 "{} got another publication's documents",
2298 feed.url
2299 );
2300 }
2301 }
2302
2303 #[tokio::test]
2309 async fn a_publication_subscription_delivers_entries_end_to_end() {
2310 const A0: &str = "did:plc:acceptanceaaaaaaaaaaaaaa";
2311 let site = format!("at://{A0}/{}/a0pub", nsid::STANDARD_PUBLICATION);
2312 let document = |title: &str, path: &str| {
2313 json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
2314 "path": path, "site": site, "textContent": "body" })
2315 };
2316 let (plc, _hits) = serve_repo(
2317 A0,
2318 vec![
2319 (
2320 nsid::STANDARD_PUBLICATION,
2321 "a0pub",
2322 json!({ "name": "A0 Journal", "url": "https://a0.example" }),
2323 ),
2324 (
2325 nsid::STANDARD_DOCUMENT,
2326 "3l2a0aaaaaa2a",
2327 document("First post", "/first"),
2328 ),
2329 (
2330 nsid::STANDARD_DOCUMENT,
2331 "3l2a0aaaaaa2b",
2332 document("Second post", "/second"),
2333 ),
2334 ],
2335 )
2336 .await;
2337
2338 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2339 crate::store::upsert_feed(
2340 &pool,
2341 &crate::store::NewFeed {
2342 url: site.clone(),
2343 ..Default::default()
2344 },
2345 )
2346 .await
2347 .unwrap();
2348 let mut config = crate::config::Config::default();
2350 config.oauth.plc_directory = plc;
2351
2352 let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
2353 let feed =
2355 crate::store::due_feeds_of_kind(&pool, &now, crate::feed::FeedKind::Publication, 50)
2356 .await
2357 .unwrap()
2358 .into_iter()
2359 .find(|f| f.url == site)
2360 .expect("the publication is not handed to the publication poller");
2361
2362 let client = crate::feed::build_client().unwrap();
2363 let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
2364 .await
2365 .unwrap();
2366 assert!(
2367 matches!(
2368 outcome,
2369 crate::feed::PollOutcome::Updated { new_entries: 2 }
2370 ),
2371 "expected two new entries, got {outcome:?}"
2372 );
2373 let titles: Vec<String> =
2374 sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ? ORDER BY title")
2375 .bind(feed.id)
2376 .fetch_all(&pool)
2377 .await
2378 .unwrap();
2379 assert_eq!(titles, vec!["First post", "Second post"]);
2380 }
2381
2382 fn canonical(rkey: &str) -> String {
2383 format!("at://{DID}/{}/{rkey}", nsid::STANDARD_PUBLICATION)
2384 }
2385
2386 #[test]
2395 fn the_site_filter_uses_the_uri_the_pds_minted() {
2396 let records = vec![publication("p", "https://scanash.com")];
2397 let (site, pubn) = publication_from_records("p", &records).expect("publication not found");
2398 assert_eq!(site, canonical("p"), "did not take the PDS's canonical URI");
2399
2400 let docs = vec![document("d1", &canonical("p"), "Hello", "/hello")];
2401 let entries = entries_from_records(&site, &pubn, &docs);
2402 assert_eq!(
2403 entries.len(),
2404 1,
2405 "a canonical-site document was not matched"
2406 );
2407 }
2408
2409 #[test]
2412 fn an_entry_maps_onto_the_stores_row() {
2413 let records = vec![publication("p", "https://example.com")];
2414 let (site, pubn) = publication_from_records("p", &records).unwrap();
2415 let docs = vec![document("rk1", &site, "Hello", "/hello")];
2416 let row: crate::store::NewEntry = entries_from_records(&site, &pubn, &docs)
2417 .pop()
2418 .unwrap()
2419 .into();
2420 assert_eq!(
2421 row.guid,
2422 format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
2423 );
2424 assert_eq!(row.url.as_deref(), Some("https://example.com/hello"));
2425 assert_eq!(row.title.as_deref(), Some("Hello"));
2426 assert_eq!(row.published.as_deref(), Some("2026-07-11T00:00:00Z"));
2427 assert_eq!(row.content_html.as_deref(), Some("body"));
2428 assert_eq!(row.author, None);
2429 assert_eq!(row.fetched_at, None);
2430 }
2431
2432 #[test]
2435 fn documents_are_filtered_by_their_site_field() {
2436 let records = vec![publication("mine", "https://example.com")];
2437 let (site, pubn) = publication_from_records("mine", &records).unwrap();
2438 let docs = vec![
2439 document("a", &site, "Mine", "/a"),
2440 document("b", &canonical("theirs"), "Theirs", "/b"),
2441 document("c", &site, "Mine again", "/c"),
2442 ];
2443 let titles: Vec<String> = entries_from_records(&site, &pubn, &docs)
2444 .into_iter()
2445 .map(|e| e.title)
2446 .collect();
2447 assert_eq!(titles, ["Mine", "Mine again"]);
2448 }
2449
2450 #[test]
2457 fn a_publication_with_a_hostile_url_is_refused() {
2458 for hostile in [
2459 "javascript:alert(1)",
2460 "data:text/html,<script>",
2461 "file:///etc/passwd",
2462 "",
2463 ] {
2464 let records = vec![publication("p", hostile)];
2465 assert!(
2466 publication_from_records("p", &records).is_none(),
2467 "accepted a publication whose url is {hostile:?}",
2468 );
2469 }
2470 }
2471
2472 #[test]
2478 fn entry_urls_are_joined_against_the_publication_base() {
2479 let records = vec![publication("p", "https://example.com/blog")];
2480 let (site, pubn) = publication_from_records("p", &records).unwrap();
2481 let docs = vec![
2482 document("a", &site, "Relative", "/a"),
2483 document("b", &site, "Absolute-looking", "https://evil.example/x"),
2484 ];
2485 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
2486 .into_iter()
2487 .map(|e| e.url)
2488 .collect();
2489 assert_eq!(urls[0].as_deref(), Some("https://example.com/a"));
2490 assert_eq!(
2498 urls[1], None,
2499 "a document path that escapes its publication's origin must yield no URL"
2500 );
2501 }
2502
2503 #[test]
2506 fn a_documents_text_fields_are_bounded_before_they_are_stored() {
2507 let site = canonical("pub");
2508 let big = "x".repeat(8 * 1024 * 1024);
2509 let records = vec![
2510 publication("pub", "https://scanash.com"),
2511 rec(
2512 nsid::STANDARD_DOCUMENT,
2513 "3l2bigaaaaa2a",
2514 json!({ "title": big, "publishedAt": "2026-07-11T00:00:00Z",
2515 "path": format!("/{}", "p".repeat(20_000)), "site": site,
2516 "textContent": "<".repeat(3 * 1024 * 1024) }),
2517 ),
2518 ];
2519 let (_, publication) = publication_from_records("pub", &records).unwrap();
2520 let entries = entries_from_records(&site, &publication, &records);
2521 let e = &entries[0];
2522 assert!(
2523 e.title.len() <= crate::feed::MAX_TITLE_BYTES,
2524 "title: {}",
2525 e.title.len()
2526 );
2527 let url = e
2528 .url
2529 .as_ref()
2530 .expect("an overlong path is truncated, not dropped");
2531 assert!(
2532 url.len() <= crate::feed::MAX_URL_BYTES,
2533 "url: {}",
2534 url.len()
2535 );
2536 let summary = e.summary.as_ref().unwrap();
2537 assert!(
2538 summary.len() <= crate::feed::MAX_CONTENT_HTML_BYTES,
2539 "the ESCAPED summary is what is stored: {}",
2540 summary.len()
2541 );
2542 }
2543
2544 #[test]
2548 fn a_documents_uri_is_bounded_as_an_entry_id() {
2549 let site = canonical("pub");
2550 let records = vec![
2551 publication("pub", "https://scanash.com"),
2552 rec(
2553 nsid::STANDARD_DOCUMENT,
2554 &"k".repeat(100_000),
2555 json!({ "title": "t", "publishedAt": "2026-07-11T00:00:00Z",
2556 "path": "/p", "site": site }),
2557 ),
2558 ];
2559 let (_, publication) = publication_from_records("pub", &records).unwrap();
2560 let entries = entries_from_records(&site, &publication, &records);
2561 let stored: crate::store::NewEntry = entries[0].clone().into();
2562 assert!(
2563 stored.guid.len() <= crate::feed::MAX_GUID_BYTES,
2564 "guid: {}",
2565 stored.guid.len()
2566 );
2567 }
2568
2569 #[test]
2570 fn a_publications_own_name_is_bounded() {
2571 let records = vec![rec(
2572 nsid::STANDARD_PUBLICATION,
2573 "pub",
2574 json!({ "name": "n".repeat(100_000), "url": "https://scanash.com" }),
2575 )];
2576 let (_, publication) = publication_from_records("pub", &records).unwrap();
2577 assert!(publication.name.unwrap().len() <= crate::feed::MAX_TITLE_BYTES);
2578 }
2579
2580 #[test]
2582 fn a_document_with_neither_summary_field_still_yields_an_entry() {
2583 let records = vec![publication("p", "https://example.com/")];
2584 let (site, pubn) = publication_from_records("p", &records).unwrap();
2585 let bare = rec(
2586 nsid::STANDARD_DOCUMENT,
2587 "bare",
2588 json!({
2589 "title": "Bare",
2590 "publishedAt": "2026-07-11T00:00:00Z",
2591 "path": "/bare",
2592 "site": site,
2593 }),
2594 );
2595 let entries = entries_from_records(&site, &pubn, &[bare]);
2596 assert_eq!(entries.len(), 1);
2597 assert_eq!(entries[0].summary, None);
2598 assert_eq!(entries[0].url.as_deref(), Some("https://example.com/bare"));
2599 }
2600
2601 #[test]
2604 fn an_empty_description_does_not_shadow_the_body() {
2605 let records = vec![publication("p", "https://example.com")];
2606 let (site, pubn) = publication_from_records("p", &records).unwrap();
2607 let doc = rec(
2608 nsid::STANDARD_DOCUMENT,
2609 "d",
2610 json!({
2611 "title": "T",
2612 "publishedAt": "2026-07-11T00:00:00Z",
2613 "path": "/d",
2614 "site": site,
2615 "description": " ",
2616 "textContent": "the real body",
2617 }),
2618 );
2619 let entries = entries_from_records(&site, &pubn, &[doc]);
2620 assert_eq!(entries[0].summary.as_deref(), Some("the real body"));
2621 }
2622
2623 #[test]
2630 fn published_at_is_normalised_or_dropped() {
2631 let records = vec![publication("p", "https://example.com")];
2632 let (site, pubn) = publication_from_records("p", &records).unwrap();
2633 let with = |rkey: &str, published_at: serde_json::Value| {
2634 rec(
2635 nsid::STANDARD_DOCUMENT,
2636 rkey,
2637 json!({ "title": "T", "publishedAt": published_at, "path": "/x", "site": site }),
2638 )
2639 };
2640 let docs = vec![
2644 with("a", json!("2026-07-11T09:30:00.123+02:00")),
2645 with("b", json!("yesterday-ish")),
2646 with("c", json!("2026-07-11T00:00:00Z")),
2647 ];
2648 let published: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
2649 .into_iter()
2650 .map(|e| e.published)
2651 .collect();
2652 assert_eq!(
2653 published,
2654 vec![
2655 Some("2026-07-11T07:30:00Z".to_string()),
2656 None,
2657 Some("2026-07-11T00:00:00Z".to_string()),
2658 ],
2659 "publishedAt was not normalised to the store's spelling"
2660 );
2661 }
2662
2663 #[test]
2675 fn an_undated_document_is_dated_from_its_tid_rkey() {
2676 let records = vec![publication("p", "https://example.com")];
2677 let (site, pubn) = publication_from_records("p", &records).unwrap();
2678 let docs = vec![rec(
2679 nsid::STANDARD_DOCUMENT,
2680 PAST_TID,
2681 json!({ "title": "T", "path": "/x", "site": site }),
2682 )];
2683 assert_eq!(
2684 entries_from_records(&site, &pubn, &docs)
2685 .into_iter()
2686 .next()
2687 .expect("the document is an entry")
2688 .published,
2689 Some(PAST_TID_WRITTEN_AT.to_string()),
2690 "the date must come from the record key, and be spelled the way the store spells dates"
2691 );
2692 }
2693
2694 #[test]
2695 fn an_unparseable_published_at_falls_back_to_the_tid_rkey() {
2696 let records = vec![publication("p", "https://example.com")];
2697 let (site, pubn) = publication_from_records("p", &records).unwrap();
2698 let docs = vec![rec(
2699 nsid::STANDARD_DOCUMENT,
2700 PAST_TID,
2701 json!({ "title": "T", "publishedAt": "yesterday-ish", "path": "/x", "site": site }),
2702 )];
2703 assert_eq!(
2704 entries_from_records(&site, &pubn, &docs)
2705 .into_iter()
2706 .next()
2707 .expect("the document is an entry")
2708 .published,
2709 Some(PAST_TID_WRITTEN_AT.to_string()),
2710 "a date the parser cannot read is no date at all, so the rkey must stand in"
2711 );
2712 }
2713
2714 #[test]
2715 fn a_stated_date_outranks_the_rkey() {
2716 let records = vec![publication("p", "https://example.com")];
2717 let (site, pubn) = publication_from_records("p", &records).unwrap();
2718 let docs = vec![rec(
2721 nsid::STANDARD_DOCUMENT,
2722 PAST_TID,
2723 json!({ "title": "T", "publishedAt": "2020-01-02T00:00:00Z", "path": "/x", "site": site }),
2724 )];
2725 assert_eq!(
2726 entries_from_records(&site, &pubn, &docs)
2727 .into_iter()
2728 .next()
2729 .expect("the document is an entry")
2730 .published,
2731 Some("2020-01-02T00:00:00Z".to_string()),
2732 "the rkey records when the file was written, which is not when the post was published"
2733 );
2734 }
2735
2736 #[test]
2745 fn a_future_dated_document_falls_back_to_its_rkey() {
2746 let records = vec![publication("p", "https://example.com")];
2747 let (site, pubn) = publication_from_records("p", &records).unwrap();
2748 let docs = vec![rec(
2749 nsid::STANDARD_DOCUMENT,
2750 PAST_TID,
2751 json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
2752 )];
2753 assert_eq!(
2754 entries_from_records(&site, &pubn, &docs)
2755 .into_iter()
2756 .next()
2757 .expect("the document is an entry")
2758 .published,
2759 Some(PAST_TID_WRITTEN_AT.to_string()),
2760 "the date must be the record's write time, not the hour the poll happened to run"
2761 );
2762 }
2763
2764 #[test]
2765 fn a_future_dated_document_without_a_tid_rkey_is_undated() {
2766 let records = vec![publication("p", "https://example.com")];
2767 let (site, pubn) = publication_from_records("p", &records).unwrap();
2768 let docs = vec![rec(
2769 nsid::STANDARD_DOCUMENT,
2770 "self",
2771 json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
2772 )];
2773 assert_eq!(
2774 entries_from_records(&site, &pubn, &docs)
2775 .into_iter()
2776 .next()
2777 .expect("the document is an entry")
2778 .published,
2779 None,
2780 "with nothing credible to date it by, the row falls to fetched_at, which holds still"
2781 );
2782 }
2783
2784 #[test]
2792 fn a_stated_date_a_little_ahead_of_our_clock_is_still_believed() {
2793 let records = vec![publication("p", "https://example.com")];
2794 let (site, pubn) = publication_from_records("p", &records).unwrap();
2795 let slightly_ahead =
2796 crate::feed::fmt_time(chrono::Utc::now() + chrono::Duration::seconds(10));
2797 let docs = vec![rec(
2798 nsid::STANDARD_DOCUMENT,
2799 "self",
2800 json!({ "title": "T", "publishedAt": slightly_ahead, "path": "/x", "site": site }),
2801 )];
2802 assert_eq!(
2803 entries_from_records(&site, &pubn, &docs)
2804 .into_iter()
2805 .next()
2806 .expect("the document is an entry")
2807 .published,
2808 Some(slightly_ahead),
2809 "a few seconds of clock skew must not cost the entry its date"
2810 );
2811 }
2812
2813 #[test]
2814 fn a_document_with_neither_a_date_nor_a_tid_rkey_stays_undated() {
2815 let records = vec![publication("p", "https://example.com")];
2816 let (site, pubn) = publication_from_records("p", &records).unwrap();
2817 let docs = vec![
2824 rec(
2825 nsid::STANDARD_DOCUMENT,
2826 "my-first-post",
2827 json!({ "title": "T", "path": "/x", "site": site }),
2828 ),
2829 rec(
2830 nsid::STANDARD_DOCUMENT,
2831 "abcdefghijklm",
2832 json!({ "title": "T", "path": "/y", "site": site }),
2833 ),
2834 ];
2835 assert_eq!(
2836 entries_from_records(&site, &pubn, &docs)
2837 .into_iter()
2838 .map(|e| e.published)
2839 .collect::<Vec<_>>(),
2840 vec![None, None],
2841 "an invented date is worse than no date; the store decides what to do with undated rows"
2842 );
2843 }
2844
2845 #[test]
2855 fn summaries_are_escaped_as_plain_text_not_sanitised_as_markup() {
2856 let records = vec![publication("p", "https://example.com")];
2857 let (site, pubn) = publication_from_records("p", &records).unwrap();
2858 let doc = |rkey: &str, body: &str| {
2859 rec(
2860 nsid::STANDARD_DOCUMENT,
2861 rkey,
2862 json!({
2863 "title": "T",
2864 "publishedAt": "2026-07-11T00:00:00Z",
2865 "path": "/d",
2866 "site": site,
2867 "textContent": body,
2868 }),
2869 )
2870 };
2871 let summaries: Vec<String> = entries_from_records(
2872 &site,
2873 &pubn,
2874 &[
2875 doc("a", "Vec<String> is a type"),
2876 doc("b", "<script>alert(1)</script>"),
2877 ],
2878 )
2879 .into_iter()
2880 .filter_map(|e| e.summary)
2881 .collect();
2882 assert_eq!(
2883 summaries[0], "Vec<String> is a type",
2884 "prose was eaten by an HTML parser"
2885 );
2886 assert!(
2887 !summaries[1].contains("<script"),
2888 "escaping failed: {}",
2889 summaries[1]
2890 );
2891 }
2892
2893 #[test]
2901 fn an_entry_url_is_never_an_unvetted_path() {
2902 let pubn = Publication {
2903 name: None,
2904 url: "not a url".to_string(),
2907 };
2908 let site = canonical("p");
2909 let docs = vec![
2910 document("a", &site, "Hostile", "javascript:alert(1)"),
2911 document("b", &site, "Fine", "https://example.com/ok"),
2912 ];
2913 let entries = entries_from_records(&site, &pubn, &docs);
2914 assert_eq!(
2915 entries[0].url, None,
2916 "an unvetted path became an entry link"
2917 );
2918 assert_eq!(
2923 entries[1].url, None,
2924 "an off-origin absolute URL was published under the publication's name"
2925 );
2926 }
2927
2928 #[test]
2932 fn a_document_without_published_at_is_still_an_entry() {
2933 let records = vec![publication("p", "https://example.com")];
2934 let (site, pubn) = publication_from_records("p", &records).unwrap();
2935 let doc = rec(
2936 nsid::STANDARD_DOCUMENT,
2937 "d",
2938 json!({ "title": "T", "path": "/d", "site": site }),
2939 );
2940 let entries = entries_from_records(&site, &pubn, &[doc]);
2941 assert_eq!(
2942 entries.len(),
2943 1,
2944 "a missing publishedAt dropped the document"
2945 );
2946 assert_eq!(entries[0].published, None);
2947 }
2948
2949 #[tokio::test]
2955 async fn fetch_refuses_a_uri_for_another_collection() {
2956 let uri = AtUri::parse(&format!("at://{DID}/app.bsky.feed.post/3lab")).unwrap();
2957 let err = fetch(&reqwest::Client::new(), "https://plc.example", &uri)
2958 .await
2959 .expect_err("read a feed post as a publication");
2960 assert!(
2961 format!("{err:#}").contains(nsid::STANDARD_PUBLICATION),
2962 "failed for the wrong reason: {err:#}"
2963 );
2964 }
2965
2966 #[test]
2971 fn a_subpath_publication_keeps_its_base_path() {
2972 let records = vec![publication("p", "https://example.com/blog")];
2973 let (site, pubn) = publication_from_records("p", &records).unwrap();
2974 let docs = vec![document("a", &site, "Relative", "posts/a")];
2975 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
2976 .into_iter()
2977 .map(|e| e.url)
2978 .collect();
2979 assert_eq!(urls[0].as_deref(), Some("https://example.com/blog/posts/a"));
2980 }
2981
2982 #[test]
2986 fn no_parseable_base_means_no_url_not_any_url() {
2987 let pubn = Publication {
2988 name: None,
2989 url: "not a url".to_string(),
2990 };
2991 let site = canonical("p");
2992 let docs = vec![document("a", &site, "Absolute", "https://evil.example/x")];
2993 let entries = entries_from_records(&site, &pubn, &docs);
2994 assert_eq!(
2995 entries[0].url, None,
2996 "an off-origin absolute URL was published"
2997 );
2998 }
2999
3000 #[test]
3006 fn an_off_origin_path_yields_no_url_rather_than_the_homepage() {
3007 let records = vec![publication("p", "https://example.com/blog")];
3008 let (site, pubn) = publication_from_records("p", &records).unwrap();
3009 let docs = vec![
3010 document("a", &site, "Elsewhere", "https://www.example.com/post"),
3011 document("b", &site, "Home", "/ok"),
3012 ];
3013 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3014 .into_iter()
3015 .map(|e| e.url)
3016 .collect();
3017 assert_eq!(
3018 urls[0], None,
3019 "an off-origin path was rewritten to the base"
3020 );
3021 assert_eq!(urls[1].as_deref(), Some("https://example.com/ok"));
3022 }
3023
3024 #[test]
3029 fn a_blank_path_yields_no_url() {
3030 let records = vec![publication("p", "https://example.com/blog")];
3031 let (site, pubn) = publication_from_records("p", &records).unwrap();
3032 let docs = vec![
3033 document("a", &site, "Blank", ""),
3034 document("b", &site, "Spaces", " "),
3035 document("c", &site, "Real", "/real"),
3036 ];
3037 let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3038 .into_iter()
3039 .map(|e| e.url)
3040 .collect();
3041 assert_eq!(urls[0], None, "a blank path became the homepage");
3042 assert_eq!(urls[1], None, "a whitespace path became the homepage");
3043 assert_eq!(urls[2].as_deref(), Some("https://example.com/real"));
3044 }
3045
3046 #[test]
3049 fn a_document_without_a_path_is_still_an_entry() {
3050 let records = vec![publication("p", "https://example.com")];
3051 let (site, pubn) = publication_from_records("p", &records).unwrap();
3052 let doc = rec(
3053 nsid::STANDARD_DOCUMENT,
3054 "d",
3055 json!({ "title": "T", "publishedAt": "2026-07-11T00:00:00Z", "site": site }),
3056 );
3057 let entries = entries_from_records(&site, &pubn, &[doc]);
3058 assert_eq!(
3059 entries.len(),
3060 1,
3061 "a missing path dropped the whole document"
3062 );
3063 assert_eq!(entries[0].title, "T");
3064 assert_eq!(entries[0].url, None);
3065 }
3066
3067 #[test]
3073 fn documents_are_classified_keep_sibling_or_orphan() {
3074 let pubs = [
3075 publication("a", "https://example.com"),
3076 publication("b", "https://b.example"),
3077 ];
3078 let known: std::collections::HashSet<&str> = pubs.iter().map(|p| p.uri.as_str()).collect();
3079 let mine = canonical("a");
3080 let wanted: std::collections::HashMap<String, usize> = [(mine.clone(), 0)].into();
3081 let fate = |d: &crate::atproto::RecordEntry| classify_document(d, &wanted, &known);
3082
3083 assert_eq!(
3084 fate(&document("d1", &mine, "Mine", "/1")),
3085 DocumentFate::Keep(0)
3086 );
3087 assert_eq!(
3088 fate(&document("d2", &canonical("b"), "B's", "/2")),
3089 DocumentFate::Sibling,
3090 "a sibling publication's document is not an orphan"
3091 );
3092 assert_eq!(
3093 fate(&document(
3094 "d3",
3095 "at://did:plc:other/site.standard.publication/x",
3096 "?",
3097 "/3"
3098 )),
3099 DocumentFate::Orphan
3100 );
3101 assert_eq!(
3102 fate(&rec(
3103 nsid::STANDARD_DOCUMENT,
3104 "d4",
3105 json!({"title": "no rest"})
3106 )),
3107 DocumentFate::Malformed
3108 );
3109 }
3110
3111 #[test]
3113 fn at_uri_parsing_uses_the_shared_prefix() {
3114 let uri = format!(
3115 "{}{DID}/{}/abc",
3116 crate::atproto::AT_URI_PREFIX,
3117 nsid::STANDARD_PUBLICATION
3118 );
3119 assert!(AtUri::parse(&uri).is_some());
3120 }
3121
3122 #[test]
3125 fn the_guid_is_the_record_uri_not_the_path() {
3126 let records = vec![publication("p", "https://example.com")];
3127 let (site, pubn) = publication_from_records("p", &records).unwrap();
3128 let docs = vec![document("rk1", &site, "T", "/moved")];
3129 let entries = entries_from_records(&site, &pubn, &docs);
3130 assert_eq!(
3131 entries[0].guid,
3132 format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
3133 );
3134 }
3135
3136 #[test]
3138 fn a_malformed_document_is_skipped_rather_than_fatal() {
3139 let records = vec![publication("p", "https://example.com")];
3140 let (site, pubn) = publication_from_records("p", &records).unwrap();
3141 let docs = vec![
3142 rec(
3143 nsid::STANDARD_DOCUMENT,
3144 "bad",
3145 json!({ "title": "no rest" }),
3146 ),
3147 document("ok", &site, "Good", "/good"),
3148 ];
3149 let entries = entries_from_records(&site, &pubn, &docs);
3150 assert_eq!(entries.len(), 1);
3151 assert_eq!(entries[0].title, "Good");
3152 }
3153
3154 #[test]
3156 fn a_missing_publication_is_none() {
3157 let records = vec![publication("other", "https://example.com")];
3158 assert!(publication_from_records("p", &records).is_none());
3159 }
3160
3161 #[test]
3162 fn at_uris_parse_in_both_forms_and_reject_malformed_ones() {
3163 let did = AtUri::parse(&format!("at://{DID}/site.standard.publication/abc")).unwrap();
3164 assert_eq!(did.authority, DID);
3165 assert_eq!(did.rkey, "abc");
3166 assert_eq!(
3167 did.to_string(),
3168 format!("at://{DID}/site.standard.publication/abc")
3169 );
3170 assert!(AtUri::parse("at://alice.example.com/site.standard.publication/abc").is_some());
3171 for bad in [
3172 "at://",
3173 "at://only-authority",
3174 "at://authority/collection",
3175 "at://authority/collection/",
3176 "at:///collection/rkey",
3177 "at://authority/collection/rkey/extra",
3178 "https://example.com/feed.xml",
3179 "at:authority/collection/rkey",
3180 ] {
3181 assert!(AtUri::parse(bad).is_none(), "parsed {bad:?}");
3182 }
3183 }
3184}