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