Skip to main content

feather_reader/oauth/
xrpc.rs

1//! DPoP-bound `com.atproto.repo.*` calls against the user's PDS.
2//!
3//! The layer the 31-method reader surface is built on. It returns the same
4//! [`RecordEntry`] / [`WriteResult`] / [`WriteOp`] types the existing
5//! [`crate::atproto`] client does, so the typed wrappers above it (subscriptions,
6//! folders, saved, read-state) transfer at cutover rather than being rewritten.
7//!
8//! Two properties that are not obvious from the call shapes:
9//!
10//! * **The PDS comes from the session's `aud`**, never re-derived per call. It
11//!   is a property of the token set — a token is valid at one PDS — and it is
12//!   AAD-bound in the session row precisely so it cannot be repointed.
13//! * **Every request carries both the proof and the token.** A resource request
14//!   is `Authorization: DPoP <token>` plus a `DPoP` proof whose `ath` binds to
15//!   that token; either alone is useless.
16
17use 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
29/// Cap on how many pages `list_all_records` will walk.
30///
31/// The same bound the existing client uses: a repo is user-controlled, and an
32/// unbounded walk is a denial-of-service against ourselves.
33const MAX_LIST_PAGES: usize = 50;
34
35/// Hard cap on records accumulated by one walk — 50 pages x the 100 we request.
36/// See [`crate::atproto::extend_bounded`] for why exceeding it is an error
37/// rather than a truncation.
38const MAX_LIST_RECORDS: usize = 5_000;
39
40/// An authenticated handle on one account's repo.
41pub struct Repo<'a> {
42    pub http: &'a Client,
43    pub pool: &'a SqlitePool,
44    pub session: &'a OAuthSession,
45    /// The session's DPoP key, already unsealed.
46    pub key: &'a SigningKey,
47}
48
49impl Repo<'_> {
50    /// The XRPC endpoint for a method, on THIS session's PDS.
51    fn url(&self, nsid: &str) -> String {
52        format!("{}/xrpc/{nsid}", self.session.aud.trim_end_matches('/'))
53    }
54
55    /// Send, and fail loudly on a non-2xx with the XRPC error if there is one.
56    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                // A repo call is repeated on a nonce challenge, on the
72                // ASSUMPTION that a PDS answering with a challenge has not acted
73                // on the request. That is how conformant servers behave and what
74                // the reference relies on, but it is not a guarantee we can
75                // verify: a PDS or intermediary that emitted `use_dpop_nonce`
76                // after a write committed would yield a duplicate record from
77                // `createRecord` or `applyWrites`, both of which are
78                // non-idempotent here.
79                retry: Retry::Allowed,
80            },
81        )
82        .await?;
83
84        if !outcome.is_success() {
85            // **The rejection as a value, under the same sentence as before.**
86            // The sidecar client already surfaces an `AtProtoError::Xrpc`; this
87            // one only said it in a string, so a caller that has to tell "the
88            // PDS refused this" from "the network broke" — the read-state
89            // reconcile (#241) — could match on nothing sturdier than wording.
90            // The context keeps `Display` byte-for-byte what it was.
91            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    /// A non-2xx answer as the [`crate::atproto::AtProtoError::Xrpc`] the
101    /// sidecar client produces for the same answer, so one matcher serves both.
102    ///
103    /// The same length bound as [`Self::error_fields`] before any parse, and the
104    /// same `"Unknown"` fallback `atproto::xrpc_error_from` uses for a body that
105    /// is not an error document.
106    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                // Unreachable: reqwest already parsed this status. Not a 500,
125                // so it can never read as the reference PDS's mismatch.
126                .unwrap_or(reqwest::StatusCode::BAD_GATEWAY),
127            error,
128            message,
129        }
130    }
131
132    /// [`send_raw`](Self::send_raw), then the body as JSON.
133    ///
134    /// Every call that wants a `Value` goes through here; `listRecords` does
135    /// not, because turning its body into a `Value` and then into a page
136    /// materialises the records twice.
137    async fn send(&self, url: &str, body: DpopBody<'_>, nsid: &str) -> Result<Value> {
138        let outcome = self.send_raw(url, body, nsid).await?;
139        // A 200 with an empty body is legitimate for deleteRecord/applyWrites.
140        if outcome.body.is_empty() {
141            return Ok(Value::Null);
142        }
143        outcome.json()
144    }
145
146    /// Render an XRPC error body for a message.
147    ///
148    /// Only ever called on a NON-success, where the body is an error document
149    /// rather than a token or a record — and even then only the `error` and
150    /// `message` fields, never the raw bytes.
151    ///
152    /// **Bounded by length before it is parsed.** This is the error twin of
153    /// `send`, whose success branch goes through `PostOutcome::json` and its node
154    /// guard — so without this a hostile PDS answering **500** instead of 200 got
155    /// the whole amplification the guard exists to stop, on every write and on
156    /// every listing failure. An error document is a few dozen bytes; see
157    /// [`super::MAX_ERROR_BODY`].
158    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    /// One page of a collection.
171    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    /// [`Self::list_records`] with the whole page, including how many records
182    /// were malformed and left out (#177) — which the tuple form cannot carry.
183    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        // **Raw bytes, parsed once.** This is the live backend's list walk, the
204        // largest body this client reads, and routing it through `Value` first
205        // held two copies of every page at the same time.
206        let outcome = self
207            .send_raw(
208                url.as_str(),
209                DpopBody::Query,
210                "com.atproto.repo.listRecords",
211            )
212            .await?;
213        // A 2xx carrying an error envelope is NOT an empty page: `records`
214        // defaulting to `[]` turned a PDS failure into `Ok(vec![])`, which
215        // `resolve_subscriptions` reads as "this DID follows nothing" and
216        // `sync_sub_refs` then writes through, revoking every `sub_ref`.
217        // Both invariants, through the one function every listRecords caller
218        // shares — this check was added here first and had to be fitted to the
219        // other two clients a round later. It reads bytes now rather than a
220        // `Value`; the invariants are the same ones, expressed as fields.
221        crate::atproto::parse_list_records(&outcome.body)
222    }
223
224    /// Every record in a collection, following the cursor.
225    ///
226    /// Bounded at `MAX_LIST_PAGES`: the repo is user-controlled, so an
227    /// unbounded walk is a denial-of-service against ourselves. A cursor that
228    /// does not advance on a non-empty page is refused rather than followed or
229    /// stopped on: either would hand back an unfinished list (#203).
230    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    /// [`list_all_records`](Self::list_all_records) against a caller's budget.
239    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            // This walk is the live backend's, and its result reaches
254            // `replace_sub_refs`: a skipped record is a dropped subscription.
255            crate::atproto::refuse_malformed(&listed, collection)?;
256            let (page, next) = (listed.records, listed.cursor);
257            let got = page.len();
258            // Same cap, same reason as `atproto::extend_bounded`: MAX_LIST_PAGES
259            // bounds requests, not memory, unless the server honours our limit.
260            // This is the LIVE walk on `backend=rust` — its result reaches
261            // `replace_sub_refs`, so a truncation here is revoked access.
262            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                // `got > 0` is not defensive tidiness -- it is a whole round
274                // trip. This project's own PDS returns a cursor ALONGSIDE a
275                // short page, so without it every list fetches a second, empty
276                // page before stopping. Both of the other clients have this
277                // guard; measured, its absence was most of the remaining gap
278                // against the sidecar.
279                //
280                // The cursor-repeat check is the separate concern: a server that
281                // hands back the same cursor forever would otherwise loop. It is
282                // REFUSED rather than ended on (#203): stopping there hands
283                // `replace_sub_refs` a list the server said was unfinished.
284                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                    // `break` with the flag cleared, rather than an early
297                    // `return`: an early return makes the post-loop check
298                    // unreachable, so the flag reads as dead state and a
299                    // reviewer hunts for the case that clears it. It also left
300                    // the "every walk refuses" mutation alive here.
301                    more_offered = false;
302                    break;
303                }
304            }
305        }
306        // **Running out of pages is a refusal, not a short answer.** Falling out
307        // of the loop used to return `Ok(out)`, so a repo bigger than the page
308        // budget produced a truncated list indistinguishable from a complete
309        // one — and `resolve_subscriptions` needs an `Err` for its fail-closed
310        // branch. Given `Ok`, it hands the short list to `replace_sub_refs`,
311        // which DELETEs the reader's whole `sub_ref` projection and reinserts
312        // only what it was given. `extend_bounded` cannot catch this either:
313        // `MAX_LIST_PAGES` x the 100 we request is `MAX_LIST_RECORDS`, so against
314        // a server that honours our limit the page budget runs out first.
315        //
316        // **The cap is on REQUESTS, though, so the record count it bites at is
317        // the server's page size x the budget — not a number our constants fix.**
318        // A PDS answering 50 a page reaches half as far; one answering more than
319        // asked trips `extend_bounded` instead. And the last allowed page is a
320        // FALSE refusal: terminating costs one extra request when a short page
321        // still carries a cursor, so a walk holding every record it will ever
322        // hold still refuses on a cursor it never followed.
323        //
324        // This client's caps are a QUARTER of the direct client's (50 pages,
325        // 5 000 records), so with `limit=100` honoured the same reader refuses
326        // here at ~4 900 records and works to ~19 900 on the sidecar. `Saved`
327        // walks this too, one record per starred article, where 4 900 is a
328        // plausible number for a real reader.
329        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    /// Create a record, letting the PDS assign the key.
341    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    /// Create or replace a record at a known key.
354    ///
355    /// `swap_record` is the CID the caller read the record at, sent as
356    /// `swapRecord`, so a record another client wrote in between is refused
357    /// with `InvalidSwap` instead of overwritten (#149); recognise that with
358    /// [`crate::atproto::is_invalid_swap`]. `None` omits the field.
359    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    /// A batch of writes, in as few round trips as the PDS's limits allow.
391    ///
392    /// Chunked by `crate::atproto::apply_writes_chunked`, which says what a
393    /// failure part-way means: the batch is atomic per CALL, not as a whole.
394    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    /// One `applyWrites` call, unchunked. Private: an oversized call is a
399    /// refused call, so nothing reaches this except through
400    /// [`apply_writes`](Self::apply_writes).
401    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
426/// Format an XRPC failure for an error message.
427fn 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
434/// The reader's typed surface, over [`Repo`].
435///
436/// Deliberately thin: each method is one repo call plus a parse, and the
437/// orderings come from [`crate::lexicon::sort`], SHARED with the sidecar client
438/// so the two cannot disagree across the cutover. A divergence there would not
439/// be subtle — it would reorder the user's feed list the moment the
440/// implementation swapped.
441impl Repo<'_> {
442    /// List a collection and parse each record into `T`, paired with its rkey.
443    ///
444    /// An unparseable record is SKIPPED with a warning rather than failing the
445    /// list. Records are written by other clients and by future versions of this
446    /// one; one record this build cannot read must not black out the whole feed
447    /// list.
448    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    /// [`list_typed`](Self::list_typed), keeping the CID each record was
461    /// listed at — the value a compare-and-swap write names (#149). One parse
462    /// path for both, so the CID listing skips exactly what the plain one does.
463    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    // ── subscriptions ────────────────────────────────────────────────────────
485
486    pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
487        self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
488    }
489
490    /// Every subscription with the CID it was listed at, unsorted — the read
491    /// half of a read-modify-write that puts with `swapRecord` (#149).
492    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    /// Every subscription, in the reader's deterministic order.
500    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    /// Subscribe to a feed. Returns the new record's rkey so the caller can
509    /// address it (rename, delete) without re-listing.
510    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    /// Replace a subscription in place — retitle, refile, change cadence.
526    /// Returns the [`WriteResult`] rather than discarding it: the caller logs
527    /// the resulting record URI, and a `()` here would have silently dropped
528    /// that field from the log line after the cutover.
529    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    /// Batch-add many subscriptions via `applyWrites` (chunked) — the OPML-import path.
540    ///
541    /// One round trip per 200 feeds rather than one per feed: an import of
542    /// several hundred feeds is the case this exists for. More than one call
543    /// means the import can part-land; on an error,
544    /// [`crate::atproto::ApplyWritesIncomplete::of`] gives `landed`, and the
545    /// first `landed` of `subs` are in the repo.
546    /// Returns the new rkeys, which are assigned HERE rather than by the server:
547    /// client-side TIDs keep the imported feeds in input order and make the
548    /// batch reproducible. The sidecar client does the same, with the same
549    /// generator.
550    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                // Propagated, NOT defaulted: `unwrap_or(Value::Null)` here would
563                // write a null record into the user's repo on a serialization
564                // failure rather than failing the import.
565                value: serde_json::to_value(sub)?,
566            });
567            rkeys.push(rkey);
568        }
569        self.apply_writes(&writes).await?;
570        Ok(rkeys)
571    }
572
573    // ── folders ──────────────────────────────────────────────────────────────
574
575    pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
576        self.list_typed(crate::lexicon::nsid::FOLDER).await
577    }
578
579    /// Every folder with the CID it was listed at, unsorted (#268).
580    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    /// Delete a folder. Subscriptions referencing it are left alone; a dangling
601    /// reference reads as "unfiled", which is the same behaviour the sidecar
602    /// client has.
603    pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
604        self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
605    }
606
607    /// Returns the [`WriteResult`], matching the sidecar client — the caller
608    /// logs the record URI from it. `swap_record` is the CID the folder was
609    /// read at, so a write by another client since is refused rather than
610    /// overwritten (#268).
611    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    // ── saved ────────────────────────────────────────────────────────────────
622
623    pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
624        self.list_typed(crate::lexicon::nsid::SAVED).await
625    }
626
627    /// Saved entries, newest first.
628    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    // ── read state ───────────────────────────────────────────────────────────
646
647    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    /// Upsert one read cursor at its feed-derived rkey.
652    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    /// Flush many dirty read cursors via `applyWrites` (chunked).
663    ///
664    /// Read state changes on nearly every page view, so this is the hottest
665    /// write path in the app; one round trip per flush rather than per feed is
666    /// the whole point. A flush past the op or byte bound is several calls,
667    /// and a failure part-way leaves the first
668    /// [`landed`](crate::atproto::ApplyWritesIncomplete::landed) cursors
669    /// written — see `crate::atproto::apply_writes_chunked`.
670    /// The `bool` is whether the record already exists in the PDS. It is not
671    /// optional bookkeeping: an `#update` on a missing record ERRORS, and
672    /// `applyWrites` is atomic per repo, so one not-yet-created cursor in the
673    /// batch would drop the whole DID's flush. The op builder is shared with the
674    /// sidecar client so the two cannot decide create-vs-update differently.
675    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    /// **Both clients must build the SAME write ops.**
692    ///
693    /// `flush_read_states` is the hottest write path in the app, and the
694    /// create-vs-update choice is the one part of it that cannot be got wrong
695    /// quietly: an `#update` on a record that does not exist errors, and
696    /// `applyWrites` is atomic per repo, so a single first-flush cursor in the
697    /// batch takes the whole DID's flush down with it.
698    ///
699    /// The first version of this method here ignored the flag and emitted
700    /// `#create` unconditionally, which would have broken every feed's first
701    /// flush after the cutover. Sharing the builder is what makes that
702    /// impossible rather than merely fixed.
703    #[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    /// **Bulk import, through the real `Repo`, asserted on the bytes.** The
730    /// test this replaces exercised `TidGenerator` directly and never called
731    /// `add_subscriptions_bulk`; the function could stop assigning rkeys and
732    /// return garbage with the suite green.
733    #[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    /// **The PDS comes from the session**, so a call cannot be pointed at
810    /// another host by a caller who passes the wrong base.
811    #[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    /// An XRPC error document is summarised by its `error`/`message` fields
827    /// only — never by echoing the raw body, which on other paths holds tokens.
828    #[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        // A body that is not an XRPC error contributes nothing but the status.
836        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    /// Every repo call must fail closed on an internal target, like every other
843    /// outbound path — asserted on the guard's own error.
844    #[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    /// **A 200 carrying an error envelope must not read as an empty repo.**
865    ///
866    /// `records` was taken off the JSON with `unwrap_or(Array([]))`, so a PDS
867    /// answering `200 {"error": …}` produced `Ok(vec![])`. That is not the
868    /// fail-closed branch in `web::resolve_subscriptions`: `sync_sub_refs`
869    /// writes the empty set through and `replace_sub_refs` DELETEs the DID's
870    /// entire `sub_ref` projection — one bad response revokes the reader's
871    /// access to every feed they have. Driven through the real client against
872    /// a real server, because the bug was the missing CALL, not the check.
873    #[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    /// **The live walk spends its budget across pages too.**
910    ///
911    /// This is the `backend=rust` walk whose result reaches `replace_sub_refs`,
912    /// so a bound that silently failed to accumulate here would revoke a
913    /// reader's access to every feed past the cut. Three pages, a two-page
914    /// budget: the walk must refuse, and must say it kept two.
915    #[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    /// **The live walk refuses a short list.** Its result reaches
960    /// `replace_sub_refs`, so returning a truncated list as a complete one
961    /// deletes every subscription past the page budget.
962    #[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    /// #177 on the live backend: a malformed record in the reader's own repo
1010    /// refuses the walk, by type, rather than dropping that subscription.
1011    #[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    /// The live walk's clean finish must still be a success — see the sidecar's
1050    /// twin for why this direction is the dangerous one.
1051    #[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        // **The terminator CARRIES a record**, because a real PDS ends on a
1064        // partial page and those are the records an off-by-one loses. Ending on
1065        // an empty page kept "drop the last page's records" alive: the count
1066        // below was right either way.
1067        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    /// **The page cap is pinned exactly, not to within one.**
1107    ///
1108    /// `the_live_walk_that_runs_out_of_pages_refuses` serves `MAX_LIST_PAGES + 1`
1109    /// pages, so a budget one page SHORT refuses too and that mutation survives
1110    /// it. A walk whose last allowed request is the terminating one must come
1111    /// back `Ok` — and on this backend an `Err` is `replace_sub_refs` never
1112    /// running, which is the direction that costs a reader their subscriptions.
1113    #[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    /// **The live walk refuses a repeated cursor (#203).** A server that hands
1165    /// back the same cursor with a non-empty page said it had more and then
1166    /// did not move. Stopping there with `Ok` handed a short list to
1167    /// `replace_sub_refs`; following it would loop. Refusing is the only exit
1168    /// that neither spins nor deletes.
1169    #[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    /// **The ERROR twin of `send`, which a review found unguarded.**
1208    ///
1209    /// `send`'s success branch goes through `PostOutcome::json` and its cap; its
1210    /// failure branch goes to `xrpc_error` → `error_fields`, which deserialised the
1211    /// whole body. So a hostile PDS answering **500** instead of 200 with the same
1212    /// explosion got the full amplification on every write and on every listing
1213    /// failure — the guard bypassed by a status code.
1214    ///
1215    /// Both directions: a small error body still yields its reason, because that
1216    /// string is what a reader's log needs to tell "your PDS said no" from "we
1217    /// broke".
1218    #[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    /// **A PDS rejection is structured, not only a sentence (#241).**
1241    ///
1242    /// The read-state flusher has to tell "the PDS refused this batch" from "the
1243    /// network broke", and the sidecar client already said so with
1244    /// [`crate::atproto::AtProtoError::Xrpc`]. This client only said it in a
1245    /// string, so the one backend production runs could be matched on nothing
1246    /// sturdier than its wording. The rendered message must not change: it is
1247    /// what every existing log line and assertion reads.
1248    #[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    /// **The write path parses a PDS body too, and it had no bound before the
1300    /// parse.**
1301    ///
1302    /// `listRecords` is guarded by `parse_list_records`; every OTHER body this
1303    /// client turns into a `Value` goes through `PostOutcome::json` — the repo
1304    /// writers' responses, the PAR response, the token response, the session
1305    /// refresh. `read_capped` bounds the WIRE at 8 MB, which is the *input* to the
1306    /// amplification rather than a limit on it: 8 MB of the cheapest node shape
1307    /// measured 824 MB retained on a 512 MB box.
1308    ///
1309    /// Drives `delete_record`, which is the shortest route from a handler to
1310    /// `send` → `json`.
1311    #[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    /// The other direction: an ordinary write response still parses. Without this
1354    /// a guard that refused every body would pass the test above.
1355    #[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    /// **A duplicated `records` key must not be able to empty a page.**
1385    ///
1386    /// This is the live `backend=rust` walk, so an empty page here reaches
1387    /// `replace_sub_refs` and deletes the reader's subscriptions. serde refuses
1388    /// a repeated field outright; a `serde_json::Value` takes the last one
1389    /// silently, so routing this body through a `Value` first turns a smuggled
1390    /// second key into a successful, empty listing. The test exists to pin
1391    /// which of the two this client uses.
1392    #[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    /// **An empty 2xx body is not an empty repo either.** `send` maps a
1429    /// zero-length 2xx to `Value::Null` — deliberately, for `deleteRecord` and
1430    /// `applyWrites` — and while `listRecords` still went through it, that
1431    /// slipped past the envelope guard
1432    /// and `records.unwrap_or(Array([]))` produced `Ok(vec![])`: the same
1433    /// `sub_ref` wipe the guard was added to prevent, through the sibling door.
1434    #[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    /// The write paths discarded the body too: `delete_record` and
1467    /// `apply_writes` are `send(..).await?; Ok(())`, so a 200 carrying an
1468    /// error envelope reported success for a delete that did not happen.
1469    #[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    /// An empty batch must not produce a request at all — an `applyWrites` with
1511    /// no writes is a round trip that can only fail.
1512    #[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        // A target that would fail loudly if it were ever contacted.
1519        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    // ── applyWrites chunking (#240) ──────────────────────────────────────────
1525    //
1526    // This is the production backend, so the chunking is asserted here on the
1527    // bytes it sends to a fake PDS that refuses what the reference PDS refuses
1528    // (see `crate::atproto::tests::serve_apply_writes`).
1529
1530    /// A session pointed at the strict fake, reached through a per-port host
1531    /// override (the override table is process-wide and tests run in parallel).
1532    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    /// **An OPML import of 201 feeds is two calls, 200 then 1, in order** — on
1562    /// the backend production runs. One call of 201 is what the reference PDS
1563    /// refuses with `Too many writes. Max: 200`.
1564    #[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    /// 500 feeds — the default per-DID cap — are three calls; exactly 200 is one.
1581    #[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    /// **A failed chunk stops the run**, and chunk 3 is never sent.
1595    #[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    /// **The byte bound splits a read-state flush well under 200 ops.** Ten
1612    /// cursors at the lexicon's 1,000-id cap are ~500 KB — one call of that is
1613    /// refused by every reference PDS older than atproto#4989.
1614    #[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    // -- #149: compare-and-swap putRecord ------------------------------------
1642
1643    /// A session whose audience is `pds` — the one place this client takes
1644    /// its host from.
1645    fn session_at(pds: &str) -> OAuthSession {
1646        let mut s = session();
1647        s.aud = pds.to_string();
1648        s
1649    }
1650
1651    /// **`swapRecord` is on the wire when given, and absent when not** — the
1652    /// live backend's half of the CAS. Asserted on the request body the PDS
1653    /// received, since a parameter accepted and then dropped is exactly what
1654    /// a happy-path test cannot see.
1655    #[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    /// The PDS's `InvalidSwap`, as this client reports it, is recognised —
1690    /// and a different 400 from the same client is not.
1691    #[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    /// The folder CID listing pairs each folder with the CID it was listed at
1721    /// (#268).
1722    #[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    /// The CID listing pairs each record with the CID it was listed at.
1741    #[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}