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>(
348 &self,
349 collection: &str,
350 rkey: &str,
351 record: &T,
352 swap_record: Option<&str>,
353 ) -> Result<WriteResult> {
354 let mut body = json!({
355 "repo": self.session.sub,
356 "collection": collection,
357 "rkey": rkey,
358 "record": record,
359 });
360 if let Some(cid) = swap_record {
361 body["swapRecord"] = json!(cid);
362 }
363 self.write("com.atproto.repo.putRecord", body).await
364 }
365
366 pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
367 let body = json!({ "repo": self.session.sub, "collection": collection, "rkey": rkey });
368 self.send(
369 &self.url("com.atproto.repo.deleteRecord"),
370 DpopBody::Json(serde_json::to_vec(&body)?),
371 "com.atproto.repo.deleteRecord",
372 )
373 .await
374 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
375 Ok(())
376 }
377
378 pub async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
383 crate::atproto::apply_writes_chunked(writes, |chunk| self.apply_writes_once(chunk)).await
384 }
385
386 async fn apply_writes_once(&self, writes: &[WriteOp]) -> Result<()> {
390 let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
391 let body = json!({ "repo": self.session.sub, "writes": ops });
392 self.send(
393 &self.url("com.atproto.repo.applyWrites"),
394 DpopBody::Json(serde_json::to_vec(&body)?),
395 "com.atproto.repo.applyWrites",
396 )
397 .await
398 .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
399 Ok(())
400 }
401
402 async fn write(&self, nsid: &str, body: Value) -> Result<WriteResult> {
403 let value = self
404 .send(
405 &self.url(nsid),
406 DpopBody::Json(serde_json::to_vec(&body)?),
407 nsid,
408 )
409 .await?;
410 serde_json::from_value(value).with_context(|| format!("{nsid} returned no usable result"))
411 }
412}
413
414fn xrpc_error(body: &[u8], status: u16) -> String {
416 match Repo::error_fields(body) {
417 Some(detail) => format!("status {status} ({detail})"),
418 None => format!("status {status}"),
419 }
420}
421
422impl Repo<'_> {
430 async fn list_typed<T: serde::de::DeserializeOwned>(
437 &self,
438 collection: &str,
439 ) -> Result<Vec<(String, T)>> {
440 Ok(self
441 .list_typed_with_cids(collection)
442 .await?
443 .into_iter()
444 .map(|(rkey, _cid, value)| (rkey, value))
445 .collect())
446 }
447
448 async fn list_typed_with_cids<T: serde::de::DeserializeOwned>(
452 &self,
453 collection: &str,
454 ) -> Result<Vec<(String, Option<String>, T)>> {
455 let records = self.list_all_records(collection).await?;
456 let mut out = Vec::with_capacity(records.len());
457 for record in records {
458 let rkey = record.rkey().unwrap_or_default().to_string();
459 match record.parse::<T>() {
460 Ok(value) => out.push((rkey, record.cid, value)),
461 Err(err) => tracing::warn!(
462 collection,
463 uri = %record.uri,
464 error = %err,
465 "skipping unparseable record in collection"
466 ),
467 }
468 }
469 Ok(out)
470 }
471
472 pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
475 self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
476 }
477
478 pub async fn list_subscriptions_with_cids(
481 &self,
482 ) -> Result<Vec<(String, Option<String>, crate::lexicon::Subscription)>> {
483 self.list_typed_with_cids(crate::lexicon::nsid::SUBSCRIPTION)
484 .await
485 }
486
487 pub async fn list_subscriptions_sorted(
489 &self,
490 ) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
491 let mut subs = self.list_subscriptions().await?;
492 subs.sort_by(crate::lexicon::sort::subscriptions);
493 Ok(subs)
494 }
495
496 pub async fn add_subscription(
499 &self,
500 sub: &crate::vetted::VettedSubscription,
501 ) -> Result<String> {
502 Ok(self
503 .create_record(crate::lexicon::nsid::SUBSCRIPTION, sub)
504 .await?
505 .into_rkey())
506 }
507
508 pub async fn remove_subscription(&self, rkey: &str) -> Result<()> {
509 self.delete_record(crate::lexicon::nsid::SUBSCRIPTION, rkey)
510 .await
511 }
512
513 pub async fn update_subscription(
518 &self,
519 rkey: &str,
520 sub: &crate::vetted::VettedSubscription,
521 swap_record: Option<&str>,
522 ) -> Result<WriteResult> {
523 self.put_record(crate::lexicon::nsid::SUBSCRIPTION, rkey, sub, swap_record)
524 .await
525 }
526
527 pub async fn add_subscriptions_bulk(
539 &self,
540 subs: &[crate::vetted::VettedSubscription],
541 ) -> Result<Vec<String>> {
542 let mut gen = crate::atproto::TidGenerator::new();
543 let mut rkeys = Vec::with_capacity(subs.len());
544 let mut writes = Vec::with_capacity(subs.len());
545 for sub in subs {
546 let rkey = gen.next();
547 writes.push(WriteOp::Create {
548 collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
549 rkey: Some(rkey.clone()),
550 value: serde_json::to_value(sub)?,
554 });
555 rkeys.push(rkey);
556 }
557 self.apply_writes(&writes).await?;
558 Ok(rkeys)
559 }
560
561 pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
564 self.list_typed(crate::lexicon::nsid::FOLDER).await
565 }
566
567 pub async fn list_folders_with_cids(
569 &self,
570 ) -> Result<Vec<(String, Option<String>, crate::lexicon::Folder)>> {
571 self.list_typed_with_cids(crate::lexicon::nsid::FOLDER)
572 .await
573 }
574
575 pub async fn list_folders_sorted(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
576 let mut folders = self.list_folders().await?;
577 folders.sort_by(crate::lexicon::sort::folders);
578 Ok(folders)
579 }
580
581 pub async fn add_folder(&self, folder: &crate::lexicon::Folder) -> Result<String> {
582 Ok(self
583 .create_record(crate::lexicon::nsid::FOLDER, folder)
584 .await?
585 .into_rkey())
586 }
587
588 pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
592 self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
593 }
594
595 pub async fn rename_folder(
600 &self,
601 rkey: &str,
602 folder: &crate::lexicon::Folder,
603 swap_record: Option<&str>,
604 ) -> Result<WriteResult> {
605 self.put_record(crate::lexicon::nsid::FOLDER, rkey, folder, swap_record)
606 .await
607 }
608
609 pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
612 self.list_typed(crate::lexicon::nsid::SAVED).await
613 }
614
615 pub async fn list_saved_sorted(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
617 let mut saved = self.list_saved().await?;
618 saved.sort_by(crate::lexicon::sort::saved);
619 Ok(saved)
620 }
621
622 pub async fn add_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<String> {
623 Ok(self
624 .create_record(crate::lexicon::nsid::SAVED, saved)
625 .await?
626 .into_rkey())
627 }
628
629 pub async fn remove_saved(&self, rkey: &str) -> Result<()> {
630 self.delete_record(crate::lexicon::nsid::SAVED, rkey).await
631 }
632
633 pub async fn list_read_states(&self) -> Result<Vec<(String, crate::lexicon::ReadState)>> {
636 self.list_typed(crate::lexicon::nsid::READ_STATE).await
637 }
638
639 pub async fn put_read_state(
641 &self,
642 rkey: &str,
643 state: &crate::lexicon::ReadState,
644 ) -> Result<()> {
645 self.put_record(crate::lexicon::nsid::READ_STATE, rkey, state, None)
646 .await?;
647 Ok(())
648 }
649
650 pub async fn flush_read_states(
664 &self,
665 cursors: &[(String, crate::lexicon::ReadState, bool)],
666 ) -> Result<()> {
667 if cursors.is_empty() {
668 return Ok(());
669 }
670 let writes = crate::atproto::read_state_write_ops(cursors)?;
671 self.apply_writes(&writes).await
672 }
673}
674
675#[cfg(test)]
676mod tests {
677 use super::*;
678
679 #[test]
692 fn read_state_writes_choose_create_or_update_per_cursor() {
693 let state = crate::lexicon::ReadState::new(
694 "https://example.com/feed",
695 Some("2026-01-01T00:00:00Z".to_string()),
696 "2026-01-01T00:00:00Z",
697 );
698 let cursors = vec![
699 ("existing".to_string(), state.clone(), true),
700 ("brand-new".to_string(), state.clone(), false),
701 ];
702
703 let ops = crate::atproto::read_state_write_ops(&cursors).expect("ops build");
704 assert_eq!(ops.len(), 2);
705
706 let rendered: Vec<Value> = ops.iter().map(|op| op.to_json()).collect();
707 assert_eq!(
708 rendered[0]["$type"], "com.atproto.repo.applyWrites#update",
709 "an existing record must be UPDATED, not re-created"
710 );
711 assert_eq!(
712 rendered[1]["$type"], "com.atproto.repo.applyWrites#create",
713 "a first flush must CREATE, or the whole atomic batch fails"
714 );
715 }
716
717 #[tokio::test]
722 async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
723 let (base, log) = crate::net::tests::serve_json_capturing(b"{}".to_vec()).await;
724 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
725 crate::net::test_host_override(
726 "bulk-pds.test",
727 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
728 );
729 let http = Client::new();
730 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
731 crate::store::init_schema(&pool).await.unwrap();
732 let key = SigningKey::generate("k");
733 let mut s = session();
734 s.aud = format!("http://bulk-pds.test:{port}");
735 let repo = repo(&http, &pool, &s, &key);
736 let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
737 .map(|i| {
738 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
739 format!("https://f{i}.example/feed.xml"),
740 "2026-07-12T00:00:00.000Z",
741 ))
742 })
743 .collect();
744
745 let rkeys = repo
746 .add_subscriptions_bulk(&subs)
747 .await
748 .expect("bulk write failed");
749
750 let sent = log.lock().unwrap().clone();
751 assert_eq!(
752 sent.len(),
753 1,
754 "expected one applyWrites request, got {sent:?}"
755 );
756 let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
757 .expect("request body is JSON");
758 let writes = body["writes"].as_array().expect("writes array");
759 assert_eq!(writes.len(), 3);
760 for (i, w) in writes.iter().enumerate() {
761 assert_eq!(w["collection"], crate::lexicon::nsid::SUBSCRIPTION);
762 assert_eq!(w["rkey"].as_str(), Some(rkeys[i].as_str()));
763 }
764 let mut sorted = rkeys.clone();
765 sorted.sort();
766 assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
767 }
768
769 fn session() -> OAuthSession {
770 OAuthSession {
771 sub: "did:plc:ewvi7nxzyoun6zhxrhs64oiz".into(),
772 issuer: "https://pds.example.com".into(),
773 aud: "https://pds.example.com".into(),
774 dpop_key_jwk: "{}".into(),
775 access_token: "tok".into(),
776 refresh_token: "ref".into(),
777 token_type: "DPoP".into(),
778 granted_scope: "atproto".into(),
779 expires_at: None,
780 }
781 }
782
783 fn repo<'a>(
784 http: &'a Client,
785 pool: &'a SqlitePool,
786 session: &'a OAuthSession,
787 key: &'a SigningKey,
788 ) -> Repo<'a> {
789 Repo {
790 http,
791 pool,
792 session,
793 key,
794 }
795 }
796
797 #[tokio::test]
800 async fn endpoints_are_built_from_the_sessions_audience() {
801 let http = Client::new();
802 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
803 let key = SigningKey::generate("k");
804 let mut s = session();
805 s.aud = "https://pds.example.com/".into();
806 let repo = repo(&http, &pool, &s, &key);
807 assert_eq!(
808 repo.url("com.atproto.repo.listRecords"),
809 "https://pds.example.com/xrpc/com.atproto.repo.listRecords",
810 "a trailing slash on the audience must not double the separator"
811 );
812 }
813
814 #[test]
817 fn an_xrpc_error_is_summarised_not_echoed() {
818 let body = br#"{"error":"InvalidRequest","message":"unknown collection"}"#;
819 let rendered = xrpc_error(body, 400);
820 assert!(rendered.contains("InvalidRequest"));
821 assert!(rendered.contains("unknown collection"));
822
823 let opaque = xrpc_error(br#"{"access_token":"SECRET"}"#, 500);
825 assert_eq!(opaque, "status 500");
826 assert!(!opaque.contains("SECRET"));
827 assert_eq!(xrpc_error(b"<html>oops</html>", 502), "status 502");
828 }
829
830 #[tokio::test]
833 async fn repo_calls_fail_closed_on_an_internal_pds() {
834 let http = Client::new();
835 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
836 super::super::store::init_schema(&pool).await.unwrap();
837 let key = SigningKey::generate("k");
838 let mut s = session();
839 s.aud = "http://127.0.0.1:2583".into();
840 let repo = repo(&http, &pool, &s, &key);
841
842 let err = repo
843 .list_records("app.feather.subscription", None, None)
844 .await
845 .expect_err("must refuse a loopback PDS");
846 assert!(
847 format!("{err:#}").contains("forbidden (internal) address"),
848 "failed for the wrong reason: {err:#}"
849 );
850 }
851
852 #[tokio::test]
862 async fn a_200_error_envelope_is_not_an_empty_repo() {
863 let base = crate::net::tests::serve_body(
864 br#"{"error":"InvalidRequest","message":"bad cursor"}"#.to_vec(),
865 )
866 .await;
867 let port: u16 = base
868 .trim_end_matches('/')
869 .rsplit(':')
870 .next()
871 .unwrap()
872 .parse()
873 .unwrap();
874 crate::net::test_host_override(
875 "envelope-pds.test",
876 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
877 );
878
879 let http = Client::new();
880 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
881 crate::store::init_schema(&pool).await.unwrap();
882 let key = SigningKey::generate("k");
883 let mut s = session();
884 s.aud = format!("http://envelope-pds.test:{port}");
885 let repo = repo(&http, &pool, &s, &key);
886
887 let err = repo
888 .list_records("app.feather.subscription", None, None)
889 .await
890 .expect_err("an error envelope was read as an empty page");
891 assert!(
892 format!("{err:#}").contains("InvalidRequest"),
893 "failed for the wrong reason: {err:#}"
894 );
895 }
896
897 #[tokio::test]
904 async fn the_live_walk_spends_its_budget_across_pages() {
905 let (bodies, per_page) = crate::atproto::tests::paged_bodies(3, 4096, false);
906 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
907 let port: u16 = base
908 .trim_end_matches('/')
909 .rsplit(':')
910 .next()
911 .unwrap()
912 .parse()
913 .unwrap();
914 crate::net::test_host_override(
915 "live-budget-pages.test",
916 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
917 );
918
919 let http = Client::new();
920 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
921 crate::store::init_schema(&pool).await.unwrap();
922 let key = SigningKey::generate("k");
923 let mut s = session();
924 s.aud = format!("http://live-budget-pages.test:{port}");
925 let repo = repo(&http, &pool, &s, &key);
926
927 let err = repo
928 .list_all_records_within(
929 "app.feather.subscription",
930 &mut crate::atproto::ByteBudget::new(per_page * 2),
931 )
932 .await
933 .expect_err("three pages cannot fit in a two-page budget");
934 let msg = format!("{err:#}");
935 assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
936 assert!(
937 msg.contains("2 held"),
938 "the live walk did not accumulate across pages: {msg}"
939 );
940 }
941
942 #[tokio::test]
946 async fn the_live_walk_that_runs_out_of_pages_refuses() {
947 let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
948 .map(|i| {
949 serde_json::json!({
950 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
951 "cursor": format!("p{}", i + 1),
952 })
953 .to_string()
954 .into_bytes()
955 })
956 .collect();
957 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
958 let port: u16 = base
959 .trim_end_matches('/')
960 .rsplit(':')
961 .next()
962 .unwrap()
963 .parse()
964 .unwrap();
965 crate::net::test_host_override(
966 "pages-exhausted-live.test",
967 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
968 );
969 let http = Client::new();
970 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
971 crate::store::init_schema(&pool).await.unwrap();
972 let key = SigningKey::generate("k");
973 let mut s = session();
974 s.aud = format!("http://pages-exhausted-live.test:{port}");
975 let repo = repo(&http, &pool, &s, &key);
976
977 let err = repo
978 .list_all_records("app.feather.subscription")
979 .await
980 .expect_err("a truncated list was returned as a complete one");
981 assert!(
982 format!("{err:#}").contains("did not finish"),
983 "failed for the wrong reason: {err:#}"
984 );
985 }
986
987 #[tokio::test]
990 async fn the_live_walk_refuses_a_page_with_a_malformed_record() {
991 let body = serde_json::json!({ "records": [
992 { "uri": "at://did:plc:x/c/3labGOOD", "value": {} },
993 { "cid": "bafy", "value": {} },
994 ]})
995 .to_string()
996 .into_bytes();
997 let base = crate::net::tests::serve_bodies_in_sequence(vec![body]).await;
998 let port: u16 = base
999 .trim_end_matches('/')
1000 .rsplit(':')
1001 .next()
1002 .unwrap()
1003 .parse()
1004 .unwrap();
1005 crate::net::test_host_override(
1006 "malformed-live.test",
1007 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1008 );
1009 let http = Client::new();
1010 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1011 crate::store::init_schema(&pool).await.unwrap();
1012 let key = SigningKey::generate("k");
1013 let mut s = session();
1014 s.aud = format!("http://malformed-live.test:{port}");
1015 let repo = repo(&http, &pool, &s, &key);
1016 let err = repo
1017 .list_all_records("app.feather.subscription")
1018 .await
1019 .expect_err("a page with a malformed record was accepted");
1020 assert!(
1021 err.downcast_ref::<crate::atproto::MalformedRecords>()
1022 .is_some(),
1023 "refused for the wrong reason: {err:#}"
1024 );
1025 }
1026
1027 #[tokio::test]
1030 async fn the_live_walk_that_finishes_cleanly_returns_the_records() {
1031 let mut bodies: Vec<Vec<u8>> = (0..3)
1032 .map(|i| {
1033 serde_json::json!({
1034 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
1035 "cursor": format!("p{}", i + 1),
1036 })
1037 .to_string()
1038 .into_bytes()
1039 })
1040 .collect();
1041 bodies.push(
1046 serde_json::json!({
1047 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1048 })
1049 .to_string()
1050 .into_bytes(),
1051 );
1052 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1053 let port: u16 = base
1054 .trim_end_matches('/')
1055 .rsplit(':')
1056 .next()
1057 .unwrap()
1058 .parse()
1059 .unwrap();
1060 crate::net::test_host_override(
1061 "clean-finish-live.test",
1062 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1063 );
1064 let http = Client::new();
1065 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1066 crate::store::init_schema(&pool).await.unwrap();
1067 let key = SigningKey::generate("k");
1068 let mut s = session();
1069 s.aud = format!("http://clean-finish-live.test:{port}");
1070 let repo = repo(&http, &pool, &s, &key);
1071
1072 let records = repo
1073 .list_all_records("app.feather.subscription")
1074 .await
1075 .expect("a walk that ran out of records is not a short list");
1076 assert_eq!(records.len(), 4);
1077 assert!(
1078 records.iter().any(|r| r.uri.ends_with("3labLAST")),
1079 "the LAST page's records were dropped: {:?}",
1080 records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
1081 );
1082 }
1083
1084 #[tokio::test]
1092 async fn a_live_walk_that_terminates_on_its_last_allowed_page_succeeds() {
1093 let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
1094 .map(|i| {
1095 serde_json::json!({
1096 "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
1097 "cursor": format!("p{}", i + 1),
1098 })
1099 .to_string()
1100 .into_bytes()
1101 })
1102 .collect();
1103 bodies.push(
1104 serde_json::json!({
1105 "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1106 })
1107 .to_string()
1108 .into_bytes(),
1109 );
1110 assert_eq!(bodies.len(), MAX_LIST_PAGES);
1111 let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1112 let port: u16 = base
1113 .trim_end_matches('/')
1114 .rsplit(':')
1115 .next()
1116 .unwrap()
1117 .parse()
1118 .unwrap();
1119 crate::net::test_host_override(
1120 "last-allowed-page-live.test",
1121 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1122 );
1123 let http = Client::new();
1124 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1125 crate::store::init_schema(&pool).await.unwrap();
1126 let key = SigningKey::generate("k");
1127 let mut s = session();
1128 s.aud = format!("http://last-allowed-page-live.test:{port}");
1129 let repo = repo(&http, &pool, &s, &key);
1130
1131 let records = repo
1132 .list_all_records("app.feather.subscription")
1133 .await
1134 .expect("a walk that terminated inside its budget is not a short list");
1135 assert_eq!(
1136 records.len(),
1137 MAX_LIST_PAGES,
1138 "a walk that used its whole page budget and finished lost records",
1139 );
1140 }
1141
1142 #[test]
1154 fn an_oversized_error_body_is_not_parsed_for_its_reason() {
1155 let small = br#"{"error":"InvalidSwap","message":"record changed"}"#;
1156 assert_eq!(
1157 Repo::error_fields(small).as_deref(),
1158 Some("InvalidSwap: record changed"),
1159 "a real error body must still render its reason",
1160 );
1161
1162 let mut huge = String::from(r#"{"error":"InvalidSwap","pad":["#);
1163 while huge.len() < crate::oauth::MAX_ERROR_BODY + 1_024 {
1164 huge.push_str("{},");
1165 }
1166 huge.push_str("{}]}");
1167 assert!(huge.len() > crate::oauth::MAX_ERROR_BODY);
1168 assert_eq!(
1169 Repo::error_fields(huge.as_bytes()),
1170 None,
1171 "an oversized error body was deserialised to fish out one string",
1172 );
1173 }
1174
1175 #[tokio::test]
1184 async fn a_rejected_write_carries_the_status_and_error_name() {
1185 use axum::response::IntoResponse as _;
1186 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1187 let addr = listener.local_addr().unwrap();
1188 let host = format!("rejecting-{}.xrpc.test", addr.port());
1189 crate::net::test_host_override(&host, addr);
1190 let app = axum::Router::new().fallback(|| async {
1191 (
1192 axum::http::StatusCode::INTERNAL_SERVER_ERROR,
1193 axum::Json(
1194 json!({ "error": "InternalServerError", "message": "Internal Server Error" }),
1195 ),
1196 )
1197 .into_response()
1198 });
1199 tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1200
1201 let http = Client::new();
1202 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1203 crate::store::init_schema(&pool).await.unwrap();
1204 let key = SigningKey::generate("k");
1205 let mut s = session();
1206 s.aud = format!("http://{host}:{}", addr.port());
1207 let repo = repo(&http, &pool, &s, &key);
1208
1209 let err = repo
1210 .apply_writes(&[WriteOp::Delete {
1211 collection: crate::lexicon::nsid::READ_STATE.into(),
1212 rkey: "rs-0".into(),
1213 }])
1214 .await
1215 .expect_err("a 500 is a failure");
1216 assert_eq!(
1217 err.to_string(),
1218 "com.atproto.repo.applyWrites failed: status 500 \
1219 (InternalServerError: Internal Server Error)",
1220 "the rendered message changed",
1221 );
1222 let xrpc = err
1223 .chain()
1224 .find_map(|cause| cause.downcast_ref::<crate::atproto::AtProtoError>());
1225 match xrpc {
1226 Some(crate::atproto::AtProtoError::Xrpc { status, error, .. }) => {
1227 assert_eq!(status.as_u16(), 500);
1228 assert_eq!(error, "InternalServerError");
1229 }
1230 other => panic!("no structured XRPC error in the chain: {other:?}"),
1231 }
1232 }
1233
1234 #[tokio::test]
1247 async fn the_live_write_path_refuses_a_node_explosion() {
1248 let mut body = String::from(r#"{"uri":"at://d/c/r","value":["#);
1249 for _ in 0..1_200_000 {
1250 body.push_str("{},");
1251 }
1252 body.push_str("{}]}");
1253 assert!(
1254 crate::atproto::count_structural_chars(body.as_bytes())
1255 > crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
1256 "the probe body is not over the cap, so this test proves nothing",
1257 );
1258 let base = crate::net::tests::serve_body(body.into_bytes()).await;
1259 let port: u16 = base
1260 .trim_end_matches('/')
1261 .rsplit(':')
1262 .next()
1263 .unwrap()
1264 .parse()
1265 .unwrap();
1266 crate::net::test_host_override(
1267 "write-explosion.test",
1268 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1269 );
1270 let http = Client::new();
1271 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1272 crate::store::init_schema(&pool).await.unwrap();
1273 let key = SigningKey::generate("k");
1274 let mut s = session();
1275 s.aud = format!("http://write-explosion.test:{port}");
1276 let repo = repo(&http, &pool, &s, &key);
1277
1278 let err = repo
1279 .delete_record("c", "r")
1280 .await
1281 .expect_err("a node explosion on the write path was parsed rather than refused");
1282 assert!(
1283 format!("{err:#}").contains("structural characters"),
1284 "failed for the wrong reason: {err:#}"
1285 );
1286 }
1287
1288 #[tokio::test]
1291 async fn the_live_write_path_accepts_an_ordinary_response() {
1292 let base =
1293 crate::net::tests::serve_body(br#"{"commit":{"cid":"bafy","rev":"3lab"}}"#.to_vec())
1294 .await;
1295 let port: u16 = base
1296 .trim_end_matches('/')
1297 .rsplit(':')
1298 .next()
1299 .unwrap()
1300 .parse()
1301 .unwrap();
1302 crate::net::test_host_override(
1303 "write-ordinary.test",
1304 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1305 );
1306 let http = Client::new();
1307 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1308 crate::store::init_schema(&pool).await.unwrap();
1309 let key = SigningKey::generate("k");
1310 let mut s = session();
1311 s.aud = format!("http://write-ordinary.test:{port}");
1312 let repo = repo(&http, &pool, &s, &key);
1313
1314 repo.delete_record("c", "r")
1315 .await
1316 .expect("an ordinary write response was refused");
1317 }
1318
1319 #[tokio::test]
1328 async fn the_live_walk_refuses_a_duplicated_records_key() {
1329 let base = crate::net::tests::serve_body(
1330 br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#.to_vec(),
1331 )
1332 .await;
1333 let port: u16 = base
1334 .trim_end_matches('/')
1335 .rsplit(':')
1336 .next()
1337 .unwrap()
1338 .parse()
1339 .unwrap();
1340 crate::net::test_host_override(
1341 "dup-records.test",
1342 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1343 );
1344
1345 let http = Client::new();
1346 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1347 crate::store::init_schema(&pool).await.unwrap();
1348 let key = SigningKey::generate("k");
1349 let mut s = session();
1350 s.aud = format!("http://dup-records.test:{port}");
1351 let repo = repo(&http, &pool, &s, &key);
1352
1353 let err = repo
1354 .list_records("app.feather.subscription", None, None)
1355 .await
1356 .expect_err("a duplicated records key was read as an empty page");
1357 assert!(
1358 format!("{err:#}").contains("duplicate"),
1359 "failed for the wrong reason: {err:#}"
1360 );
1361 }
1362
1363 #[tokio::test]
1370 async fn an_empty_200_body_is_not_an_empty_repo() {
1371 let base = crate::net::tests::serve_body(Vec::new()).await;
1372 let port: u16 = base
1373 .trim_end_matches('/')
1374 .rsplit(':')
1375 .next()
1376 .unwrap()
1377 .parse()
1378 .unwrap();
1379 crate::net::test_host_override(
1380 "empty-body.test",
1381 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1382 );
1383 let http = Client::new();
1384 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1385 crate::store::init_schema(&pool).await.unwrap();
1386 let key = SigningKey::generate("k");
1387 let mut s = session();
1388 s.aud = format!("http://empty-body.test:{port}");
1389 let repo = repo(&http, &pool, &s, &key);
1390
1391 let err = repo
1392 .list_records("app.feather.subscription", None, None)
1393 .await
1394 .expect_err("an empty body was read as an empty repo");
1395 assert!(
1396 format!("{err:#}").contains("no records"),
1397 "failed for the wrong reason: {err:#}"
1398 );
1399 }
1400
1401 #[tokio::test]
1405 async fn a_200_error_envelope_is_not_a_successful_write() {
1406 let base = crate::net::tests::serve_body(
1407 br#"{"error":"InvalidRequest","message":"nope"}"#.to_vec(),
1408 )
1409 .await;
1410 let port: u16 = base
1411 .trim_end_matches('/')
1412 .rsplit(':')
1413 .next()
1414 .unwrap()
1415 .parse()
1416 .unwrap();
1417 crate::net::test_host_override(
1418 "envelope-write.test",
1419 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1420 );
1421 let http = Client::new();
1422 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1423 crate::store::init_schema(&pool).await.unwrap();
1424 let key = SigningKey::generate("k");
1425 let mut s = session();
1426 s.aud = format!("http://envelope-write.test:{port}");
1427 let repo = repo(&http, &pool, &s, &key);
1428
1429 let err = repo
1430 .delete_record("app.feather.subscription", "rk1")
1431 .await
1432 .expect_err("a failed delete was reported as success");
1433 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1434
1435 let err = repo
1436 .apply_writes(&[crate::atproto::WriteOp::Delete {
1437 collection: "app.feather.subscription".to_string(),
1438 rkey: "rk1".to_string(),
1439 }])
1440 .await
1441 .expect_err("a failed batch was reported as success");
1442 assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1443 }
1444
1445 #[tokio::test]
1448 async fn an_empty_batch_is_not_sent() {
1449 let http = Client::new();
1450 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1451 let key = SigningKey::generate("k");
1452 let mut s = session();
1453 s.aud = "http://127.0.0.1:2583".into();
1455 let repo = repo(&http, &pool, &s, &key);
1456 assert!(repo.apply_writes(&[]).await.is_ok());
1457 }
1458
1459 async fn strict_pds(
1468 fail_call: Option<usize>,
1469 ) -> (
1470 OAuthSession,
1471 SqlitePool,
1472 crate::atproto::tests::ApplyWritesLog,
1473 ) {
1474 let (base, log) = crate::atproto::tests::serve_apply_writes(fail_call).await;
1475 let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
1476 let host = format!("chunk-oauth-{port}.test");
1477 crate::net::test_host_override(&host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
1478 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1479 crate::store::init_schema(&pool).await.unwrap();
1480 let mut s = session();
1481 s.aud = format!("http://{host}:{port}");
1482 (s, pool, log)
1483 }
1484
1485 fn subs(n: usize) -> Vec<crate::vetted::VettedSubscription> {
1486 (0..n)
1487 .map(|i| {
1488 crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
1489 format!("https://f{i}.example/feed.xml"),
1490 "2026-07-12T00:00:00.000Z",
1491 ))
1492 })
1493 .collect()
1494 }
1495
1496 #[tokio::test]
1500 async fn bulk_subscribe_of_201_is_two_calls_in_order() {
1501 let (s, pool, log) = strict_pds(None).await;
1502 let (http, key) = (Client::new(), SigningKey::generate("k"));
1503 let rkeys = repo(&http, &pool, &s, &key)
1504 .add_subscriptions_bulk(&subs(201))
1505 .await
1506 .expect("a 201-feed import must succeed against a PDS that caps at 200");
1507 assert_eq!(crate::atproto::tests::call_sizes(&log), vec![200, 1]);
1508 assert_eq!(
1509 crate::atproto::tests::sent_rkeys(&log),
1510 rkeys,
1511 "every feed, once, in input order"
1512 );
1513 }
1514
1515 #[tokio::test]
1517 async fn bulk_subscribe_splits_at_200_and_not_before() {
1518 for (n, want) in [(500, vec![200, 200, 100]), (200, vec![200])] {
1519 let (s, pool, log) = strict_pds(None).await;
1520 let (http, key) = (Client::new(), SigningKey::generate("k"));
1521 repo(&http, &pool, &s, &key)
1522 .add_subscriptions_bulk(&subs(n))
1523 .await
1524 .expect("bulk write");
1525 assert_eq!(crate::atproto::tests::call_sizes(&log), want, "{n} feeds");
1526 }
1527 }
1528
1529 #[tokio::test]
1531 async fn bulk_subscribe_stops_at_the_first_failed_chunk() {
1532 let (s, pool, log) = strict_pds(Some(2)).await;
1533 let (http, key) = (Client::new(), SigningKey::generate("k"));
1534 let err = repo(&http, &pool, &s, &key)
1535 .add_subscriptions_bulk(&subs(500))
1536 .await
1537 .expect_err("a failed chunk must fail the call");
1538 assert_eq!(
1539 crate::atproto::tests::call_sizes(&log),
1540 vec![200, 200],
1541 "chunk 3 must NOT be sent"
1542 );
1543 assert!(format!("{err:#}").contains("boom"), "{err:#}");
1544 }
1545
1546 #[tokio::test]
1550 async fn read_state_flush_splits_on_bytes_under_200_ops() {
1551 let (s, pool, log) = strict_pds(None).await;
1552 let (http, key) = (Client::new(), SigningKey::generate("k"));
1553 let cursors: Vec<(String, crate::lexicon::ReadState, bool)> = (0..10)
1554 .map(|i| {
1555 let mut state = crate::lexicon::ReadState::new(
1556 format!("https://f{i}.example/feed.xml"),
1557 None,
1558 "2026-07-12T00:00:00.000Z",
1559 );
1560 state.read_ids = (0..crate::lexicon::ReadState::MAX_IDS)
1561 .map(|j| format!("https://f{i}.example/posts/{j:04}/an-entry-permalink"))
1562 .collect();
1563 (format!("rk{i:04}"), state, false)
1564 })
1565 .collect();
1566 repo(&http, &pool, &s, &key)
1567 .flush_read_states(&cursors)
1568 .await
1569 .expect("a byte-heavy flush must succeed in chunks");
1570 let sizes = crate::atproto::tests::call_sizes(&log);
1571 assert!(sizes.len() > 1, "one call for ~500 KB: {sizes:?}");
1572 let want: Vec<String> = cursors.iter().map(|(rkey, _, _)| rkey.clone()).collect();
1573 assert_eq!(crate::atproto::tests::sent_rkeys(&log), want);
1574 }
1575
1576 fn session_at(pds: &str) -> OAuthSession {
1581 let mut s = session();
1582 s.aud = pds.to_string();
1583 s
1584 }
1585
1586 #[tokio::test]
1591 async fn put_record_sends_swap_record_only_when_given() {
1592 use crate::atproto::tests::{serve_status_json, swap_sub, write_ok, OLD_CID};
1593 let (_, pds, log) = serve_status_json(200, write_ok()).await;
1594 let http = Client::new();
1595 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1596 crate::store::init_schema(&pool).await.unwrap();
1597 let key = SigningKey::generate("k");
1598 let s = session_at(&pds);
1599 let repo = repo(&http, &pool, &s, &key);
1600
1601 repo.update_subscription("rk", &swap_sub(), Some(OLD_CID))
1602 .await
1603 .expect("put with a swap");
1604 repo.update_subscription("rk", &swap_sub(), None)
1605 .await
1606 .expect("put without a swap");
1607
1608 let sent = log.lock().unwrap().clone();
1609 assert_eq!(sent.len(), 2, "{sent:?}");
1610 assert_eq!(sent[0]["rkey"], "rk", "captured no usable body: {sent:?}");
1611 assert_eq!(
1612 sent[0]["swapRecord"], OLD_CID,
1613 "the CID the caller read never reached the PDS: {}",
1614 sent[0]
1615 );
1616 assert_eq!(sent[1]["rkey"], "rk");
1617 assert!(
1618 sent[1].get("swapRecord").is_none(),
1619 "no swap was asked for, so none may be sent: {}",
1620 sent[1]
1621 );
1622 }
1623
1624 #[tokio::test]
1627 async fn an_invalid_swap_from_the_pds_is_recognised() {
1628 use crate::atproto::tests::{invalid_swap_xrpc, serve_status_json, swap_sub, OLD_CID};
1629 let http = Client::new();
1630 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1631 crate::store::init_schema(&pool).await.unwrap();
1632 let key = SigningKey::generate("k");
1633
1634 let (_, pds, _) = serve_status_json(400, invalid_swap_xrpc()).await;
1635 let s = session_at(&pds);
1636 let err = repo(&http, &pool, &s, &key)
1637 .update_subscription("rk", &swap_sub(), Some(OLD_CID))
1638 .await
1639 .expect_err("the PDS refused the swap");
1640 assert!(crate::atproto::is_invalid_swap(&err), "{err:#}");
1641
1642 let (_, pds, _) = serve_status_json(
1643 400,
1644 serde_json::json!({ "error": "InvalidRequest", "message": "bad record" }),
1645 )
1646 .await;
1647 let s = session_at(&pds);
1648 let err = repo(&http, &pool, &s, &key)
1649 .update_subscription("rk", &swap_sub(), Some(OLD_CID))
1650 .await
1651 .expect_err("refused");
1652 assert!(!crate::atproto::is_invalid_swap(&err), "{err:#}");
1653 }
1654
1655 #[tokio::test]
1658 async fn list_folders_with_cids_keeps_each_records_cid() {
1659 use crate::atproto::tests::{
1660 assert_folders_listed_with_cids, serve_status_json, two_folders_page,
1661 };
1662 let (_, pds, _) = serve_status_json(200, two_folders_page()).await;
1663 let http = Client::new();
1664 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1665 crate::store::init_schema(&pool).await.unwrap();
1666 let key = SigningKey::generate("k");
1667 let s = session_at(&pds);
1668 let listed = repo(&http, &pool, &s, &key)
1669 .list_folders_with_cids()
1670 .await
1671 .expect("listing");
1672 assert_folders_listed_with_cids(&listed);
1673 }
1674
1675 #[tokio::test]
1677 async fn list_subscriptions_with_cids_keeps_each_records_cid() {
1678 use crate::atproto::tests::{assert_listed_with_cids, serve_status_json, two_subs_page};
1679 let (_, pds, _) = serve_status_json(200, two_subs_page()).await;
1680 let http = Client::new();
1681 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1682 crate::store::init_schema(&pool).await.unwrap();
1683 let key = SigningKey::generate("k");
1684 let s = session_at(&pds);
1685 let listed = repo(&http, &pool, &s, &key)
1686 .list_subscriptions_with_cids()
1687 .await
1688 .expect("listing");
1689 assert_listed_with_cids(&listed);
1690 }
1691}