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