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>> {
231 self.list_all_records_within(
232 collection,
233 &mut crate::atproto::ByteBudget::new(crate::atproto::MAX_LIST_BYTES),
234 )
235 .await
236 }
237
238 pub(crate) async fn list_all_records_within(
240 &self,
241 collection: &str,
242 budget: &mut crate::atproto::ByteBudget,
243 ) -> Result<Vec<RecordEntry>> {
244 let mut out = Vec::new();
245 let max_bytes = budget.max();
246 let mut cursor: Option<String> = None;
247 let mut more_offered = false;
248
249 for _ in 0..MAX_LIST_PAGES {
250 let listed = self
251 .list_records_page(collection, Some(100), cursor.as_deref())
252 .await?;
253 crate::atproto::refuse_malformed(&listed, collection)?;
256 let (page, next) = (listed.records, listed.cursor);
257 let got = page.len();
258 if !budget.admit(&page) {
263 return Err(crate::atproto::ListingTooLarge::Bytes {
264 collection: collection.to_string(),
265 max_bytes,
266 held: out.len(),
267 charged: budget.used(),
268 }
269 .into());
270 }
271 crate::atproto::extend_bounded(&mut out, page, MAX_LIST_RECORDS, collection)?;
272 match next {
273 Some(next) if got > 0 && Some(&next) == cursor.as_ref() => {
285 anyhow::bail!(
286 "listRecords for {collection} returned a repeated cursor on a \
287 non-empty page ({} held) — refusing a list the server did not finish",
288 out.len(),
289 );
290 }
291 Some(next) if got > 0 => {
292 cursor = Some(next);
293 more_offered = true;
294 }
295 _ => {
296 more_offered = false;
302 break;
303 }
304 }
305 }
306 if more_offered {
330 return Err(crate::atproto::ListingTooLarge::Pages {
331 collection: collection.to_string(),
332 pages: MAX_LIST_PAGES,
333 held: out.len(),
334 }
335 .into());
336 }
337 Ok(out)
338 }
339
340 pub async fn create_record<T: crate::vetted::WritableRecord>(
342 &self,
343 collection: &str,
344 record: &T,
345 ) -> Result<WriteResult> {
346 self.write(
347 "com.atproto.repo.createRecord",
348 json!({ "repo": self.session.sub, "collection": collection, "record": record }),
349 )
350 .await
351 }
352
353 pub async fn put_record<T: crate::vetted::WritableRecord>(
360 &self,
361 collection: &str,
362 rkey: &str,
363 record: &T,
364 swap_record: Option<&str>,
365 ) -> Result<WriteResult> {
366 let mut body = json!({
367 "repo": self.session.sub,
368 "collection": collection,
369 "rkey": rkey,
370 "record": record,
371 });
372 if let Some(cid) = swap_record {
373 body["swapRecord"] = json!(cid);
374 }
375 self.write("com.atproto.repo.putRecord", body).await
376 }
377
378 pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
379 let body = json!({ "repo": self.session.sub, "collection": collection, "rkey": rkey });
380 self.send(
381 &self.url("com.atproto.repo.deleteRecord"),
382 DpopBody::Json(serde_json::to_vec(&body)?),
383 "com.atproto.repo.deleteRecord",
384 )
385 .await
386 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
387 Ok(())
388 }
389
390 pub async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
395 crate::atproto::apply_writes_chunked(writes, |chunk| self.apply_writes_once(chunk)).await
396 }
397
398 async fn apply_writes_once(&self, writes: &[WriteOp]) -> Result<()> {
402 let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
403 let body = json!({ "repo": self.session.sub, "writes": ops });
404 self.send(
405 &self.url("com.atproto.repo.applyWrites"),
406 DpopBody::Json(serde_json::to_vec(&body)?),
407 "com.atproto.repo.applyWrites",
408 )
409 .await
410 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
411 Ok(())
412 }
413
414 async fn write(&self, nsid: &str, body: Value) -> Result<WriteResult> {
415 let value = self
416 .send(
417 &self.url(nsid),
418 DpopBody::Json(serde_json::to_vec(&body)?),
419 nsid,
420 )
421 .await?;
422 serde_json::from_value(value).with_context(|| format!("{nsid} returned no usable result"))
423 }
424}
425
426fn xrpc_error(body: &[u8], status: u16) -> String {
428 match Repo::error_fields(body) {
429 Some(detail) => format!("status {status} ({detail})"),
430 None => format!("status {status}"),
431 }
432}
433
434impl Repo<'_> {
442 async fn list_typed<T: serde::de::DeserializeOwned>(
449 &self,
450 collection: &str,
451 ) -> Result<Vec<(String, T)>> {
452 Ok(self
453 .list_typed_with_cids(collection)
454 .await?
455 .into_iter()
456 .map(|(rkey, _cid, value)| (rkey, value))
457 .collect())
458 }
459
460 async fn list_typed_with_cids<T: serde::de::DeserializeOwned>(
464 &self,
465 collection: &str,
466 ) -> Result<Vec<(String, Option<String>, T)>> {
467 let records = self.list_all_records(collection).await?;
468 let mut out = Vec::with_capacity(records.len());
469 for record in records {
470 let rkey = record.rkey().unwrap_or_default().to_string();
471 match record.parse::<T>() {
472 Ok(value) => out.push((rkey, record.cid, value)),
473 Err(err) => tracing::warn!(
474 collection,
475 uri = %record.uri,
476 error = %err,
477 "skipping unparseable record in collection"
478 ),
479 }
480 }
481 Ok(out)
482 }
483
484 pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
487 self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
488 }
489
490 pub async fn list_subscriptions_with_cids(
493 &self,
494 ) -> Result<Vec<(String, Option<String>, crate::lexicon::Subscription)>> {
495 self.list_typed_with_cids(crate::lexicon::nsid::SUBSCRIPTION)
496 .await
497 }
498
499 pub async fn list_subscriptions_sorted(
501 &self,
502 ) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
503 let mut subs = self.list_subscriptions().await?;
504 subs.sort_by(crate::lexicon::sort::subscriptions);
505 Ok(subs)
506 }
507
508 pub async fn add_subscription(
511 &self,
512 sub: &crate::vetted::VettedSubscription,
513 ) -> Result<String> {
514 Ok(self
515 .create_record(crate::lexicon::nsid::SUBSCRIPTION, sub)
516 .await?
517 .into_rkey())
518 }
519
520 pub async fn remove_subscription(&self, rkey: &str) -> Result<()> {
521 self.delete_record(crate::lexicon::nsid::SUBSCRIPTION, rkey)
522 .await
523 }
524
525 pub async fn update_subscription(
530 &self,
531 rkey: &str,
532 sub: &crate::vetted::VettedSubscription,
533 swap_record: Option<&str>,
534 ) -> Result<WriteResult> {
535 self.put_record(crate::lexicon::nsid::SUBSCRIPTION, rkey, sub, swap_record)
536 .await
537 }
538
539 pub async fn add_subscriptions_bulk(
551 &self,
552 subs: &[crate::vetted::VettedSubscription],
553 ) -> Result<Vec<String>> {
554 let mut gen = crate::atproto::TidGenerator::new();
555 let mut rkeys = Vec::with_capacity(subs.len());
556 let mut writes = Vec::with_capacity(subs.len());
557 for sub in subs {
558 let rkey = gen.next();
559 writes.push(WriteOp::Create {
560 collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
561 rkey: Some(rkey.clone()),
562 value: serde_json::to_value(sub)?,
566 });
567 rkeys.push(rkey);
568 }
569 self.apply_writes(&writes).await?;
570 Ok(rkeys)
571 }
572
573 pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
576 self.list_typed(crate::lexicon::nsid::FOLDER).await
577 }
578
579 pub async fn list_folders_with_cids(
581 &self,
582 ) -> Result<Vec<(String, Option<String>, crate::lexicon::Folder)>> {
583 self.list_typed_with_cids(crate::lexicon::nsid::FOLDER)
584 .await
585 }
586
587 pub async fn list_folders_sorted(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
588 let mut folders = self.list_folders().await?;
589 folders.sort_by(crate::lexicon::sort::folders);
590 Ok(folders)
591 }
592
593 pub async fn add_folder(&self, folder: &crate::lexicon::Folder) -> Result<String> {
594 Ok(self
595 .create_record(crate::lexicon::nsid::FOLDER, folder)
596 .await?
597 .into_rkey())
598 }
599
600 pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
604 self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
605 }
606
607 pub async fn rename_folder(
612 &self,
613 rkey: &str,
614 folder: &crate::lexicon::Folder,
615 swap_record: Option<&str>,
616 ) -> Result<WriteResult> {
617 self.put_record(crate::lexicon::nsid::FOLDER, rkey, folder, swap_record)
618 .await
619 }
620
621 pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
624 self.list_typed(crate::lexicon::nsid::SAVED).await
625 }
626
627 pub async fn list_saved_sorted(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
629 let mut saved = self.list_saved().await?;
630 saved.sort_by(crate::lexicon::sort::saved);
631 Ok(saved)
632 }
633
634 pub async fn add_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<String> {
635 Ok(self
636 .create_record(crate::lexicon::nsid::SAVED, saved)
637 .await?
638 .into_rkey())
639 }
640
641 pub async fn remove_saved(&self, rkey: &str) -> Result<()> {
642 self.delete_record(crate::lexicon::nsid::SAVED, rkey).await
643 }
644
645 pub async fn list_read_states(&self) -> Result<Vec<(String, crate::lexicon::ReadState)>> {
648 self.list_typed(crate::lexicon::nsid::READ_STATE).await
649 }
650
651 pub async fn put_read_state(
653 &self,
654 rkey: &str,
655 state: &crate::lexicon::ReadState,
656 ) -> Result<()> {
657 self.put_record(crate::lexicon::nsid::READ_STATE, rkey, state, None)
658 .await?;
659 Ok(())
660 }
661
662 pub async fn flush_read_states(
676 &self,
677 cursors: &[(String, crate::lexicon::ReadState, bool)],
678 ) -> Result<()> {
679 if cursors.is_empty() {
680 return Ok(());
681 }
682 let writes = crate::atproto::read_state_write_ops(cursors)?;
683 self.apply_writes(&writes).await
684 }
685}
686
687#[cfg(test)]
688mod tests {
689 use super::*;
690
691 #[test]
704 fn read_state_writes_choose_create_or_update_per_cursor() {
705 let state = crate::lexicon::ReadState::new(
706 "https://example.com/feed",
707 Some("2026-01-01T00:00:00Z".to_string()),
708 "2026-01-01T00:00:00Z",
709 );
710 let cursors = vec![
711 ("existing".to_string(), state.clone(), true),
712 ("brand-new".to_string(), state.clone(), false),
713 ];
714
715 let ops = crate::atproto::read_state_write_ops(&cursors).expect("ops build");
716 assert_eq!(ops.len(), 2);
717
718 let rendered: Vec<Value> = ops.iter().map(|op| op.to_json()).collect();
719 assert_eq!(
720 rendered[0]["$type"], "com.atproto.repo.applyWrites#update",
721 "an existing record must be UPDATED, not re-created"
722 );
723 assert_eq!(
724 rendered[1]["$type"], "com.atproto.repo.applyWrites#create",
725 "a first flush must CREATE, or the whole atomic batch fails"
726 );
727 }
728
729 #[tokio::test]
734 async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
735 let (base, log) = crate::net::tests::serve_json_capturing(b"{}".to_vec()).await;
736 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
737 crate::net::test_host_override(
738 "bulk-pds.test",
739 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
740 );
741 let http = Client::new();
742 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
743 crate::store::init_schema(&pool).await.unwrap();
744 let key = SigningKey::generate("k");
745 let mut s = session();
746 s.aud = format!("http://bulk-pds.test:{port}");
747 let repo = repo(&http, &pool, &s, &key);
748 let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
749 .map(|i| {
750 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
751 format!("https://f{i}.example/feed.xml"),
752 "2026-07-12T00:00:00.000Z",
753 ))
754 })
755 .collect();
756
757 let rkeys = repo
758 .add_subscriptions_bulk(&subs)
759 .await
760 .expect("bulk write failed");
761
762 let sent = log.lock().unwrap().clone();
763 assert_eq!(
764 sent.len(),
765 1,
766 "expected one applyWrites request, got {sent:?}"
767 );
768 let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
769 .expect("request body is JSON");
770 let writes = body["writes"].as_array().expect("writes array");
771 assert_eq!(writes.len(), 3);
772 for (i, w) in writes.iter().enumerate() {
773 assert_eq!(w["collection"], crate::lexicon::nsid::SUBSCRIPTION);
774 assert_eq!(w["rkey"].as_str(), Some(rkeys[i].as_str()));
775 }
776 let mut sorted = rkeys.clone();
777 sorted.sort();
778 assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
779 }
780
781 fn session() -> OAuthSession {
782 OAuthSession {
783 sub: "did:plc:ewvi7nxzyoun6zhxrhs64oiz".into(),
784 issuer: "https://pds.example.com".into(),
785 aud: "https://pds.example.com".into(),
786 dpop_key_jwk: "{}".into(),
787 access_token: "tok".into(),
788 refresh_token: "ref".into(),
789 token_type: "DPoP".into(),
790 granted_scope: "atproto".into(),
791 expires_at: None,
792 }
793 }
794
795 fn repo<'a>(
796 http: &'a Client,
797 pool: &'a SqlitePool,
798 session: &'a OAuthSession,
799 key: &'a SigningKey,
800 ) -> Repo<'a> {
801 Repo {
802 http,
803 pool,
804 session,
805 key,
806 }
807 }
808
809 #[tokio::test]
812 async fn endpoints_are_built_from_the_sessions_audience() {
813 let http = Client::new();
814 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
815 let key = SigningKey::generate("k");
816 let mut s = session();
817 s.aud = "https://pds.example.com/".into();
818 let repo = repo(&http, &pool, &s, &key);
819 assert_eq!(
820 repo.url("com.atproto.repo.listRecords"),
821 "https://pds.example.com/xrpc/com.atproto.repo.listRecords",
822 "a trailing slash on the audience must not double the separator"
823 );
824 }
825
826 #[test]
829 fn an_xrpc_error_is_summarised_not_echoed() {
830 let body = br#"{"error":"InvalidRequest","message":"unknown collection"}"#;
831 let rendered = xrpc_error(body, 400);
832 assert!(rendered.contains("InvalidRequest"));
833 assert!(rendered.contains("unknown collection"));
834
835 let opaque = xrpc_error(br#"{"access_token":"SECRET"}"#, 500);
837 assert_eq!(opaque, "status 500");
838 assert!(!opaque.contains("SECRET"));
839 assert_eq!(xrpc_error(b"<html>oops</html>", 502), "status 502");
840 }
841
842 #[tokio::test]
845 async fn repo_calls_fail_closed_on_an_internal_pds() {
846 let http = Client::new();
847 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
848 super::super::store::init_schema(&pool).await.unwrap();
849 let key = SigningKey::generate("k");
850 let mut s = session();
851 s.aud = "http://127.0.0.1:2583".into();
852 let repo = repo(&http, &pool, &s, &key);
853
854 let err = repo
855 .list_records("app.feather.subscription", None, None)
856 .await
857 .expect_err("must refuse a loopback PDS");
858 assert!(
859 format!("{err:#}").contains("forbidden (internal) address"),
860 "failed for the wrong reason: {err:#}"
861 );
862 }
863
864 #[tokio::test]
874 async fn a_200_error_envelope_is_not_an_empty_repo() {
875 let base = crate::net::tests::serve_body(
876 br#"{"error":"InvalidRequest","message":"bad cursor"}"#.to_vec(),
877 )
878 .await;
879 let port: u16 = base
880 .trim_end_matches('/')
881 .rsplit(':')
882 .next()
883 .unwrap()
884 .parse()
885 .unwrap();
886 crate::net::test_host_override(
887 "envelope-pds.test",
888 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
889 );
890
891 let http = Client::new();
892 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
893 crate::store::init_schema(&pool).await.unwrap();
894 let key = SigningKey::generate("k");
895 let mut s = session();
896 s.aud = format!("http://envelope-pds.test:{port}");
897 let repo = repo(&http, &pool, &s, &key);
898
899 let err = repo
900 .list_records("app.feather.subscription", None, None)
901 .await
902 .expect_err("an error envelope was read as an empty page");
903 assert!(
904 format!("{err:#}").contains("InvalidRequest"),
905 "failed for the wrong reason: {err:#}"
906 );
907 }
908
909 #[tokio::test]
916 async fn the_live_walk_spends_its_budget_across_pages() {
917 let (bodies, per_page) = crate::atproto::tests::paged_bodies(3, 4096, false);
918 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
919 let port: u16 = base
920 .trim_end_matches('/')
921 .rsplit(':')
922 .next()
923 .unwrap()
924 .parse()
925 .unwrap();
926 crate::net::test_host_override(
927 "live-budget-pages.test",
928 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
929 );
930
931 let http = Client::new();
932 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
933 crate::store::init_schema(&pool).await.unwrap();
934 let key = SigningKey::generate("k");
935 let mut s = session();
936 s.aud = format!("http://live-budget-pages.test:{port}");
937 let repo = repo(&http, &pool, &s, &key);
938
939 let err = repo
940 .list_all_records_within(
941 "app.feather.subscription",
942 &mut crate::atproto::ByteBudget::new(per_page * 2),
943 )
944 .await
945 .expect_err("three pages cannot fit in a two-page budget");
946 let msg = format!("{err:#}");
947 assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
948 assert!(
949 msg.contains("2 held"),
950 "the live walk did not accumulate across pages: {msg}"
951 );
952 assert_eq!(
953 crate::feed::publication_failure_kind(&err),
954 crate::feed::FailureKind::Body,
955 "{msg}"
956 );
957 }
958
959 #[tokio::test]
963 async fn the_live_walk_that_runs_out_of_pages_refuses() {
964 let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
965 .map(|i| {
966 serde_json::json!({
967 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
968 "cursor": format!("p{}", i + 1),
969 })
970 .to_string()
971 .into_bytes()
972 })
973 .collect();
974 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
975 let port: u16 = base
976 .trim_end_matches('/')
977 .rsplit(':')
978 .next()
979 .unwrap()
980 .parse()
981 .unwrap();
982 crate::net::test_host_override(
983 "pages-exhausted-live.test",
984 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
985 );
986 let http = Client::new();
987 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
988 crate::store::init_schema(&pool).await.unwrap();
989 let key = SigningKey::generate("k");
990 let mut s = session();
991 s.aud = format!("http://pages-exhausted-live.test:{port}");
992 let repo = repo(&http, &pool, &s, &key);
993
994 let err = repo
995 .list_all_records("app.feather.subscription")
996 .await
997 .expect_err("a truncated list was returned as a complete one");
998 assert!(
999 format!("{err:#}").contains("did not finish"),
1000 "failed for the wrong reason: {err:#}"
1001 );
1002 assert_eq!(
1003 crate::feed::publication_failure_kind(&err),
1004 crate::feed::FailureKind::Body,
1005 "{err:#}"
1006 );
1007 }
1008
1009 #[tokio::test]
1012 async fn the_live_walk_refuses_a_page_with_a_malformed_record() {
1013 let body = serde_json::json!({ "records": [
1014 { "uri": "at://did:plc:x/c/3labGOOD", "value": {} },
1015 { "cid": "bafy", "value": {} },
1016 ]})
1017 .to_string()
1018 .into_bytes();
1019 let base = crate::net::tests::serve_bodies_in_sequence(vec![body]).await;
1020 let port: u16 = base
1021 .trim_end_matches('/')
1022 .rsplit(':')
1023 .next()
1024 .unwrap()
1025 .parse()
1026 .unwrap();
1027 crate::net::test_host_override(
1028 "malformed-live.test",
1029 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1030 );
1031 let http = Client::new();
1032 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1033 crate::store::init_schema(&pool).await.unwrap();
1034 let key = SigningKey::generate("k");
1035 let mut s = session();
1036 s.aud = format!("http://malformed-live.test:{port}");
1037 let repo = repo(&http, &pool, &s, &key);
1038 let err = repo
1039 .list_all_records("app.feather.subscription")
1040 .await
1041 .expect_err("a page with a malformed record was accepted");
1042 assert!(
1043 err.downcast_ref::<crate::atproto::MalformedRecords>()
1044 .is_some(),
1045 "refused for the wrong reason: {err:#}"
1046 );
1047 }
1048
1049 #[tokio::test]
1052 async fn the_live_walk_that_finishes_cleanly_returns_the_records() {
1053 let mut bodies: Vec<Vec<u8>> = (0..3)
1054 .map(|i| {
1055 serde_json::json!({
1056 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
1057 "cursor": format!("p{}", i + 1),
1058 })
1059 .to_string()
1060 .into_bytes()
1061 })
1062 .collect();
1063 bodies.push(
1068 serde_json::json!({
1069 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1070 })
1071 .to_string()
1072 .into_bytes(),
1073 );
1074 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1075 let port: u16 = base
1076 .trim_end_matches('/')
1077 .rsplit(':')
1078 .next()
1079 .unwrap()
1080 .parse()
1081 .unwrap();
1082 crate::net::test_host_override(
1083 "clean-finish-live.test",
1084 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1085 );
1086 let http = Client::new();
1087 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1088 crate::store::init_schema(&pool).await.unwrap();
1089 let key = SigningKey::generate("k");
1090 let mut s = session();
1091 s.aud = format!("http://clean-finish-live.test:{port}");
1092 let repo = repo(&http, &pool, &s, &key);
1093
1094 let records = repo
1095 .list_all_records("app.feather.subscription")
1096 .await
1097 .expect("a walk that ran out of records is not a short list");
1098 assert_eq!(records.len(), 4);
1099 assert!(
1100 records.iter().any(|r| r.uri.ends_with("3labLAST")),
1101 "the LAST page's records were dropped: {:?}",
1102 records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
1103 );
1104 }
1105
1106 #[tokio::test]
1114 async fn a_live_walk_that_terminates_on_its_last_allowed_page_succeeds() {
1115 let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
1116 .map(|i| {
1117 serde_json::json!({
1118 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
1119 "cursor": format!("p{}", i + 1),
1120 })
1121 .to_string()
1122 .into_bytes()
1123 })
1124 .collect();
1125 bodies.push(
1126 serde_json::json!({
1127 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1128 })
1129 .to_string()
1130 .into_bytes(),
1131 );
1132 assert_eq!(bodies.len(), MAX_LIST_PAGES);
1133 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1134 let port: u16 = base
1135 .trim_end_matches('/')
1136 .rsplit(':')
1137 .next()
1138 .unwrap()
1139 .parse()
1140 .unwrap();
1141 crate::net::test_host_override(
1142 "last-allowed-page-live.test",
1143 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1144 );
1145 let http = Client::new();
1146 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1147 crate::store::init_schema(&pool).await.unwrap();
1148 let key = SigningKey::generate("k");
1149 let mut s = session();
1150 s.aud = format!("http://last-allowed-page-live.test:{port}");
1151 let repo = repo(&http, &pool, &s, &key);
1152
1153 let records = repo
1154 .list_all_records("app.feather.subscription")
1155 .await
1156 .expect("a walk that terminated inside its budget is not a short list");
1157 assert_eq!(
1158 records.len(),
1159 MAX_LIST_PAGES,
1160 "a walk that used its whole page budget and finished lost records",
1161 );
1162 }
1163
1164 #[tokio::test]
1170 async fn the_live_walk_refuses_a_repeated_cursor() {
1171 let body = serde_json::json!({
1172 "records": [{ "uri": "at://did:plc:x/c/3labONE", "value": {} }],
1173 "cursor": "same-every-time",
1174 })
1175 .to_string()
1176 .into_bytes();
1177 let base = crate::net::tests::serve_body(body).await;
1178 let port: u16 = base
1179 .trim_end_matches('/')
1180 .rsplit(':')
1181 .next()
1182 .unwrap()
1183 .parse()
1184 .unwrap();
1185 crate::net::test_host_override(
1186 "repeated-cursor-live.test",
1187 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1188 );
1189 let http = Client::new();
1190 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1191 crate::store::init_schema(&pool).await.unwrap();
1192 let key = SigningKey::generate("k");
1193 let mut s = session();
1194 s.aud = format!("http://repeated-cursor-live.test:{port}");
1195 let repo = repo(&http, &pool, &s, &key);
1196
1197 let err = repo
1198 .list_all_records("app.feather.subscription")
1199 .await
1200 .expect_err("a repeated cursor ended the walk with a short list");
1201 assert!(
1202 format!("{err:#}").contains("repeated cursor"),
1203 "refused for the wrong reason: {err:#}"
1204 );
1205 }
1206
1207 #[test]
1219 fn an_oversized_error_body_is_not_parsed_for_its_reason() {
1220 let small = br#"{"error":"InvalidSwap","message":"record changed"}"#;
1221 assert_eq!(
1222 Repo::error_fields(small).as_deref(),
1223 Some("InvalidSwap: record changed"),
1224 "a real error body must still render its reason",
1225 );
1226
1227 let mut huge = String::from(r#"{"error":"InvalidSwap","pad":["#);
1228 while huge.len() < crate::oauth::MAX_ERROR_BODY + 1_024 {
1229 huge.push_str("{},");
1230 }
1231 huge.push_str("{}]}");
1232 assert!(huge.len() > crate::oauth::MAX_ERROR_BODY);
1233 assert_eq!(
1234 Repo::error_fields(huge.as_bytes()),
1235 None,
1236 "an oversized error body was deserialised to fish out one string",
1237 );
1238 }
1239
1240 #[tokio::test]
1249 async fn a_rejected_write_carries_the_status_and_error_name() {
1250 use axum::response::IntoResponse as _;
1251 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1252 let addr = listener.local_addr().unwrap();
1253 let host = format!("rejecting-{}.xrpc.test", addr.port());
1254 crate::net::test_host_override(&host, addr);
1255 let app = axum::Router::new().fallback(|| async {
1256 (
1257 axum::http::StatusCode::INTERNAL_SERVER_ERROR,
1258 axum::Json(
1259 json!({ "error": "InternalServerError", "message": "Internal Server Error" }),
1260 ),
1261 )
1262 .into_response()
1263 });
1264 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1265
1266 let http = Client::new();
1267 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1268 crate::store::init_schema(&pool).await.unwrap();
1269 let key = SigningKey::generate("k");
1270 let mut s = session();
1271 s.aud = format!("http://{host}:{}", addr.port());
1272 let repo = repo(&http, &pool, &s, &key);
1273
1274 let err = repo
1275 .apply_writes(&[WriteOp::Delete {
1276 collection: crate::lexicon::nsid::READ_STATE.into(),
1277 rkey: "rs-0".into(),
1278 }])
1279 .await
1280 .expect_err("a 500 is a failure");
1281 assert_eq!(
1282 err.to_string(),
1283 "com.atproto.repo.applyWrites failed: status 500 \
1284 (InternalServerError: Internal Server Error)",
1285 "the rendered message changed",
1286 );
1287 let xrpc = err
1288 .chain()
1289 .find_map(|cause| cause.downcast_ref::<crate::atproto::AtProtoError>());
1290 match xrpc {
1291 Some(crate::atproto::AtProtoError::Xrpc { status, error, .. }) => {
1292 assert_eq!(status.as_u16(), 500);
1293 assert_eq!(error, "InternalServerError");
1294 }
1295 other => panic!("no structured XRPC error in the chain: {other:?}"),
1296 }
1297 }
1298
1299 #[tokio::test]
1312 async fn the_live_write_path_refuses_a_node_explosion() {
1313 let mut body = String::from(r#"{"uri":"at://d/c/r","value":["#);
1314 for _ in 0..1_200_000 {
1315 body.push_str("{},");
1316 }
1317 body.push_str("{}]}");
1318 assert!(
1319 crate::atproto::count_structural_chars(body.as_bytes())
1320 > crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
1321 "the probe body is not over the cap, so this test proves nothing",
1322 );
1323 let base = crate::net::tests::serve_body(body.into_bytes()).await;
1324 let port: u16 = base
1325 .trim_end_matches('/')
1326 .rsplit(':')
1327 .next()
1328 .unwrap()
1329 .parse()
1330 .unwrap();
1331 crate::net::test_host_override(
1332 "write-explosion.test",
1333 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1334 );
1335 let http = Client::new();
1336 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1337 crate::store::init_schema(&pool).await.unwrap();
1338 let key = SigningKey::generate("k");
1339 let mut s = session();
1340 s.aud = format!("http://write-explosion.test:{port}");
1341 let repo = repo(&http, &pool, &s, &key);
1342
1343 let err = repo
1344 .delete_record("c", "r")
1345 .await
1346 .expect_err("a node explosion on the write path was parsed rather than refused");
1347 assert!(
1348 format!("{err:#}").contains("structural characters"),
1349 "failed for the wrong reason: {err:#}"
1350 );
1351 }
1352
1353 #[tokio::test]
1356 async fn the_live_write_path_accepts_an_ordinary_response() {
1357 let base =
1358 crate::net::tests::serve_body(br#"{"commit":{"cid":"bafy","rev":"3lab"}}"#.to_vec())
1359 .await;
1360 let port: u16 = base
1361 .trim_end_matches('/')
1362 .rsplit(':')
1363 .next()
1364 .unwrap()
1365 .parse()
1366 .unwrap();
1367 crate::net::test_host_override(
1368 "write-ordinary.test",
1369 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1370 );
1371 let http = Client::new();
1372 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1373 crate::store::init_schema(&pool).await.unwrap();
1374 let key = SigningKey::generate("k");
1375 let mut s = session();
1376 s.aud = format!("http://write-ordinary.test:{port}");
1377 let repo = repo(&http, &pool, &s, &key);
1378
1379 repo.delete_record("c", "r")
1380 .await
1381 .expect("an ordinary write response was refused");
1382 }
1383
1384 #[tokio::test]
1393 async fn the_live_walk_refuses_a_duplicated_records_key() {
1394 let base = crate::net::tests::serve_body(
1395 br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#.to_vec(),
1396 )
1397 .await;
1398 let port: u16 = base
1399 .trim_end_matches('/')
1400 .rsplit(':')
1401 .next()
1402 .unwrap()
1403 .parse()
1404 .unwrap();
1405 crate::net::test_host_override(
1406 "dup-records.test",
1407 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1408 );
1409
1410 let http = Client::new();
1411 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1412 crate::store::init_schema(&pool).await.unwrap();
1413 let key = SigningKey::generate("k");
1414 let mut s = session();
1415 s.aud = format!("http://dup-records.test:{port}");
1416 let repo = repo(&http, &pool, &s, &key);
1417
1418 let err = repo
1419 .list_records("app.feather.subscription", None, None)
1420 .await
1421 .expect_err("a duplicated records key was read as an empty page");
1422 assert!(
1423 format!("{err:#}").contains("duplicate"),
1424 "failed for the wrong reason: {err:#}"
1425 );
1426 }
1427
1428 #[tokio::test]
1435 async fn an_empty_200_body_is_not_an_empty_repo() {
1436 let base = crate::net::tests::serve_body(Vec::new()).await;
1437 let port: u16 = base
1438 .trim_end_matches('/')
1439 .rsplit(':')
1440 .next()
1441 .unwrap()
1442 .parse()
1443 .unwrap();
1444 crate::net::test_host_override(
1445 "empty-body.test",
1446 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1447 );
1448 let http = Client::new();
1449 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1450 crate::store::init_schema(&pool).await.unwrap();
1451 let key = SigningKey::generate("k");
1452 let mut s = session();
1453 s.aud = format!("http://empty-body.test:{port}");
1454 let repo = repo(&http, &pool, &s, &key);
1455
1456 let err = repo
1457 .list_records("app.feather.subscription", None, None)
1458 .await
1459 .expect_err("an empty body was read as an empty repo");
1460 assert!(
1461 format!("{err:#}").contains("no records"),
1462 "failed for the wrong reason: {err:#}"
1463 );
1464 }
1465
1466 #[tokio::test]
1470 async fn a_200_error_envelope_is_not_a_successful_write() {
1471 let base = crate::net::tests::serve_body(
1472 br#"{"error":"InvalidRequest","message":"nope"}"#.to_vec(),
1473 )
1474 .await;
1475 let port: u16 = base
1476 .trim_end_matches('/')
1477 .rsplit(':')
1478 .next()
1479 .unwrap()
1480 .parse()
1481 .unwrap();
1482 crate::net::test_host_override(
1483 "envelope-write.test",
1484 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1485 );
1486 let http = Client::new();
1487 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1488 crate::store::init_schema(&pool).await.unwrap();
1489 let key = SigningKey::generate("k");
1490 let mut s = session();
1491 s.aud = format!("http://envelope-write.test:{port}");
1492 let repo = repo(&http, &pool, &s, &key);
1493
1494 let err = repo
1495 .delete_record("app.feather.subscription", "rk1")
1496 .await
1497 .expect_err("a failed delete was reported as success");
1498 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1499
1500 let err = repo
1501 .apply_writes(&[crate::atproto::WriteOp::Delete {
1502 collection: "app.feather.subscription".to_string(),
1503 rkey: "rk1".to_string(),
1504 }])
1505 .await
1506 .expect_err("a failed batch was reported as success");
1507 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1508 }
1509
1510 #[tokio::test]
1513 async fn an_empty_batch_is_not_sent() {
1514 let http = Client::new();
1515 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1516 let key = SigningKey::generate("k");
1517 let mut s = session();
1518 s.aud = "http://127.0.0.1:2583".into();
1520 let repo = repo(&http, &pool, &s, &key);
1521 assert!(repo.apply_writes(&[]).await.is_ok());
1522 }
1523
1524 async fn strict_pds(
1533 fail_call: Option<usize>,
1534 ) -> (
1535 OAuthSession,
1536 SqlitePool,
1537 crate::atproto::tests::ApplyWritesLog,
1538 ) {
1539 let (base, log) = crate::atproto::tests::serve_apply_writes(fail_call).await;
1540 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
1541 let host = format!("chunk-oauth-{port}.test");
1542 crate::net::test_host_override(&host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
1543 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1544 crate::store::init_schema(&pool).await.unwrap();
1545 let mut s = session();
1546 s.aud = format!("http://{host}:{port}");
1547 (s, pool, log)
1548 }
1549
1550 fn subs(n: usize) -> Vec<crate::vetted::VettedSubscription> {
1551 (0..n)
1552 .map(|i| {
1553 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
1554 format!("https://f{i}.example/feed.xml"),
1555 "2026-07-12T00:00:00.000Z",
1556 ))
1557 })
1558 .collect()
1559 }
1560
1561 #[tokio::test]
1565 async fn bulk_subscribe_of_201_is_two_calls_in_order() {
1566 let (s, pool, log) = strict_pds(None).await;
1567 let (http, key) = (Client::new(), SigningKey::generate("k"));
1568 let rkeys = repo(&http, &pool, &s, &key)
1569 .add_subscriptions_bulk(&subs(201))
1570 .await
1571 .expect("a 201-feed import must succeed against a PDS that caps at 200");
1572 assert_eq!(crate::atproto::tests::call_sizes(&log), vec![200, 1]);
1573 assert_eq!(
1574 crate::atproto::tests::sent_rkeys(&log),
1575 rkeys,
1576 "every feed, once, in input order"
1577 );
1578 }
1579
1580 #[tokio::test]
1582 async fn bulk_subscribe_splits_at_200_and_not_before() {
1583 for (n, want) in [(500, vec![200, 200, 100]), (200, vec![200])] {
1584 let (s, pool, log) = strict_pds(None).await;
1585 let (http, key) = (Client::new(), SigningKey::generate("k"));
1586 repo(&http, &pool, &s, &key)
1587 .add_subscriptions_bulk(&subs(n))
1588 .await
1589 .expect("bulk write");
1590 assert_eq!(crate::atproto::tests::call_sizes(&log), want, "{n} feeds");
1591 }
1592 }
1593
1594 #[tokio::test]
1596 async fn bulk_subscribe_stops_at_the_first_failed_chunk() {
1597 let (s, pool, log) = strict_pds(Some(2)).await;
1598 let (http, key) = (Client::new(), SigningKey::generate("k"));
1599 let err = repo(&http, &pool, &s, &key)
1600 .add_subscriptions_bulk(&subs(500))
1601 .await
1602 .expect_err("a failed chunk must fail the call");
1603 assert_eq!(
1604 crate::atproto::tests::call_sizes(&log),
1605 vec![200, 200],
1606 "chunk 3 must NOT be sent"
1607 );
1608 assert!(format!("{err:#}").contains("boom"), "{err:#}");
1609 }
1610
1611 #[tokio::test]
1615 async fn read_state_flush_splits_on_bytes_under_200_ops() {
1616 let (s, pool, log) = strict_pds(None).await;
1617 let (http, key) = (Client::new(), SigningKey::generate("k"));
1618 let cursors: Vec<(String, crate::lexicon::ReadState, bool)> = (0..10)
1619 .map(|i| {
1620 let mut state = crate::lexicon::ReadState::new(
1621 format!("https://f{i}.example/feed.xml"),
1622 None,
1623 "2026-07-12T00:00:00.000Z",
1624 );
1625 state.read_ids = (0..crate::lexicon::ReadState::MAX_IDS)
1626 .map(|j| format!("https://f{i}.example/posts/{j:04}/an-entry-permalink"))
1627 .collect();
1628 (format!("rk{i:04}"), state, false)
1629 })
1630 .collect();
1631 repo(&http, &pool, &s, &key)
1632 .flush_read_states(&cursors)
1633 .await
1634 .expect("a byte-heavy flush must succeed in chunks");
1635 let sizes = crate::atproto::tests::call_sizes(&log);
1636 assert!(sizes.len() > 1, "one call for ~500 KB: {sizes:?}");
1637 let want: Vec<String> = cursors.iter().map(|(rkey, _, _)| rkey.clone()).collect();
1638 assert_eq!(crate::atproto::tests::sent_rkeys(&log), want);
1639 }
1640
1641 fn session_at(pds: &str) -> OAuthSession {
1646 let mut s = session();
1647 s.aud = pds.to_string();
1648 s
1649 }
1650
1651 #[tokio::test]
1656 async fn put_record_sends_swap_record_only_when_given() {
1657 use crate::atproto::tests::{serve_status_json, swap_sub, write_ok, OLD_CID};
1658 let (_, pds, log) = serve_status_json(200, write_ok()).await;
1659 let http = Client::new();
1660 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1661 crate::store::init_schema(&pool).await.unwrap();
1662 let key = SigningKey::generate("k");
1663 let s = session_at(&pds);
1664 let repo = repo(&http, &pool, &s, &key);
1665
1666 repo.update_subscription("rk", &swap_sub(), Some(OLD_CID))
1667 .await
1668 .expect("put with a swap");
1669 repo.update_subscription("rk", &swap_sub(), None)
1670 .await
1671 .expect("put without a swap");
1672
1673 let sent = log.lock().unwrap().clone();
1674 assert_eq!(sent.len(), 2, "{sent:?}");
1675 assert_eq!(sent[0]["rkey"], "rk", "captured no usable body: {sent:?}");
1676 assert_eq!(
1677 sent[0]["swapRecord"], OLD_CID,
1678 "the CID the caller read never reached the PDS: {}",
1679 sent[0]
1680 );
1681 assert_eq!(sent[1]["rkey"], "rk");
1682 assert!(
1683 sent[1].get("swapRecord").is_none(),
1684 "no swap was asked for, so none may be sent: {}",
1685 sent[1]
1686 );
1687 }
1688
1689 #[tokio::test]
1692 async fn an_invalid_swap_from_the_pds_is_recognised() {
1693 use crate::atproto::tests::{invalid_swap_xrpc, serve_status_json, swap_sub, OLD_CID};
1694 let http = Client::new();
1695 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1696 crate::store::init_schema(&pool).await.unwrap();
1697 let key = SigningKey::generate("k");
1698
1699 let (_, pds, _) = serve_status_json(400, invalid_swap_xrpc()).await;
1700 let s = session_at(&pds);
1701 let err = repo(&http, &pool, &s, &key)
1702 .update_subscription("rk", &swap_sub(), Some(OLD_CID))
1703 .await
1704 .expect_err("the PDS refused the swap");
1705 assert!(crate::atproto::is_invalid_swap(&err), "{err:#}");
1706
1707 let (_, pds, _) = serve_status_json(
1708 400,
1709 serde_json::json!({ "error": "InvalidRequest", "message": "bad record" }),
1710 )
1711 .await;
1712 let s = session_at(&pds);
1713 let err = repo(&http, &pool, &s, &key)
1714 .update_subscription("rk", &swap_sub(), Some(OLD_CID))
1715 .await
1716 .expect_err("refused");
1717 assert!(!crate::atproto::is_invalid_swap(&err), "{err:#}");
1718 }
1719
1720 #[tokio::test]
1723 async fn list_folders_with_cids_keeps_each_records_cid() {
1724 use crate::atproto::tests::{
1725 assert_folders_listed_with_cids, serve_status_json, two_folders_page,
1726 };
1727 let (_, pds, _) = serve_status_json(200, two_folders_page()).await;
1728 let http = Client::new();
1729 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1730 crate::store::init_schema(&pool).await.unwrap();
1731 let key = SigningKey::generate("k");
1732 let s = session_at(&pds);
1733 let listed = repo(&http, &pool, &s, &key)
1734 .list_folders_with_cids()
1735 .await
1736 .expect("listing");
1737 assert_folders_listed_with_cids(&listed);
1738 }
1739
1740 #[tokio::test]
1742 async fn list_subscriptions_with_cids_keeps_each_records_cid() {
1743 use crate::atproto::tests::{assert_listed_with_cids, serve_status_json, two_subs_page};
1744 let (_, pds, _) = serve_status_json(200, two_subs_page()).await;
1745 let http = Client::new();
1746 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1747 crate::store::init_schema(&pool).await.unwrap();
1748 let key = SigningKey::generate("k");
1749 let s = session_at(&pds);
1750 let listed = repo(&http, &pool, &s, &key)
1751 .list_subscriptions_with_cids()
1752 .await
1753 .expect("listing");
1754 assert_listed_with_cids(&listed);
1755 }
1756}