1use anyhow::{bail, 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 bail!(
86 "{nsid} failed: {}",
87 xrpc_error(&outcome.body, outcome.status)
88 );
89 }
90 Ok(outcome)
91 }
92
93 async fn send(&self, url: &str, body: DpopBody<'_>, nsid: &str) -> Result<Value> {
99 let outcome = self.send_raw(url, body, nsid).await?;
100 if outcome.body.is_empty() {
102 return Ok(Value::Null);
103 }
104 outcome.json()
105 }
106
107 fn error_fields(body: &[u8]) -> Option<String> {
120 if !super::error_body_worth_parsing(body) {
121 return None;
122 }
123 let value: Value = serde_json::from_slice(body).ok()?;
124 let kind = value.get("error").and_then(Value::as_str)?;
125 match value.get("message").and_then(Value::as_str) {
126 Some(message) => Some(format!("{kind}: {message}")),
127 None => Some(kind.to_string()),
128 }
129 }
130
131 pub async fn list_records(
133 &self,
134 collection: &str,
135 limit: Option<u32>,
136 cursor: Option<&str>,
137 ) -> Result<(Vec<RecordEntry>, Option<String>)> {
138 let mut url = url::Url::parse(&self.url("com.atproto.repo.listRecords"))
139 .context("building the listRecords URL")?;
140 {
141 let mut query = url.query_pairs_mut();
142 query.append_pair("repo", &self.session.sub);
143 query.append_pair("collection", collection);
144 if let Some(limit) = limit {
145 query.append_pair("limit", &limit.to_string());
146 }
147 if let Some(cursor) = cursor {
148 query.append_pair("cursor", cursor);
149 }
150 }
151
152 let outcome = self
156 .send_raw(
157 url.as_str(),
158 DpopBody::Query,
159 "com.atproto.repo.listRecords",
160 )
161 .await?;
162 let page = crate::atproto::parse_list_records(&outcome.body)?;
171 let (records, cursor) = (page.records, page.cursor);
172 Ok((records, cursor))
173 }
174
175 pub async fn list_all_records(&self, collection: &str) -> Result<Vec<RecordEntry>> {
181 self.list_all_records_within(
182 collection,
183 &mut crate::atproto::ByteBudget::new(crate::atproto::MAX_LIST_BYTES),
184 )
185 .await
186 }
187
188 pub(crate) async fn list_all_records_within(
190 &self,
191 collection: &str,
192 budget: &mut crate::atproto::ByteBudget,
193 ) -> Result<Vec<RecordEntry>> {
194 let mut out = Vec::new();
195 let max_bytes = budget.max();
196 let mut cursor: Option<String> = None;
197 let mut more_offered = false;
198
199 for _ in 0..MAX_LIST_PAGES {
200 let (page, next) = self
201 .list_records(collection, Some(100), cursor.as_deref())
202 .await?;
203 let got = page.len();
204 if !budget.admit(&page) {
209 anyhow::bail!(
210 "listRecords for {collection} exceeded the {max_bytes}-byte cap \
211 ({} held, {} bytes charged) — refusing to accumulate further",
212 out.len(),
213 budget.used(),
214 );
215 }
216 crate::atproto::extend_bounded(&mut out, page, MAX_LIST_RECORDS, collection)?;
217 match next {
218 Some(next) if got > 0 && Some(&next) != cursor.as_ref() => {
228 cursor = Some(next);
229 more_offered = true;
230 }
231 _ => {
232 more_offered = false;
238 break;
239 }
240 }
241 }
242 if more_offered {
266 anyhow::bail!(
267 "listRecords for {collection} did not finish within {MAX_LIST_PAGES} pages \
268 ({} held, and the PDS still offered more) — refusing a short list",
269 out.len(),
270 );
271 }
272 Ok(out)
273 }
274
275 pub async fn create_record<T: crate::vetted::WritableRecord>(
277 &self,
278 collection: &str,
279 record: &T,
280 ) -> Result<WriteResult> {
281 self.write(
282 "com.atproto.repo.createRecord",
283 json!({ "repo": self.session.sub, "collection": collection, "record": record }),
284 )
285 .await
286 }
287
288 pub async fn put_record<T: crate::vetted::WritableRecord>(
290 &self,
291 collection: &str,
292 rkey: &str,
293 record: &T,
294 ) -> Result<WriteResult> {
295 self.write(
296 "com.atproto.repo.putRecord",
297 json!({
298 "repo": self.session.sub,
299 "collection": collection,
300 "rkey": rkey,
301 "record": record,
302 }),
303 )
304 .await
305 }
306
307 pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
308 let body = json!({ "repo": self.session.sub, "collection": collection, "rkey": rkey });
309 self.send(
310 &self.url("com.atproto.repo.deleteRecord"),
311 DpopBody::Json(serde_json::to_vec(&body)?),
312 "com.atproto.repo.deleteRecord",
313 )
314 .await
315 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
316 Ok(())
317 }
318
319 pub async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
321 if writes.is_empty() {
322 return Ok(());
323 }
324 let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
325 let body = json!({ "repo": self.session.sub, "writes": ops });
326 self.send(
327 &self.url("com.atproto.repo.applyWrites"),
328 DpopBody::Json(serde_json::to_vec(&body)?),
329 "com.atproto.repo.applyWrites",
330 )
331 .await
332 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
333 Ok(())
334 }
335
336 async fn write(&self, nsid: &str, body: Value) -> Result<WriteResult> {
337 let value = self
338 .send(
339 &self.url(nsid),
340 DpopBody::Json(serde_json::to_vec(&body)?),
341 nsid,
342 )
343 .await?;
344 serde_json::from_value(value).with_context(|| format!("{nsid} returned no usable result"))
345 }
346}
347
348fn xrpc_error(body: &[u8], status: u16) -> String {
350 match Repo::error_fields(body) {
351 Some(detail) => format!("status {status} ({detail})"),
352 None => format!("status {status}"),
353 }
354}
355
356impl Repo<'_> {
364 async fn list_typed<T: serde::de::DeserializeOwned>(
371 &self,
372 collection: &str,
373 ) -> Result<Vec<(String, T)>> {
374 let records = self.list_all_records(collection).await?;
375 let mut out = Vec::with_capacity(records.len());
376 for record in records {
377 let rkey = record.rkey().unwrap_or_default().to_string();
378 match record.parse::<T>() {
379 Ok(value) => out.push((rkey, value)),
380 Err(err) => tracing::warn!(
381 collection,
382 uri = %record.uri,
383 error = %err,
384 "skipping unparseable record in collection"
385 ),
386 }
387 }
388 Ok(out)
389 }
390
391 pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
394 self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
395 }
396
397 pub async fn list_subscriptions_sorted(
399 &self,
400 ) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
401 let mut subs = self.list_subscriptions().await?;
402 subs.sort_by(crate::lexicon::sort::subscriptions);
403 Ok(subs)
404 }
405
406 pub async fn add_subscription(
409 &self,
410 sub: &crate::vetted::VettedSubscription,
411 ) -> Result<String> {
412 Ok(self
413 .create_record(crate::lexicon::nsid::SUBSCRIPTION, sub)
414 .await?
415 .into_rkey())
416 }
417
418 pub async fn remove_subscription(&self, rkey: &str) -> Result<()> {
419 self.delete_record(crate::lexicon::nsid::SUBSCRIPTION, rkey)
420 .await
421 }
422
423 pub async fn update_subscription(
428 &self,
429 rkey: &str,
430 sub: &crate::vetted::VettedSubscription,
431 ) -> Result<WriteResult> {
432 self.put_record(crate::lexicon::nsid::SUBSCRIPTION, rkey, sub)
433 .await
434 }
435
436 pub async fn add_subscriptions_bulk(
445 &self,
446 subs: &[crate::vetted::VettedSubscription],
447 ) -> Result<Vec<String>> {
448 let mut gen = crate::atproto::TidGenerator::new();
449 let mut rkeys = Vec::with_capacity(subs.len());
450 let mut writes = Vec::with_capacity(subs.len());
451 for sub in subs {
452 let rkey = gen.next();
453 writes.push(WriteOp::Create {
454 collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
455 rkey: Some(rkey.clone()),
456 value: serde_json::to_value(sub)?,
460 });
461 rkeys.push(rkey);
462 }
463 self.apply_writes(&writes).await?;
464 Ok(rkeys)
465 }
466
467 pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
470 self.list_typed(crate::lexicon::nsid::FOLDER).await
471 }
472
473 pub async fn list_folders_sorted(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
474 let mut folders = self.list_folders().await?;
475 folders.sort_by(crate::lexicon::sort::folders);
476 Ok(folders)
477 }
478
479 pub async fn add_folder(&self, folder: &crate::lexicon::Folder) -> Result<String> {
480 Ok(self
481 .create_record(crate::lexicon::nsid::FOLDER, folder)
482 .await?
483 .into_rkey())
484 }
485
486 pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
490 self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
491 }
492
493 pub async fn rename_folder(
496 &self,
497 rkey: &str,
498 folder: &crate::lexicon::Folder,
499 ) -> Result<WriteResult> {
500 self.put_record(crate::lexicon::nsid::FOLDER, rkey, folder)
501 .await
502 }
503
504 pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
507 self.list_typed(crate::lexicon::nsid::SAVED).await
508 }
509
510 pub async fn list_saved_sorted(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
512 let mut saved = self.list_saved().await?;
513 saved.sort_by(crate::lexicon::sort::saved);
514 Ok(saved)
515 }
516
517 pub async fn add_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<String> {
518 Ok(self
519 .create_record(crate::lexicon::nsid::SAVED, saved)
520 .await?
521 .into_rkey())
522 }
523
524 pub async fn remove_saved(&self, rkey: &str) -> Result<()> {
525 self.delete_record(crate::lexicon::nsid::SAVED, rkey).await
526 }
527
528 pub async fn list_read_states(&self) -> Result<Vec<(String, crate::lexicon::ReadState)>> {
531 self.list_typed(crate::lexicon::nsid::READ_STATE).await
532 }
533
534 pub async fn put_read_state(
536 &self,
537 rkey: &str,
538 state: &crate::lexicon::ReadState,
539 ) -> Result<()> {
540 self.put_record(crate::lexicon::nsid::READ_STATE, rkey, state)
541 .await?;
542 Ok(())
543 }
544
545 pub async fn flush_read_states(
556 &self,
557 cursors: &[(String, crate::lexicon::ReadState, bool)],
558 ) -> Result<()> {
559 if cursors.is_empty() {
560 return Ok(());
561 }
562 let writes = crate::atproto::read_state_write_ops(cursors)?;
563 self.apply_writes(&writes).await
564 }
565}
566
567#[cfg(test)]
568mod tests {
569 use super::*;
570
571 #[test]
584 fn read_state_writes_choose_create_or_update_per_cursor() {
585 let state = crate::lexicon::ReadState::new(
586 "https://example.com/feed",
587 Some("2026-01-01T00:00:00Z".to_string()),
588 "2026-01-01T00:00:00Z",
589 );
590 let cursors = vec![
591 ("existing".to_string(), state.clone(), true),
592 ("brand-new".to_string(), state.clone(), false),
593 ];
594
595 let ops = crate::atproto::read_state_write_ops(&cursors).expect("ops build");
596 assert_eq!(ops.len(), 2);
597
598 let rendered: Vec<Value> = ops.iter().map(|op| op.to_json()).collect();
599 assert_eq!(
600 rendered[0]["$type"], "com.atproto.repo.applyWrites#update",
601 "an existing record must be UPDATED, not re-created"
602 );
603 assert_eq!(
604 rendered[1]["$type"], "com.atproto.repo.applyWrites#create",
605 "a first flush must CREATE, or the whole atomic batch fails"
606 );
607 }
608
609 #[tokio::test]
614 async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
615 let (base, log) = crate::net::tests::serve_json_capturing(b"{}".to_vec()).await;
616 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
617 crate::net::test_host_override(
618 "bulk-pds.test",
619 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
620 );
621 let http = Client::new();
622 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
623 crate::store::init_schema(&pool).await.unwrap();
624 let key = SigningKey::generate("k");
625 let mut s = session();
626 s.aud = format!("http://bulk-pds.test:{port}");
627 let repo = repo(&http, &pool, &s, &key);
628 let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
629 .map(|i| {
630 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
631 format!("https://f{i}.example/feed.xml"),
632 "2026-07-12T00:00:00.000Z",
633 ))
634 })
635 .collect();
636
637 let rkeys = repo
638 .add_subscriptions_bulk(&subs)
639 .await
640 .expect("bulk write failed");
641
642 let sent = log.lock().unwrap().clone();
643 assert_eq!(
644 sent.len(),
645 1,
646 "expected one applyWrites request, got {sent:?}"
647 );
648 let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
649 .expect("request body is JSON");
650 let writes = body["writes"].as_array().expect("writes array");
651 assert_eq!(writes.len(), 3);
652 for (i, w) in writes.iter().enumerate() {
653 assert_eq!(w["collection"], crate::lexicon::nsid::SUBSCRIPTION);
654 assert_eq!(w["rkey"].as_str(), Some(rkeys[i].as_str()));
655 }
656 let mut sorted = rkeys.clone();
657 sorted.sort();
658 assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
659 }
660
661 fn session() -> OAuthSession {
662 OAuthSession {
663 sub: "did:plc:ewvi7nxzyoun6zhxrhs64oiz".into(),
664 issuer: "https://pds.example.com".into(),
665 aud: "https://pds.example.com".into(),
666 dpop_key_jwk: "{}".into(),
667 access_token: "tok".into(),
668 refresh_token: "ref".into(),
669 token_type: "DPoP".into(),
670 granted_scope: "atproto".into(),
671 expires_at: None,
672 }
673 }
674
675 fn repo<'a>(
676 http: &'a Client,
677 pool: &'a SqlitePool,
678 session: &'a OAuthSession,
679 key: &'a SigningKey,
680 ) -> Repo<'a> {
681 Repo {
682 http,
683 pool,
684 session,
685 key,
686 }
687 }
688
689 #[tokio::test]
692 async fn endpoints_are_built_from_the_sessions_audience() {
693 let http = Client::new();
694 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
695 let key = SigningKey::generate("k");
696 let mut s = session();
697 s.aud = "https://pds.example.com/".into();
698 let repo = repo(&http, &pool, &s, &key);
699 assert_eq!(
700 repo.url("com.atproto.repo.listRecords"),
701 "https://pds.example.com/xrpc/com.atproto.repo.listRecords",
702 "a trailing slash on the audience must not double the separator"
703 );
704 }
705
706 #[test]
709 fn an_xrpc_error_is_summarised_not_echoed() {
710 let body = br#"{"error":"InvalidRequest","message":"unknown collection"}"#;
711 let rendered = xrpc_error(body, 400);
712 assert!(rendered.contains("InvalidRequest"));
713 assert!(rendered.contains("unknown collection"));
714
715 let opaque = xrpc_error(br#"{"access_token":"SECRET"}"#, 500);
717 assert_eq!(opaque, "status 500");
718 assert!(!opaque.contains("SECRET"));
719 assert_eq!(xrpc_error(b"<html>oops</html>", 502), "status 502");
720 }
721
722 #[tokio::test]
725 async fn repo_calls_fail_closed_on_an_internal_pds() {
726 let http = Client::new();
727 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
728 super::super::store::init_schema(&pool).await.unwrap();
729 let key = SigningKey::generate("k");
730 let mut s = session();
731 s.aud = "http://127.0.0.1:2583".into();
732 let repo = repo(&http, &pool, &s, &key);
733
734 let err = repo
735 .list_records("app.feather.subscription", None, None)
736 .await
737 .expect_err("must refuse a loopback PDS");
738 assert!(
739 format!("{err:#}").contains("forbidden (internal) address"),
740 "failed for the wrong reason: {err:#}"
741 );
742 }
743
744 #[tokio::test]
754 async fn a_200_error_envelope_is_not_an_empty_repo() {
755 let base = crate::net::tests::serve_body(
756 br#"{"error":"InvalidRequest","message":"bad cursor"}"#.to_vec(),
757 )
758 .await;
759 let port: u16 = base
760 .trim_end_matches('/')
761 .rsplit(':')
762 .next()
763 .unwrap()
764 .parse()
765 .unwrap();
766 crate::net::test_host_override(
767 "envelope-pds.test",
768 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
769 );
770
771 let http = Client::new();
772 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
773 crate::store::init_schema(&pool).await.unwrap();
774 let key = SigningKey::generate("k");
775 let mut s = session();
776 s.aud = format!("http://envelope-pds.test:{port}");
777 let repo = repo(&http, &pool, &s, &key);
778
779 let err = repo
780 .list_records("app.feather.subscription", None, None)
781 .await
782 .expect_err("an error envelope was read as an empty page");
783 assert!(
784 format!("{err:#}").contains("InvalidRequest"),
785 "failed for the wrong reason: {err:#}"
786 );
787 }
788
789 #[tokio::test]
796 async fn the_live_walk_spends_its_budget_across_pages() {
797 let (bodies, per_page) = crate::atproto::tests::paged_bodies(3, 4096, false);
798 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
799 let port: u16 = base
800 .trim_end_matches('/')
801 .rsplit(':')
802 .next()
803 .unwrap()
804 .parse()
805 .unwrap();
806 crate::net::test_host_override(
807 "live-budget-pages.test",
808 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
809 );
810
811 let http = Client::new();
812 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
813 crate::store::init_schema(&pool).await.unwrap();
814 let key = SigningKey::generate("k");
815 let mut s = session();
816 s.aud = format!("http://live-budget-pages.test:{port}");
817 let repo = repo(&http, &pool, &s, &key);
818
819 let err = repo
820 .list_all_records_within(
821 "app.feather.subscription",
822 &mut crate::atproto::ByteBudget::new(per_page * 2),
823 )
824 .await
825 .expect_err("three pages cannot fit in a two-page budget");
826 let msg = format!("{err:#}");
827 assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
828 assert!(
829 msg.contains("2 held"),
830 "the live walk did not accumulate across pages: {msg}"
831 );
832 }
833
834 #[tokio::test]
838 async fn the_live_walk_that_runs_out_of_pages_refuses() {
839 let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
840 .map(|i| {
841 serde_json::json!({
842 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
843 "cursor": format!("p{}", i + 1),
844 })
845 .to_string()
846 .into_bytes()
847 })
848 .collect();
849 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
850 let port: u16 = base
851 .trim_end_matches('/')
852 .rsplit(':')
853 .next()
854 .unwrap()
855 .parse()
856 .unwrap();
857 crate::net::test_host_override(
858 "pages-exhausted-live.test",
859 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
860 );
861 let http = Client::new();
862 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
863 crate::store::init_schema(&pool).await.unwrap();
864 let key = SigningKey::generate("k");
865 let mut s = session();
866 s.aud = format!("http://pages-exhausted-live.test:{port}");
867 let repo = repo(&http, &pool, &s, &key);
868
869 let err = repo
870 .list_all_records("app.feather.subscription")
871 .await
872 .expect_err("a truncated list was returned as a complete one");
873 assert!(
874 format!("{err:#}").contains("did not finish"),
875 "failed for the wrong reason: {err:#}"
876 );
877 }
878
879 #[tokio::test]
882 async fn the_live_walk_that_finishes_cleanly_returns_the_records() {
883 let mut bodies: Vec<Vec<u8>> = (0..3)
884 .map(|i| {
885 serde_json::json!({
886 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
887 "cursor": format!("p{}", i + 1),
888 })
889 .to_string()
890 .into_bytes()
891 })
892 .collect();
893 bodies.push(
898 serde_json::json!({
899 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
900 })
901 .to_string()
902 .into_bytes(),
903 );
904 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
905 let port: u16 = base
906 .trim_end_matches('/')
907 .rsplit(':')
908 .next()
909 .unwrap()
910 .parse()
911 .unwrap();
912 crate::net::test_host_override(
913 "clean-finish-live.test",
914 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
915 );
916 let http = Client::new();
917 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
918 crate::store::init_schema(&pool).await.unwrap();
919 let key = SigningKey::generate("k");
920 let mut s = session();
921 s.aud = format!("http://clean-finish-live.test:{port}");
922 let repo = repo(&http, &pool, &s, &key);
923
924 let records = repo
925 .list_all_records("app.feather.subscription")
926 .await
927 .expect("a walk that ran out of records is not a short list");
928 assert_eq!(records.len(), 4);
929 assert!(
930 records.iter().any(|r| r.uri.ends_with("3labLAST")),
931 "the LAST page's records were dropped: {:?}",
932 records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
933 );
934 }
935
936 #[tokio::test]
944 async fn a_live_walk_that_terminates_on_its_last_allowed_page_succeeds() {
945 let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
946 .map(|i| {
947 serde_json::json!({
948 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
949 "cursor": format!("p{}", i + 1),
950 })
951 .to_string()
952 .into_bytes()
953 })
954 .collect();
955 bodies.push(
956 serde_json::json!({
957 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
958 })
959 .to_string()
960 .into_bytes(),
961 );
962 assert_eq!(bodies.len(), MAX_LIST_PAGES);
963 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
964 let port: u16 = base
965 .trim_end_matches('/')
966 .rsplit(':')
967 .next()
968 .unwrap()
969 .parse()
970 .unwrap();
971 crate::net::test_host_override(
972 "last-allowed-page-live.test",
973 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
974 );
975 let http = Client::new();
976 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
977 crate::store::init_schema(&pool).await.unwrap();
978 let key = SigningKey::generate("k");
979 let mut s = session();
980 s.aud = format!("http://last-allowed-page-live.test:{port}");
981 let repo = repo(&http, &pool, &s, &key);
982
983 let records = repo
984 .list_all_records("app.feather.subscription")
985 .await
986 .expect("a walk that terminated inside its budget is not a short list");
987 assert_eq!(
988 records.len(),
989 MAX_LIST_PAGES,
990 "a walk that used its whole page budget and finished lost records",
991 );
992 }
993
994 #[test]
1006 fn an_oversized_error_body_is_not_parsed_for_its_reason() {
1007 let small = br#"{"error":"InvalidSwap","message":"record changed"}"#;
1008 assert_eq!(
1009 Repo::error_fields(small).as_deref(),
1010 Some("InvalidSwap: record changed"),
1011 "a real error body must still render its reason",
1012 );
1013
1014 let mut huge = String::from(r#"{"error":"InvalidSwap","pad":["#);
1015 while huge.len() < crate::oauth::MAX_ERROR_BODY + 1_024 {
1016 huge.push_str("{},");
1017 }
1018 huge.push_str("{}]}");
1019 assert!(huge.len() > crate::oauth::MAX_ERROR_BODY);
1020 assert_eq!(
1021 Repo::error_fields(huge.as_bytes()),
1022 None,
1023 "an oversized error body was deserialised to fish out one string",
1024 );
1025 }
1026
1027 #[tokio::test]
1040 async fn the_live_write_path_refuses_a_node_explosion() {
1041 let mut body = String::from(r#"{"uri":"at://d/c/r","value":["#);
1042 for _ in 0..1_200_000 {
1043 body.push_str("{},");
1044 }
1045 body.push_str("{}]}");
1046 assert!(
1047 crate::atproto::count_structural_chars(body.as_bytes())
1048 > crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
1049 "the probe body is not over the cap, so this test proves nothing",
1050 );
1051 let base = crate::net::tests::serve_body(body.into_bytes()).await;
1052 let port: u16 = base
1053 .trim_end_matches('/')
1054 .rsplit(':')
1055 .next()
1056 .unwrap()
1057 .parse()
1058 .unwrap();
1059 crate::net::test_host_override(
1060 "write-explosion.test",
1061 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1062 );
1063 let http = Client::new();
1064 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1065 crate::store::init_schema(&pool).await.unwrap();
1066 let key = SigningKey::generate("k");
1067 let mut s = session();
1068 s.aud = format!("http://write-explosion.test:{port}");
1069 let repo = repo(&http, &pool, &s, &key);
1070
1071 let err = repo
1072 .delete_record("c", "r")
1073 .await
1074 .expect_err("a node explosion on the write path was parsed rather than refused");
1075 assert!(
1076 format!("{err:#}").contains("structural characters"),
1077 "failed for the wrong reason: {err:#}"
1078 );
1079 }
1080
1081 #[tokio::test]
1084 async fn the_live_write_path_accepts_an_ordinary_response() {
1085 let base =
1086 crate::net::tests::serve_body(br#"{"commit":{"cid":"bafy","rev":"3lab"}}"#.to_vec())
1087 .await;
1088 let port: u16 = base
1089 .trim_end_matches('/')
1090 .rsplit(':')
1091 .next()
1092 .unwrap()
1093 .parse()
1094 .unwrap();
1095 crate::net::test_host_override(
1096 "write-ordinary.test",
1097 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1098 );
1099 let http = Client::new();
1100 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1101 crate::store::init_schema(&pool).await.unwrap();
1102 let key = SigningKey::generate("k");
1103 let mut s = session();
1104 s.aud = format!("http://write-ordinary.test:{port}");
1105 let repo = repo(&http, &pool, &s, &key);
1106
1107 repo.delete_record("c", "r")
1108 .await
1109 .expect("an ordinary write response was refused");
1110 }
1111
1112 #[tokio::test]
1121 async fn the_live_walk_refuses_a_duplicated_records_key() {
1122 let base = crate::net::tests::serve_body(
1123 br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#.to_vec(),
1124 )
1125 .await;
1126 let port: u16 = base
1127 .trim_end_matches('/')
1128 .rsplit(':')
1129 .next()
1130 .unwrap()
1131 .parse()
1132 .unwrap();
1133 crate::net::test_host_override(
1134 "dup-records.test",
1135 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1136 );
1137
1138 let http = Client::new();
1139 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1140 crate::store::init_schema(&pool).await.unwrap();
1141 let key = SigningKey::generate("k");
1142 let mut s = session();
1143 s.aud = format!("http://dup-records.test:{port}");
1144 let repo = repo(&http, &pool, &s, &key);
1145
1146 let err = repo
1147 .list_records("app.feather.subscription", None, None)
1148 .await
1149 .expect_err("a duplicated records key was read as an empty page");
1150 assert!(
1151 format!("{err:#}").contains("duplicate"),
1152 "failed for the wrong reason: {err:#}"
1153 );
1154 }
1155
1156 #[tokio::test]
1163 async fn an_empty_200_body_is_not_an_empty_repo() {
1164 let base = crate::net::tests::serve_body(Vec::new()).await;
1165 let port: u16 = base
1166 .trim_end_matches('/')
1167 .rsplit(':')
1168 .next()
1169 .unwrap()
1170 .parse()
1171 .unwrap();
1172 crate::net::test_host_override(
1173 "empty-body.test",
1174 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1175 );
1176 let http = Client::new();
1177 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1178 crate::store::init_schema(&pool).await.unwrap();
1179 let key = SigningKey::generate("k");
1180 let mut s = session();
1181 s.aud = format!("http://empty-body.test:{port}");
1182 let repo = repo(&http, &pool, &s, &key);
1183
1184 let err = repo
1185 .list_records("app.feather.subscription", None, None)
1186 .await
1187 .expect_err("an empty body was read as an empty repo");
1188 assert!(
1189 format!("{err:#}").contains("no records"),
1190 "failed for the wrong reason: {err:#}"
1191 );
1192 }
1193
1194 #[tokio::test]
1198 async fn a_200_error_envelope_is_not_a_successful_write() {
1199 let base = crate::net::tests::serve_body(
1200 br#"{"error":"InvalidRequest","message":"nope"}"#.to_vec(),
1201 )
1202 .await;
1203 let port: u16 = base
1204 .trim_end_matches('/')
1205 .rsplit(':')
1206 .next()
1207 .unwrap()
1208 .parse()
1209 .unwrap();
1210 crate::net::test_host_override(
1211 "envelope-write.test",
1212 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1213 );
1214 let http = Client::new();
1215 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1216 crate::store::init_schema(&pool).await.unwrap();
1217 let key = SigningKey::generate("k");
1218 let mut s = session();
1219 s.aud = format!("http://envelope-write.test:{port}");
1220 let repo = repo(&http, &pool, &s, &key);
1221
1222 let err = repo
1223 .delete_record("app.feather.subscription", "rk1")
1224 .await
1225 .expect_err("a failed delete was reported as success");
1226 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1227
1228 let err = repo
1229 .apply_writes(&[crate::atproto::WriteOp::Delete {
1230 collection: "app.feather.subscription".to_string(),
1231 rkey: "rk1".to_string(),
1232 }])
1233 .await
1234 .expect_err("a failed batch was reported as success");
1235 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1236 }
1237
1238 #[tokio::test]
1241 async fn an_empty_batch_is_not_sent() {
1242 let http = Client::new();
1243 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1244 let key = SigningKey::generate("k");
1245 let mut s = session();
1246 s.aud = "http://127.0.0.1:2583".into();
1248 let repo = repo(&http, &pool, &s, &key);
1249 assert!(repo.apply_writes(&[]).await.is_ok());
1250 }
1251}