1use anyhow::{Context as _, Result};
18use reqwest::Client;
19use serde_json::{json, Value};
20use sqlx::SqlitePool;
21
22use crate::atproto::{RecordEntry, WriteOp, WriteResult};
23
24use super::dpop::Endpoint;
25use super::keys::SigningKey;
26use super::request::{self, DpopBody, DpopRequest, Retry};
27use super::store::OAuthSession;
28
29const MAX_LIST_PAGES: usize = 50;
34
35const MAX_LIST_RECORDS: usize = 5_000;
39
40pub struct Repo<'a> {
42 pub http: &'a Client,
43 pub pool: &'a SqlitePool,
44 pub session: &'a OAuthSession,
45 pub key: &'a SigningKey,
47}
48
49impl Repo<'_> {
50 fn url(&self, nsid: &str) -> String {
52 format!("{}/xrpc/{nsid}", self.session.aud.trim_end_matches('/'))
53 }
54
55 async fn send_raw(
57 &self,
58 url: &str,
59 body: DpopBody<'_>,
60 nsid: &str,
61 ) -> Result<request::PostOutcome> {
62 let outcome = request::send_with_dpop(
63 self.http,
64 self.pool,
65 &DpopRequest {
66 endpoint: Endpoint::ResourceServer,
67 url,
68 key: self.key,
69 access_token: Some(&self.session.access_token),
70 body,
71 retry: Retry::Allowed,
80 },
81 )
82 .await?;
83
84 if !outcome.is_success() {
85 let rejection = Self::rejection(&outcome.body, outcome.status);
92 return Err(anyhow::Error::new(rejection).context(format!(
93 "{nsid} failed: {}",
94 xrpc_error(&outcome.body, outcome.status)
95 )));
96 }
97 Ok(outcome)
98 }
99
100 fn rejection(body: &[u8], status: u16) -> crate::atproto::AtProtoError {
107 #[derive(serde::Deserialize)]
108 struct Fields {
109 error: Option<String>,
110 message: Option<String>,
111 }
112 let fields = super::error_body_worth_parsing(body)
113 .then(|| serde_json::from_slice::<Fields>(body).ok())
114 .flatten();
115 let (error, message) = match fields {
116 Some(Fields {
117 error: Some(error),
118 message,
119 }) => (error, message),
120 _ => ("Unknown".to_string(), None),
121 };
122 crate::atproto::AtProtoError::Xrpc {
123 status: reqwest::StatusCode::from_u16(status)
124 .unwrap_or(reqwest::StatusCode::BAD_GATEWAY),
127 error,
128 message,
129 }
130 }
131
132 async fn send(&self, url: &str, body: DpopBody<'_>, nsid: &str) -> Result<Value> {
138 let outcome = self.send_raw(url, body, nsid).await?;
139 if outcome.body.is_empty() {
141 return Ok(Value::Null);
142 }
143 outcome.json()
144 }
145
146 fn error_fields(body: &[u8]) -> Option<String> {
159 if !super::error_body_worth_parsing(body) {
160 return None;
161 }
162 let value: Value = serde_json::from_slice(body).ok()?;
163 let kind = value.get("error").and_then(Value::as_str)?;
164 match value.get("message").and_then(Value::as_str) {
165 Some(message) => Some(format!("{kind}: {message}")),
166 None => Some(kind.to_string()),
167 }
168 }
169
170 pub async fn list_records(
172 &self,
173 collection: &str,
174 limit: Option<u32>,
175 cursor: Option<&str>,
176 ) -> Result<(Vec<RecordEntry>, Option<String>)> {
177 let page = self.list_records_page(collection, limit, cursor).await?;
178 Ok((page.records, page.cursor))
179 }
180
181 pub(crate) async fn list_records_page(
184 &self,
185 collection: &str,
186 limit: Option<u32>,
187 cursor: Option<&str>,
188 ) -> Result<crate::atproto::ListRecordsResponse> {
189 let mut url = url::Url::parse(&self.url("com.atproto.repo.listRecords"))
190 .context("building the listRecords URL")?;
191 {
192 let mut query = url.query_pairs_mut();
193 query.append_pair("repo", &self.session.sub);
194 query.append_pair("collection", collection);
195 if let Some(limit) = limit {
196 query.append_pair("limit", &limit.to_string());
197 }
198 if let Some(cursor) = cursor {
199 query.append_pair("cursor", cursor);
200 }
201 }
202
203 let outcome = self
207 .send_raw(
208 url.as_str(),
209 DpopBody::Query,
210 "com.atproto.repo.listRecords",
211 )
212 .await?;
213 crate::atproto::parse_list_records(&outcome.body)
222 }
223
224 pub async fn list_all_records(&self, collection: &str) -> Result<Vec<RecordEntry>> {
230 self.list_all_records_within(
231 collection,
232 &mut crate::atproto::ByteBudget::new(crate::atproto::MAX_LIST_BYTES),
233 )
234 .await
235 }
236
237 pub(crate) async fn list_all_records_within(
239 &self,
240 collection: &str,
241 budget: &mut crate::atproto::ByteBudget,
242 ) -> Result<Vec<RecordEntry>> {
243 let mut out = Vec::new();
244 let max_bytes = budget.max();
245 let mut cursor: Option<String> = None;
246 let mut more_offered = false;
247
248 for _ in 0..MAX_LIST_PAGES {
249 let listed = self
250 .list_records_page(collection, Some(100), cursor.as_deref())
251 .await?;
252 crate::atproto::refuse_malformed(&listed, collection)?;
255 let (page, next) = (listed.records, listed.cursor);
256 let got = page.len();
257 if !budget.admit(&page) {
262 anyhow::bail!(
263 "listRecords for {collection} exceeded the {max_bytes}-byte cap \
264 ({} held, {} bytes charged) — refusing to accumulate further",
265 out.len(),
266 budget.used(),
267 );
268 }
269 crate::atproto::extend_bounded(&mut out, page, MAX_LIST_RECORDS, collection)?;
270 match next {
271 Some(next) if got > 0 && Some(&next) != cursor.as_ref() => {
281 cursor = Some(next);
282 more_offered = true;
283 }
284 _ => {
285 more_offered = false;
291 break;
292 }
293 }
294 }
295 if more_offered {
319 anyhow::bail!(
320 "listRecords for {collection} did not finish within {MAX_LIST_PAGES} pages \
321 ({} held, and the PDS still offered more) — refusing a short list",
322 out.len(),
323 );
324 }
325 Ok(out)
326 }
327
328 pub async fn create_record<T: crate::vetted::WritableRecord>(
330 &self,
331 collection: &str,
332 record: &T,
333 ) -> Result<WriteResult> {
334 self.write(
335 "com.atproto.repo.createRecord",
336 json!({ "repo": self.session.sub, "collection": collection, "record": record }),
337 )
338 .await
339 }
340
341 pub async fn put_record<T: crate::vetted::WritableRecord>(
343 &self,
344 collection: &str,
345 rkey: &str,
346 record: &T,
347 ) -> Result<WriteResult> {
348 self.write(
349 "com.atproto.repo.putRecord",
350 json!({
351 "repo": self.session.sub,
352 "collection": collection,
353 "rkey": rkey,
354 "record": record,
355 }),
356 )
357 .await
358 }
359
360 pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
361 let body = json!({ "repo": self.session.sub, "collection": collection, "rkey": rkey });
362 self.send(
363 &self.url("com.atproto.repo.deleteRecord"),
364 DpopBody::Json(serde_json::to_vec(&body)?),
365 "com.atproto.repo.deleteRecord",
366 )
367 .await
368 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
369 Ok(())
370 }
371
372 pub async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
377 crate::atproto::apply_writes_chunked(writes, |chunk| self.apply_writes_once(chunk)).await
378 }
379
380 async fn apply_writes_once(&self, writes: &[WriteOp]) -> Result<()> {
384 let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
385 let body = json!({ "repo": self.session.sub, "writes": ops });
386 self.send(
387 &self.url("com.atproto.repo.applyWrites"),
388 DpopBody::Json(serde_json::to_vec(&body)?),
389 "com.atproto.repo.applyWrites",
390 )
391 .await
392 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
393 Ok(())
394 }
395
396 async fn write(&self, nsid: &str, body: Value) -> Result<WriteResult> {
397 let value = self
398 .send(
399 &self.url(nsid),
400 DpopBody::Json(serde_json::to_vec(&body)?),
401 nsid,
402 )
403 .await?;
404 serde_json::from_value(value).with_context(|| format!("{nsid} returned no usable result"))
405 }
406}
407
408fn xrpc_error(body: &[u8], status: u16) -> String {
410 match Repo::error_fields(body) {
411 Some(detail) => format!("status {status} ({detail})"),
412 None => format!("status {status}"),
413 }
414}
415
416impl Repo<'_> {
424 async fn list_typed<T: serde::de::DeserializeOwned>(
431 &self,
432 collection: &str,
433 ) -> Result<Vec<(String, T)>> {
434 let records = self.list_all_records(collection).await?;
435 let mut out = Vec::with_capacity(records.len());
436 for record in records {
437 let rkey = record.rkey().unwrap_or_default().to_string();
438 match record.parse::<T>() {
439 Ok(value) => out.push((rkey, value)),
440 Err(err) => tracing::warn!(
441 collection,
442 uri = %record.uri,
443 error = %err,
444 "skipping unparseable record in collection"
445 ),
446 }
447 }
448 Ok(out)
449 }
450
451 pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
454 self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
455 }
456
457 pub async fn list_subscriptions_sorted(
459 &self,
460 ) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
461 let mut subs = self.list_subscriptions().await?;
462 subs.sort_by(crate::lexicon::sort::subscriptions);
463 Ok(subs)
464 }
465
466 pub async fn add_subscription(
469 &self,
470 sub: &crate::vetted::VettedSubscription,
471 ) -> Result<String> {
472 Ok(self
473 .create_record(crate::lexicon::nsid::SUBSCRIPTION, sub)
474 .await?
475 .into_rkey())
476 }
477
478 pub async fn remove_subscription(&self, rkey: &str) -> Result<()> {
479 self.delete_record(crate::lexicon::nsid::SUBSCRIPTION, rkey)
480 .await
481 }
482
483 pub async fn update_subscription(
488 &self,
489 rkey: &str,
490 sub: &crate::vetted::VettedSubscription,
491 ) -> Result<WriteResult> {
492 self.put_record(crate::lexicon::nsid::SUBSCRIPTION, rkey, sub)
493 .await
494 }
495
496 pub async fn add_subscriptions_bulk(
508 &self,
509 subs: &[crate::vetted::VettedSubscription],
510 ) -> Result<Vec<String>> {
511 let mut gen = crate::atproto::TidGenerator::new();
512 let mut rkeys = Vec::with_capacity(subs.len());
513 let mut writes = Vec::with_capacity(subs.len());
514 for sub in subs {
515 let rkey = gen.next();
516 writes.push(WriteOp::Create {
517 collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
518 rkey: Some(rkey.clone()),
519 value: serde_json::to_value(sub)?,
523 });
524 rkeys.push(rkey);
525 }
526 self.apply_writes(&writes).await?;
527 Ok(rkeys)
528 }
529
530 pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
533 self.list_typed(crate::lexicon::nsid::FOLDER).await
534 }
535
536 pub async fn list_folders_sorted(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
537 let mut folders = self.list_folders().await?;
538 folders.sort_by(crate::lexicon::sort::folders);
539 Ok(folders)
540 }
541
542 pub async fn add_folder(&self, folder: &crate::lexicon::Folder) -> Result<String> {
543 Ok(self
544 .create_record(crate::lexicon::nsid::FOLDER, folder)
545 .await?
546 .into_rkey())
547 }
548
549 pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
553 self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
554 }
555
556 pub async fn rename_folder(
559 &self,
560 rkey: &str,
561 folder: &crate::lexicon::Folder,
562 ) -> Result<WriteResult> {
563 self.put_record(crate::lexicon::nsid::FOLDER, rkey, folder)
564 .await
565 }
566
567 pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
570 self.list_typed(crate::lexicon::nsid::SAVED).await
571 }
572
573 pub async fn list_saved_sorted(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
575 let mut saved = self.list_saved().await?;
576 saved.sort_by(crate::lexicon::sort::saved);
577 Ok(saved)
578 }
579
580 pub async fn add_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<String> {
581 Ok(self
582 .create_record(crate::lexicon::nsid::SAVED, saved)
583 .await?
584 .into_rkey())
585 }
586
587 pub async fn remove_saved(&self, rkey: &str) -> Result<()> {
588 self.delete_record(crate::lexicon::nsid::SAVED, rkey).await
589 }
590
591 pub async fn list_read_states(&self) -> Result<Vec<(String, crate::lexicon::ReadState)>> {
594 self.list_typed(crate::lexicon::nsid::READ_STATE).await
595 }
596
597 pub async fn put_read_state(
599 &self,
600 rkey: &str,
601 state: &crate::lexicon::ReadState,
602 ) -> Result<()> {
603 self.put_record(crate::lexicon::nsid::READ_STATE, rkey, state)
604 .await?;
605 Ok(())
606 }
607
608 pub async fn flush_read_states(
622 &self,
623 cursors: &[(String, crate::lexicon::ReadState, bool)],
624 ) -> Result<()> {
625 if cursors.is_empty() {
626 return Ok(());
627 }
628 let writes = crate::atproto::read_state_write_ops(cursors)?;
629 self.apply_writes(&writes).await
630 }
631}
632
633#[cfg(test)]
634mod tests {
635 use super::*;
636
637 #[test]
650 fn read_state_writes_choose_create_or_update_per_cursor() {
651 let state = crate::lexicon::ReadState::new(
652 "https://example.com/feed",
653 Some("2026-01-01T00:00:00Z".to_string()),
654 "2026-01-01T00:00:00Z",
655 );
656 let cursors = vec![
657 ("existing".to_string(), state.clone(), true),
658 ("brand-new".to_string(), state.clone(), false),
659 ];
660
661 let ops = crate::atproto::read_state_write_ops(&cursors).expect("ops build");
662 assert_eq!(ops.len(), 2);
663
664 let rendered: Vec<Value> = ops.iter().map(|op| op.to_json()).collect();
665 assert_eq!(
666 rendered[0]["$type"], "com.atproto.repo.applyWrites#update",
667 "an existing record must be UPDATED, not re-created"
668 );
669 assert_eq!(
670 rendered[1]["$type"], "com.atproto.repo.applyWrites#create",
671 "a first flush must CREATE, or the whole atomic batch fails"
672 );
673 }
674
675 #[tokio::test]
680 async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
681 let (base, log) = crate::net::tests::serve_json_capturing(b"{}".to_vec()).await;
682 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
683 crate::net::test_host_override(
684 "bulk-pds.test",
685 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
686 );
687 let http = Client::new();
688 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
689 crate::store::init_schema(&pool).await.unwrap();
690 let key = SigningKey::generate("k");
691 let mut s = session();
692 s.aud = format!("http://bulk-pds.test:{port}");
693 let repo = repo(&http, &pool, &s, &key);
694 let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
695 .map(|i| {
696 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
697 format!("https://f{i}.example/feed.xml"),
698 "2026-07-12T00:00:00.000Z",
699 ))
700 })
701 .collect();
702
703 let rkeys = repo
704 .add_subscriptions_bulk(&subs)
705 .await
706 .expect("bulk write failed");
707
708 let sent = log.lock().unwrap().clone();
709 assert_eq!(
710 sent.len(),
711 1,
712 "expected one applyWrites request, got {sent:?}"
713 );
714 let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
715 .expect("request body is JSON");
716 let writes = body["writes"].as_array().expect("writes array");
717 assert_eq!(writes.len(), 3);
718 for (i, w) in writes.iter().enumerate() {
719 assert_eq!(w["collection"], crate::lexicon::nsid::SUBSCRIPTION);
720 assert_eq!(w["rkey"].as_str(), Some(rkeys[i].as_str()));
721 }
722 let mut sorted = rkeys.clone();
723 sorted.sort();
724 assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
725 }
726
727 fn session() -> OAuthSession {
728 OAuthSession {
729 sub: "did:plc:ewvi7nxzyoun6zhxrhs64oiz".into(),
730 issuer: "https://pds.example.com".into(),
731 aud: "https://pds.example.com".into(),
732 dpop_key_jwk: "{}".into(),
733 access_token: "tok".into(),
734 refresh_token: "ref".into(),
735 token_type: "DPoP".into(),
736 granted_scope: "atproto".into(),
737 expires_at: None,
738 }
739 }
740
741 fn repo<'a>(
742 http: &'a Client,
743 pool: &'a SqlitePool,
744 session: &'a OAuthSession,
745 key: &'a SigningKey,
746 ) -> Repo<'a> {
747 Repo {
748 http,
749 pool,
750 session,
751 key,
752 }
753 }
754
755 #[tokio::test]
758 async fn endpoints_are_built_from_the_sessions_audience() {
759 let http = Client::new();
760 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
761 let key = SigningKey::generate("k");
762 let mut s = session();
763 s.aud = "https://pds.example.com/".into();
764 let repo = repo(&http, &pool, &s, &key);
765 assert_eq!(
766 repo.url("com.atproto.repo.listRecords"),
767 "https://pds.example.com/xrpc/com.atproto.repo.listRecords",
768 "a trailing slash on the audience must not double the separator"
769 );
770 }
771
772 #[test]
775 fn an_xrpc_error_is_summarised_not_echoed() {
776 let body = br#"{"error":"InvalidRequest","message":"unknown collection"}"#;
777 let rendered = xrpc_error(body, 400);
778 assert!(rendered.contains("InvalidRequest"));
779 assert!(rendered.contains("unknown collection"));
780
781 let opaque = xrpc_error(br#"{"access_token":"SECRET"}"#, 500);
783 assert_eq!(opaque, "status 500");
784 assert!(!opaque.contains("SECRET"));
785 assert_eq!(xrpc_error(b"<html>oops</html>", 502), "status 502");
786 }
787
788 #[tokio::test]
791 async fn repo_calls_fail_closed_on_an_internal_pds() {
792 let http = Client::new();
793 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
794 super::super::store::init_schema(&pool).await.unwrap();
795 let key = SigningKey::generate("k");
796 let mut s = session();
797 s.aud = "http://127.0.0.1:2583".into();
798 let repo = repo(&http, &pool, &s, &key);
799
800 let err = repo
801 .list_records("app.feather.subscription", None, None)
802 .await
803 .expect_err("must refuse a loopback PDS");
804 assert!(
805 format!("{err:#}").contains("forbidden (internal) address"),
806 "failed for the wrong reason: {err:#}"
807 );
808 }
809
810 #[tokio::test]
820 async fn a_200_error_envelope_is_not_an_empty_repo() {
821 let base = crate::net::tests::serve_body(
822 br#"{"error":"InvalidRequest","message":"bad cursor"}"#.to_vec(),
823 )
824 .await;
825 let port: u16 = base
826 .trim_end_matches('/')
827 .rsplit(':')
828 .next()
829 .unwrap()
830 .parse()
831 .unwrap();
832 crate::net::test_host_override(
833 "envelope-pds.test",
834 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
835 );
836
837 let http = Client::new();
838 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
839 crate::store::init_schema(&pool).await.unwrap();
840 let key = SigningKey::generate("k");
841 let mut s = session();
842 s.aud = format!("http://envelope-pds.test:{port}");
843 let repo = repo(&http, &pool, &s, &key);
844
845 let err = repo
846 .list_records("app.feather.subscription", None, None)
847 .await
848 .expect_err("an error envelope was read as an empty page");
849 assert!(
850 format!("{err:#}").contains("InvalidRequest"),
851 "failed for the wrong reason: {err:#}"
852 );
853 }
854
855 #[tokio::test]
862 async fn the_live_walk_spends_its_budget_across_pages() {
863 let (bodies, per_page) = crate::atproto::tests::paged_bodies(3, 4096, false);
864 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
865 let port: u16 = base
866 .trim_end_matches('/')
867 .rsplit(':')
868 .next()
869 .unwrap()
870 .parse()
871 .unwrap();
872 crate::net::test_host_override(
873 "live-budget-pages.test",
874 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
875 );
876
877 let http = Client::new();
878 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
879 crate::store::init_schema(&pool).await.unwrap();
880 let key = SigningKey::generate("k");
881 let mut s = session();
882 s.aud = format!("http://live-budget-pages.test:{port}");
883 let repo = repo(&http, &pool, &s, &key);
884
885 let err = repo
886 .list_all_records_within(
887 "app.feather.subscription",
888 &mut crate::atproto::ByteBudget::new(per_page * 2),
889 )
890 .await
891 .expect_err("three pages cannot fit in a two-page budget");
892 let msg = format!("{err:#}");
893 assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
894 assert!(
895 msg.contains("2 held"),
896 "the live walk did not accumulate across pages: {msg}"
897 );
898 }
899
900 #[tokio::test]
904 async fn the_live_walk_that_runs_out_of_pages_refuses() {
905 let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
906 .map(|i| {
907 serde_json::json!({
908 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
909 "cursor": format!("p{}", i + 1),
910 })
911 .to_string()
912 .into_bytes()
913 })
914 .collect();
915 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
916 let port: u16 = base
917 .trim_end_matches('/')
918 .rsplit(':')
919 .next()
920 .unwrap()
921 .parse()
922 .unwrap();
923 crate::net::test_host_override(
924 "pages-exhausted-live.test",
925 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
926 );
927 let http = Client::new();
928 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
929 crate::store::init_schema(&pool).await.unwrap();
930 let key = SigningKey::generate("k");
931 let mut s = session();
932 s.aud = format!("http://pages-exhausted-live.test:{port}");
933 let repo = repo(&http, &pool, &s, &key);
934
935 let err = repo
936 .list_all_records("app.feather.subscription")
937 .await
938 .expect_err("a truncated list was returned as a complete one");
939 assert!(
940 format!("{err:#}").contains("did not finish"),
941 "failed for the wrong reason: {err:#}"
942 );
943 }
944
945 #[tokio::test]
948 async fn the_live_walk_refuses_a_page_with_a_malformed_record() {
949 let body = serde_json::json!({ "records": [
950 { "uri": "at://did:plc:x/c/3labGOOD", "value": {} },
951 { "cid": "bafy", "value": {} },
952 ]})
953 .to_string()
954 .into_bytes();
955 let base = crate::net::tests::serve_bodies_in_sequence(vec![body]).await;
956 let port: u16 = base
957 .trim_end_matches('/')
958 .rsplit(':')
959 .next()
960 .unwrap()
961 .parse()
962 .unwrap();
963 crate::net::test_host_override(
964 "malformed-live.test",
965 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
966 );
967 let http = Client::new();
968 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
969 crate::store::init_schema(&pool).await.unwrap();
970 let key = SigningKey::generate("k");
971 let mut s = session();
972 s.aud = format!("http://malformed-live.test:{port}");
973 let repo = repo(&http, &pool, &s, &key);
974 let err = repo
975 .list_all_records("app.feather.subscription")
976 .await
977 .expect_err("a page with a malformed record was accepted");
978 assert!(
979 err.downcast_ref::<crate::atproto::MalformedRecords>()
980 .is_some(),
981 "refused for the wrong reason: {err:#}"
982 );
983 }
984
985 #[tokio::test]
988 async fn the_live_walk_that_finishes_cleanly_returns_the_records() {
989 let mut bodies: Vec<Vec<u8>> = (0..3)
990 .map(|i| {
991 serde_json::json!({
992 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
993 "cursor": format!("p{}", i + 1),
994 })
995 .to_string()
996 .into_bytes()
997 })
998 .collect();
999 bodies.push(
1004 serde_json::json!({
1005 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1006 })
1007 .to_string()
1008 .into_bytes(),
1009 );
1010 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1011 let port: u16 = base
1012 .trim_end_matches('/')
1013 .rsplit(':')
1014 .next()
1015 .unwrap()
1016 .parse()
1017 .unwrap();
1018 crate::net::test_host_override(
1019 "clean-finish-live.test",
1020 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1021 );
1022 let http = Client::new();
1023 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1024 crate::store::init_schema(&pool).await.unwrap();
1025 let key = SigningKey::generate("k");
1026 let mut s = session();
1027 s.aud = format!("http://clean-finish-live.test:{port}");
1028 let repo = repo(&http, &pool, &s, &key);
1029
1030 let records = repo
1031 .list_all_records("app.feather.subscription")
1032 .await
1033 .expect("a walk that ran out of records is not a short list");
1034 assert_eq!(records.len(), 4);
1035 assert!(
1036 records.iter().any(|r| r.uri.ends_with("3labLAST")),
1037 "the LAST page's records were dropped: {:?}",
1038 records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
1039 );
1040 }
1041
1042 #[tokio::test]
1050 async fn a_live_walk_that_terminates_on_its_last_allowed_page_succeeds() {
1051 let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
1052 .map(|i| {
1053 serde_json::json!({
1054 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
1055 "cursor": format!("p{}", i + 1),
1056 })
1057 .to_string()
1058 .into_bytes()
1059 })
1060 .collect();
1061 bodies.push(
1062 serde_json::json!({
1063 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1064 })
1065 .to_string()
1066 .into_bytes(),
1067 );
1068 assert_eq!(bodies.len(), MAX_LIST_PAGES);
1069 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1070 let port: u16 = base
1071 .trim_end_matches('/')
1072 .rsplit(':')
1073 .next()
1074 .unwrap()
1075 .parse()
1076 .unwrap();
1077 crate::net::test_host_override(
1078 "last-allowed-page-live.test",
1079 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1080 );
1081 let http = Client::new();
1082 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1083 crate::store::init_schema(&pool).await.unwrap();
1084 let key = SigningKey::generate("k");
1085 let mut s = session();
1086 s.aud = format!("http://last-allowed-page-live.test:{port}");
1087 let repo = repo(&http, &pool, &s, &key);
1088
1089 let records = repo
1090 .list_all_records("app.feather.subscription")
1091 .await
1092 .expect("a walk that terminated inside its budget is not a short list");
1093 assert_eq!(
1094 records.len(),
1095 MAX_LIST_PAGES,
1096 "a walk that used its whole page budget and finished lost records",
1097 );
1098 }
1099
1100 #[test]
1112 fn an_oversized_error_body_is_not_parsed_for_its_reason() {
1113 let small = br#"{"error":"InvalidSwap","message":"record changed"}"#;
1114 assert_eq!(
1115 Repo::error_fields(small).as_deref(),
1116 Some("InvalidSwap: record changed"),
1117 "a real error body must still render its reason",
1118 );
1119
1120 let mut huge = String::from(r#"{"error":"InvalidSwap","pad":["#);
1121 while huge.len() < crate::oauth::MAX_ERROR_BODY + 1_024 {
1122 huge.push_str("{},");
1123 }
1124 huge.push_str("{}]}");
1125 assert!(huge.len() > crate::oauth::MAX_ERROR_BODY);
1126 assert_eq!(
1127 Repo::error_fields(huge.as_bytes()),
1128 None,
1129 "an oversized error body was deserialised to fish out one string",
1130 );
1131 }
1132
1133 #[tokio::test]
1142 async fn a_rejected_write_carries_the_status_and_error_name() {
1143 use axum::response::IntoResponse as _;
1144 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1145 let addr = listener.local_addr().unwrap();
1146 let host = format!("rejecting-{}.xrpc.test", addr.port());
1147 crate::net::test_host_override(&host, addr);
1148 let app = axum::Router::new().fallback(|| async {
1149 (
1150 axum::http::StatusCode::INTERNAL_SERVER_ERROR,
1151 axum::Json(
1152 json!({ "error": "InternalServerError", "message": "Internal Server Error" }),
1153 ),
1154 )
1155 .into_response()
1156 });
1157 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1158
1159 let http = Client::new();
1160 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1161 crate::store::init_schema(&pool).await.unwrap();
1162 let key = SigningKey::generate("k");
1163 let mut s = session();
1164 s.aud = format!("http://{host}:{}", addr.port());
1165 let repo = repo(&http, &pool, &s, &key);
1166
1167 let err = repo
1168 .apply_writes(&[WriteOp::Delete {
1169 collection: crate::lexicon::nsid::READ_STATE.into(),
1170 rkey: "rs-0".into(),
1171 }])
1172 .await
1173 .expect_err("a 500 is a failure");
1174 assert_eq!(
1175 err.to_string(),
1176 "com.atproto.repo.applyWrites failed: status 500 \
1177 (InternalServerError: Internal Server Error)",
1178 "the rendered message changed",
1179 );
1180 let xrpc = err
1181 .chain()
1182 .find_map(|cause| cause.downcast_ref::<crate::atproto::AtProtoError>());
1183 match xrpc {
1184 Some(crate::atproto::AtProtoError::Xrpc { status, error, .. }) => {
1185 assert_eq!(status.as_u16(), 500);
1186 assert_eq!(error, "InternalServerError");
1187 }
1188 other => panic!("no structured XRPC error in the chain: {other:?}"),
1189 }
1190 }
1191
1192 #[tokio::test]
1205 async fn the_live_write_path_refuses_a_node_explosion() {
1206 let mut body = String::from(r#"{"uri":"at://d/c/r","value":["#);
1207 for _ in 0..1_200_000 {
1208 body.push_str("{},");
1209 }
1210 body.push_str("{}]}");
1211 assert!(
1212 crate::atproto::count_structural_chars(body.as_bytes())
1213 > crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
1214 "the probe body is not over the cap, so this test proves nothing",
1215 );
1216 let base = crate::net::tests::serve_body(body.into_bytes()).await;
1217 let port: u16 = base
1218 .trim_end_matches('/')
1219 .rsplit(':')
1220 .next()
1221 .unwrap()
1222 .parse()
1223 .unwrap();
1224 crate::net::test_host_override(
1225 "write-explosion.test",
1226 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1227 );
1228 let http = Client::new();
1229 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1230 crate::store::init_schema(&pool).await.unwrap();
1231 let key = SigningKey::generate("k");
1232 let mut s = session();
1233 s.aud = format!("http://write-explosion.test:{port}");
1234 let repo = repo(&http, &pool, &s, &key);
1235
1236 let err = repo
1237 .delete_record("c", "r")
1238 .await
1239 .expect_err("a node explosion on the write path was parsed rather than refused");
1240 assert!(
1241 format!("{err:#}").contains("structural characters"),
1242 "failed for the wrong reason: {err:#}"
1243 );
1244 }
1245
1246 #[tokio::test]
1249 async fn the_live_write_path_accepts_an_ordinary_response() {
1250 let base =
1251 crate::net::tests::serve_body(br#"{"commit":{"cid":"bafy","rev":"3lab"}}"#.to_vec())
1252 .await;
1253 let port: u16 = base
1254 .trim_end_matches('/')
1255 .rsplit(':')
1256 .next()
1257 .unwrap()
1258 .parse()
1259 .unwrap();
1260 crate::net::test_host_override(
1261 "write-ordinary.test",
1262 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1263 );
1264 let http = Client::new();
1265 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1266 crate::store::init_schema(&pool).await.unwrap();
1267 let key = SigningKey::generate("k");
1268 let mut s = session();
1269 s.aud = format!("http://write-ordinary.test:{port}");
1270 let repo = repo(&http, &pool, &s, &key);
1271
1272 repo.delete_record("c", "r")
1273 .await
1274 .expect("an ordinary write response was refused");
1275 }
1276
1277 #[tokio::test]
1286 async fn the_live_walk_refuses_a_duplicated_records_key() {
1287 let base = crate::net::tests::serve_body(
1288 br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#.to_vec(),
1289 )
1290 .await;
1291 let port: u16 = base
1292 .trim_end_matches('/')
1293 .rsplit(':')
1294 .next()
1295 .unwrap()
1296 .parse()
1297 .unwrap();
1298 crate::net::test_host_override(
1299 "dup-records.test",
1300 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1301 );
1302
1303 let http = Client::new();
1304 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1305 crate::store::init_schema(&pool).await.unwrap();
1306 let key = SigningKey::generate("k");
1307 let mut s = session();
1308 s.aud = format!("http://dup-records.test:{port}");
1309 let repo = repo(&http, &pool, &s, &key);
1310
1311 let err = repo
1312 .list_records("app.feather.subscription", None, None)
1313 .await
1314 .expect_err("a duplicated records key was read as an empty page");
1315 assert!(
1316 format!("{err:#}").contains("duplicate"),
1317 "failed for the wrong reason: {err:#}"
1318 );
1319 }
1320
1321 #[tokio::test]
1328 async fn an_empty_200_body_is_not_an_empty_repo() {
1329 let base = crate::net::tests::serve_body(Vec::new()).await;
1330 let port: u16 = base
1331 .trim_end_matches('/')
1332 .rsplit(':')
1333 .next()
1334 .unwrap()
1335 .parse()
1336 .unwrap();
1337 crate::net::test_host_override(
1338 "empty-body.test",
1339 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1340 );
1341 let http = Client::new();
1342 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1343 crate::store::init_schema(&pool).await.unwrap();
1344 let key = SigningKey::generate("k");
1345 let mut s = session();
1346 s.aud = format!("http://empty-body.test:{port}");
1347 let repo = repo(&http, &pool, &s, &key);
1348
1349 let err = repo
1350 .list_records("app.feather.subscription", None, None)
1351 .await
1352 .expect_err("an empty body was read as an empty repo");
1353 assert!(
1354 format!("{err:#}").contains("no records"),
1355 "failed for the wrong reason: {err:#}"
1356 );
1357 }
1358
1359 #[tokio::test]
1363 async fn a_200_error_envelope_is_not_a_successful_write() {
1364 let base = crate::net::tests::serve_body(
1365 br#"{"error":"InvalidRequest","message":"nope"}"#.to_vec(),
1366 )
1367 .await;
1368 let port: u16 = base
1369 .trim_end_matches('/')
1370 .rsplit(':')
1371 .next()
1372 .unwrap()
1373 .parse()
1374 .unwrap();
1375 crate::net::test_host_override(
1376 "envelope-write.test",
1377 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1378 );
1379 let http = Client::new();
1380 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1381 crate::store::init_schema(&pool).await.unwrap();
1382 let key = SigningKey::generate("k");
1383 let mut s = session();
1384 s.aud = format!("http://envelope-write.test:{port}");
1385 let repo = repo(&http, &pool, &s, &key);
1386
1387 let err = repo
1388 .delete_record("app.feather.subscription", "rk1")
1389 .await
1390 .expect_err("a failed delete was reported as success");
1391 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1392
1393 let err = repo
1394 .apply_writes(&[crate::atproto::WriteOp::Delete {
1395 collection: "app.feather.subscription".to_string(),
1396 rkey: "rk1".to_string(),
1397 }])
1398 .await
1399 .expect_err("a failed batch was reported as success");
1400 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1401 }
1402
1403 #[tokio::test]
1406 async fn an_empty_batch_is_not_sent() {
1407 let http = Client::new();
1408 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1409 let key = SigningKey::generate("k");
1410 let mut s = session();
1411 s.aud = "http://127.0.0.1:2583".into();
1413 let repo = repo(&http, &pool, &s, &key);
1414 assert!(repo.apply_writes(&[]).await.is_ok());
1415 }
1416
1417 async fn strict_pds(
1426 fail_call: Option<usize>,
1427 ) -> (
1428 OAuthSession,
1429 SqlitePool,
1430 crate::atproto::tests::ApplyWritesLog,
1431 ) {
1432 let (base, log) = crate::atproto::tests::serve_apply_writes(fail_call).await;
1433 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
1434 let host = format!("chunk-oauth-{port}.test");
1435 crate::net::test_host_override(&host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
1436 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1437 crate::store::init_schema(&pool).await.unwrap();
1438 let mut s = session();
1439 s.aud = format!("http://{host}:{port}");
1440 (s, pool, log)
1441 }
1442
1443 fn subs(n: usize) -> Vec<crate::vetted::VettedSubscription> {
1444 (0..n)
1445 .map(|i| {
1446 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
1447 format!("https://f{i}.example/feed.xml"),
1448 "2026-07-12T00:00:00.000Z",
1449 ))
1450 })
1451 .collect()
1452 }
1453
1454 #[tokio::test]
1458 async fn bulk_subscribe_of_201_is_two_calls_in_order() {
1459 let (s, pool, log) = strict_pds(None).await;
1460 let (http, key) = (Client::new(), SigningKey::generate("k"));
1461 let rkeys = repo(&http, &pool, &s, &key)
1462 .add_subscriptions_bulk(&subs(201))
1463 .await
1464 .expect("a 201-feed import must succeed against a PDS that caps at 200");
1465 assert_eq!(crate::atproto::tests::call_sizes(&log), vec![200, 1]);
1466 assert_eq!(
1467 crate::atproto::tests::sent_rkeys(&log),
1468 rkeys,
1469 "every feed, once, in input order"
1470 );
1471 }
1472
1473 #[tokio::test]
1475 async fn bulk_subscribe_splits_at_200_and_not_before() {
1476 for (n, want) in [(500, vec![200, 200, 100]), (200, vec![200])] {
1477 let (s, pool, log) = strict_pds(None).await;
1478 let (http, key) = (Client::new(), SigningKey::generate("k"));
1479 repo(&http, &pool, &s, &key)
1480 .add_subscriptions_bulk(&subs(n))
1481 .await
1482 .expect("bulk write");
1483 assert_eq!(crate::atproto::tests::call_sizes(&log), want, "{n} feeds");
1484 }
1485 }
1486
1487 #[tokio::test]
1489 async fn bulk_subscribe_stops_at_the_first_failed_chunk() {
1490 let (s, pool, log) = strict_pds(Some(2)).await;
1491 let (http, key) = (Client::new(), SigningKey::generate("k"));
1492 let err = repo(&http, &pool, &s, &key)
1493 .add_subscriptions_bulk(&subs(500))
1494 .await
1495 .expect_err("a failed chunk must fail the call");
1496 assert_eq!(
1497 crate::atproto::tests::call_sizes(&log),
1498 vec![200, 200],
1499 "chunk 3 must NOT be sent"
1500 );
1501 assert!(format!("{err:#}").contains("boom"), "{err:#}");
1502 }
1503
1504 #[tokio::test]
1508 async fn read_state_flush_splits_on_bytes_under_200_ops() {
1509 let (s, pool, log) = strict_pds(None).await;
1510 let (http, key) = (Client::new(), SigningKey::generate("k"));
1511 let cursors: Vec<(String, crate::lexicon::ReadState, bool)> = (0..10)
1512 .map(|i| {
1513 let mut state = crate::lexicon::ReadState::new(
1514 format!("https://f{i}.example/feed.xml"),
1515 None,
1516 "2026-07-12T00:00:00.000Z",
1517 );
1518 state.read_ids = (0..crate::lexicon::ReadState::MAX_IDS)
1519 .map(|j| format!("https://f{i}.example/posts/{j:04}/an-entry-permalink"))
1520 .collect();
1521 (format!("rk{i:04}"), state, false)
1522 })
1523 .collect();
1524 repo(&http, &pool, &s, &key)
1525 .flush_read_states(&cursors)
1526 .await
1527 .expect("a byte-heavy flush must succeed in chunks");
1528 let sizes = crate::atproto::tests::call_sizes(&log);
1529 assert!(sizes.len() > 1, "one call for ~500 KB: {sizes:?}");
1530 let want: Vec<String> = cursors.iter().map(|(rkey, _, _)| rkey.clone()).collect();
1531 assert_eq!(crate::atproto::tests::sent_rkeys(&log), want);
1532 }
1533}