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 also terminates the walk rather than spinning.
229    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    /// [`list_all_records`](Self::list_all_records) against a caller's budget.
238    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            // This walk is the live backend's, and its result reaches
253            // `replace_sub_refs`: a skipped record is a dropped subscription.
254            crate::atproto::refuse_malformed(&listed, collection)?;
255            let (page, next) = (listed.records, listed.cursor);
256            let got = page.len();
257            // Same cap, same reason as `atproto::extend_bounded`: MAX_LIST_PAGES
258            // bounds requests, not memory, unless the server honours our limit.
259            // This is the LIVE walk on `backend=rust` — its result reaches
260            // `replace_sub_refs`, so a truncation here is revoked access.
261            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                // `got > 0` is not defensive tidiness -- it is a whole round
272                // trip. This project's own PDS returns a cursor ALONGSIDE a
273                // short page, so without it every list fetches a second, empty
274                // page before stopping. Both of the other clients have this
275                // guard; measured, its absence was most of the remaining gap
276                // against the sidecar.
277                //
278                // The cursor-repeat check is the separate concern: a server that
279                // hands back the same cursor forever would otherwise loop.
280                Some(next) if got > 0 && Some(&next) != cursor.as_ref() => {
281                    cursor = Some(next);
282                    more_offered = true;
283                }
284                _ => {
285                    // `break` with the flag cleared, rather than an early
286                    // `return`: an early return makes the post-loop check
287                    // unreachable, so the flag reads as dead state and a
288                    // reviewer hunts for the case that clears it. It also left
289                    // the "every walk refuses" mutation alive here.
290                    more_offered = false;
291                    break;
292                }
293            }
294        }
295        // **Running out of pages is a refusal, not a short answer.** Falling out
296        // of the loop used to return `Ok(out)`, so a repo bigger than the page
297        // budget produced a truncated list indistinguishable from a complete
298        // one — and `resolve_subscriptions` needs an `Err` for its fail-closed
299        // branch. Given `Ok`, it hands the short list to `replace_sub_refs`,
300        // which DELETEs the reader's whole `sub_ref` projection and reinserts
301        // only what it was given. `extend_bounded` cannot catch this either:
302        // `MAX_LIST_PAGES` x the 100 we request is `MAX_LIST_RECORDS`, so against
303        // a server that honours our limit the page budget runs out first.
304        //
305        // **The cap is on REQUESTS, though, so the record count it bites at is
306        // the server's page size x the budget — not a number our constants fix.**
307        // A PDS answering 50 a page reaches half as far; one answering more than
308        // asked trips `extend_bounded` instead. And the last allowed page is a
309        // FALSE refusal: terminating costs one extra request when a short page
310        // still carries a cursor, so a walk holding every record it will ever
311        // hold still refuses on a cursor it never followed.
312        //
313        // This client's caps are a QUARTER of the direct client's (50 pages,
314        // 5 000 records), so with `limit=100` honoured the same reader refuses
315        // here at ~4 900 records and works to ~19 900 on the sidecar. `Saved`
316        // walks this too, one record per starred article, where 4 900 is a
317        // plausible number for a real reader.
318        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    /// Create a record, letting the PDS assign the key.
329    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    /// Create or replace a record at a known key.
342    ///
343    /// `swap_record` is the CID the caller read the record at, sent as
344    /// `swapRecord`, so a record another client wrote in between is refused
345    /// with `InvalidSwap` instead of overwritten (#149); recognise that with
346    /// [`crate::atproto::is_invalid_swap`]. `None` omits the field.
347    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    /// A batch of writes, in as few round trips as the PDS's limits allow.
379    ///
380    /// Chunked by `crate::atproto::apply_writes_chunked`, which says what a
381    /// failure part-way means: the batch is atomic per CALL, not as a whole.
382    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    /// One `applyWrites` call, unchunked. Private: an oversized call is a
387    /// refused call, so nothing reaches this except through
388    /// [`apply_writes`](Self::apply_writes).
389    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
414/// Format an XRPC failure for an error message.
415fn 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
422/// The reader's typed surface, over [`Repo`].
423///
424/// Deliberately thin: each method is one repo call plus a parse, and the
425/// orderings come from [`crate::lexicon::sort`], SHARED with the sidecar client
426/// so the two cannot disagree across the cutover. A divergence there would not
427/// be subtle — it would reorder the user's feed list the moment the
428/// implementation swapped.
429impl Repo<'_> {
430    /// List a collection and parse each record into `T`, paired with its rkey.
431    ///
432    /// An unparseable record is SKIPPED with a warning rather than failing the
433    /// list. Records are written by other clients and by future versions of this
434    /// one; one record this build cannot read must not black out the whole feed
435    /// list.
436    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    /// [`list_typed`](Self::list_typed), keeping the CID each record was
449    /// listed at — the value a compare-and-swap write names (#149). One parse
450    /// path for both, so the CID listing skips exactly what the plain one does.
451    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    // ── subscriptions ────────────────────────────────────────────────────────
473
474    pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
475        self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
476    }
477
478    /// Every subscription with the CID it was listed at, unsorted — the read
479    /// half of a read-modify-write that puts with `swapRecord` (#149).
480    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    /// Every subscription, in the reader's deterministic order.
488    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    /// Subscribe to a feed. Returns the new record's rkey so the caller can
497    /// address it (rename, delete) without re-listing.
498    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    /// Replace a subscription in place — retitle, refile, change cadence.
514    /// Returns the [`WriteResult`] rather than discarding it: the caller logs
515    /// the resulting record URI, and a `()` here would have silently dropped
516    /// that field from the log line after the cutover.
517    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    /// Batch-add many subscriptions via `applyWrites` (chunked) — the OPML-import path.
528    ///
529    /// One round trip per 200 feeds rather than one per feed: an import of
530    /// several hundred feeds is the case this exists for. More than one call
531    /// means the import can part-land; on an error,
532    /// [`crate::atproto::ApplyWritesIncomplete::of`] gives `landed`, and the
533    /// first `landed` of `subs` are in the repo.
534    /// Returns the new rkeys, which are assigned HERE rather than by the server:
535    /// client-side TIDs keep the imported feeds in input order and make the
536    /// batch reproducible. The sidecar client does the same, with the same
537    /// generator.
538    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                // Propagated, NOT defaulted: `unwrap_or(Value::Null)` here would
551                // write a null record into the user's repo on a serialization
552                // failure rather than failing the import.
553                value: serde_json::to_value(sub)?,
554            });
555            rkeys.push(rkey);
556        }
557        self.apply_writes(&writes).await?;
558        Ok(rkeys)
559    }
560
561    // ── folders ──────────────────────────────────────────────────────────────
562
563    pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
564        self.list_typed(crate::lexicon::nsid::FOLDER).await
565    }
566
567    /// Every folder with the CID it was listed at, unsorted (#268).
568    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    /// Delete a folder. Subscriptions referencing it are left alone; a dangling
589    /// reference reads as "unfiled", which is the same behaviour the sidecar
590    /// client has.
591    pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
592        self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
593    }
594
595    /// Returns the [`WriteResult`], matching the sidecar client — the caller
596    /// logs the record URI from it. `swap_record` is the CID the folder was
597    /// read at, so a write by another client since is refused rather than
598    /// overwritten (#268).
599    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    // ── saved ────────────────────────────────────────────────────────────────
610
611    pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
612        self.list_typed(crate::lexicon::nsid::SAVED).await
613    }
614
615    /// Saved entries, newest first.
616    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    // ── read state ───────────────────────────────────────────────────────────
634
635    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    /// Upsert one read cursor at its feed-derived rkey.
640    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    /// Flush many dirty read cursors via `applyWrites` (chunked).
651    ///
652    /// Read state changes on nearly every page view, so this is the hottest
653    /// write path in the app; one round trip per flush rather than per feed is
654    /// the whole point. A flush past the op or byte bound is several calls,
655    /// and a failure part-way leaves the first
656    /// [`landed`](crate::atproto::ApplyWritesIncomplete::landed) cursors
657    /// written — see `crate::atproto::apply_writes_chunked`.
658    /// The `bool` is whether the record already exists in the PDS. It is not
659    /// optional bookkeeping: an `#update` on a missing record ERRORS, and
660    /// `applyWrites` is atomic per repo, so one not-yet-created cursor in the
661    /// batch would drop the whole DID's flush. The op builder is shared with the
662    /// sidecar client so the two cannot decide create-vs-update differently.
663    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    /// **Both clients must build the SAME write ops.**
680    ///
681    /// `flush_read_states` is the hottest write path in the app, and the
682    /// create-vs-update choice is the one part of it that cannot be got wrong
683    /// quietly: an `#update` on a record that does not exist errors, and
684    /// `applyWrites` is atomic per repo, so a single first-flush cursor in the
685    /// batch takes the whole DID's flush down with it.
686    ///
687    /// The first version of this method here ignored the flag and emitted
688    /// `#create` unconditionally, which would have broken every feed's first
689    /// flush after the cutover. Sharing the builder is what makes that
690    /// impossible rather than merely fixed.
691    #[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    /// **Bulk import, through the real `Repo`, asserted on the bytes.** The
718    /// test this replaces exercised `TidGenerator` directly and never called
719    /// `add_subscriptions_bulk`; the function could stop assigning rkeys and
720    /// return garbage with the suite green.
721    #[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    /// **The PDS comes from the session**, so a call cannot be pointed at
798    /// another host by a caller who passes the wrong base.
799    #[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    /// An XRPC error document is summarised by its `error`/`message` fields
815    /// only — never by echoing the raw body, which on other paths holds tokens.
816    #[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        // A body that is not an XRPC error contributes nothing but the status.
824        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    /// Every repo call must fail closed on an internal target, like every other
831    /// outbound path — asserted on the guard's own error.
832    #[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    /// **A 200 carrying an error envelope must not read as an empty repo.**
853    ///
854    /// `records` was taken off the JSON with `unwrap_or(Array([]))`, so a PDS
855    /// answering `200 {"error": …}` produced `Ok(vec![])`. That is not the
856    /// fail-closed branch in `web::resolve_subscriptions`: `sync_sub_refs`
857    /// writes the empty set through and `replace_sub_refs` DELETEs the DID's
858    /// entire `sub_ref` projection — one bad response revokes the reader's
859    /// access to every feed they have. Driven through the real client against
860    /// a real server, because the bug was the missing CALL, not the check.
861    #[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    /// **The live walk spends its budget across pages too.**
898    ///
899    /// This is the `backend=rust` walk whose result reaches `replace_sub_refs`,
900    /// so a bound that silently failed to accumulate here would revoke a
901    /// reader's access to every feed past the cut. Three pages, a two-page
902    /// budget: the walk must refuse, and must say it kept two.
903    #[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    /// **The live walk refuses a short list.** Its result reaches
943    /// `replace_sub_refs`, so returning a truncated list as a complete one
944    /// deletes every subscription past the page budget.
945    #[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    /// #177 on the live backend: a malformed record in the reader's own repo
988    /// refuses the walk, by type, rather than dropping that subscription.
989    #[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    /// The live walk's clean finish must still be a success — see the sidecar's
1028    /// twin for why this direction is the dangerous one.
1029    #[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        // **The terminator CARRIES a record**, because a real PDS ends on a
1042        // partial page and those are the records an off-by-one loses. Ending on
1043        // an empty page kept "drop the last page's records" alive: the count
1044        // below was right either way.
1045        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    /// **The page cap is pinned exactly, not to within one.**
1085    ///
1086    /// `the_live_walk_that_runs_out_of_pages_refuses` serves `MAX_LIST_PAGES + 1`
1087    /// pages, so a budget one page SHORT refuses too and that mutation survives
1088    /// it. A walk whose last allowed request is the terminating one must come
1089    /// back `Ok` — and on this backend an `Err` is `replace_sub_refs` never
1090    /// running, which is the direction that costs a reader their subscriptions.
1091    #[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    /// **The ERROR twin of `send`, which a review found unguarded.**
1143    ///
1144    /// `send`'s success branch goes through `PostOutcome::json` and its cap; its
1145    /// failure branch goes to `xrpc_error` → `error_fields`, which deserialised the
1146    /// whole body. So a hostile PDS answering **500** instead of 200 with the same
1147    /// explosion got the full amplification on every write and on every listing
1148    /// failure — the guard bypassed by a status code.
1149    ///
1150    /// Both directions: a small error body still yields its reason, because that
1151    /// string is what a reader's log needs to tell "your PDS said no" from "we
1152    /// broke".
1153    #[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    /// **A PDS rejection is structured, not only a sentence (#241).**
1176    ///
1177    /// The read-state flusher has to tell "the PDS refused this batch" from "the
1178    /// network broke", and the sidecar client already said so with
1179    /// [`crate::atproto::AtProtoError::Xrpc`]. This client only said it in a
1180    /// string, so the one backend production runs could be matched on nothing
1181    /// sturdier than its wording. The rendered message must not change: it is
1182    /// what every existing log line and assertion reads.
1183    #[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    /// **The write path parses a PDS body too, and it had no bound before the
1235    /// parse.**
1236    ///
1237    /// `listRecords` is guarded by `parse_list_records`; every OTHER body this
1238    /// client turns into a `Value` goes through `PostOutcome::json` — the repo
1239    /// writers' responses, the PAR response, the token response, the session
1240    /// refresh. `read_capped` bounds the WIRE at 8 MB, which is the *input* to the
1241    /// amplification rather than a limit on it: 8 MB of the cheapest node shape
1242    /// measured 824 MB retained on a 512 MB box.
1243    ///
1244    /// Drives `delete_record`, which is the shortest route from a handler to
1245    /// `send` → `json`.
1246    #[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    /// The other direction: an ordinary write response still parses. Without this
1289    /// a guard that refused every body would pass the test above.
1290    #[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    /// **A duplicated `records` key must not be able to empty a page.**
1320    ///
1321    /// This is the live `backend=rust` walk, so an empty page here reaches
1322    /// `replace_sub_refs` and deletes the reader's subscriptions. serde refuses
1323    /// a repeated field outright; a `serde_json::Value` takes the last one
1324    /// silently, so routing this body through a `Value` first turns a smuggled
1325    /// second key into a successful, empty listing. The test exists to pin
1326    /// which of the two this client uses.
1327    #[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    /// **An empty 2xx body is not an empty repo either.** `send` maps a
1364    /// zero-length 2xx to `Value::Null` — deliberately, for `deleteRecord` and
1365    /// `applyWrites` — and while `listRecords` still went through it, that
1366    /// slipped past the envelope guard
1367    /// and `records.unwrap_or(Array([]))` produced `Ok(vec![])`: the same
1368    /// `sub_ref` wipe the guard was added to prevent, through the sibling door.
1369    #[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    /// The write paths discarded the body too: `delete_record` and
1402    /// `apply_writes` are `send(..).await?; Ok(())`, so a 200 carrying an
1403    /// error envelope reported success for a delete that did not happen.
1404    #[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    /// An empty batch must not produce a request at all — an `applyWrites` with
1446    /// no writes is a round trip that can only fail.
1447    #[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        // A target that would fail loudly if it were ever contacted.
1454        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    // ── applyWrites chunking (#240) ──────────────────────────────────────────
1460    //
1461    // This is the production backend, so the chunking is asserted here on the
1462    // bytes it sends to a fake PDS that refuses what the reference PDS refuses
1463    // (see `crate::atproto::tests::serve_apply_writes`).
1464
1465    /// A session pointed at the strict fake, reached through a per-port host
1466    /// override (the override table is process-wide and tests run in parallel).
1467    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    /// **An OPML import of 201 feeds is two calls, 200 then 1, in order** — on
1497    /// the backend production runs. One call of 201 is what the reference PDS
1498    /// refuses with `Too many writes. Max: 200`.
1499    #[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    /// 500 feeds — the default per-DID cap — are three calls; exactly 200 is one.
1516    #[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    /// **A failed chunk stops the run**, and chunk 3 is never sent.
1530    #[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    /// **The byte bound splits a read-state flush well under 200 ops.** Ten
1547    /// cursors at the lexicon's 1,000-id cap are ~500 KB — one call of that is
1548    /// refused by every reference PDS older than atproto#4989.
1549    #[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    // -- #149: compare-and-swap putRecord ------------------------------------
1577
1578    /// A session whose audience is `pds` — the one place this client takes
1579    /// its host from.
1580    fn session_at(pds: &str) -> OAuthSession {
1581        let mut s = session();
1582        s.aud = pds.to_string();
1583        s
1584    }
1585
1586    /// **`swapRecord` is on the wire when given, and absent when not** — the
1587    /// live backend's half of the CAS. Asserted on the request body the PDS
1588    /// received, since a parameter accepted and then dropped is exactly what
1589    /// a happy-path test cannot see.
1590    #[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    /// The PDS's `InvalidSwap`, as this client reports it, is recognised —
1625    /// and a different 400 from the same client is not.
1626    #[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    /// The folder CID listing pairs each folder with the CID it was listed at
1656    /// (#268).
1657    #[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    /// The CID listing pairs each record with the CID it was listed at.
1676    #[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}