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    pub async fn put_record<T: crate::vetted::WritableRecord>(
343        &self,
344        collection: &str,
345        rkey: &str,
346        record: &T,
347    ) -> Result<WriteResult> {
348        self.write(
349            "com.atproto.repo.putRecord",
350            json!({
351                "repo": self.session.sub,
352                "collection": collection,
353                "rkey": rkey,
354                "record": record,
355            }),
356        )
357        .await
358    }
359
360    pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
361        let body = json!({ "repo": self.session.sub, "collection": collection, "rkey": rkey });
362        self.send(
363            &self.url("com.atproto.repo.deleteRecord"),
364            DpopBody::Json(serde_json::to_vec(&body)?),
365            "com.atproto.repo.deleteRecord",
366        )
367        .await
368        .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
369        Ok(())
370    }
371
372    /// A batch of writes, in as few round trips as the PDS's limits allow.
373    ///
374    /// Chunked by `crate::atproto::apply_writes_chunked`, which says what a
375    /// failure part-way means: the batch is atomic per CALL, not as a whole.
376    pub async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
377        crate::atproto::apply_writes_chunked(writes, |chunk| self.apply_writes_once(chunk)).await
378    }
379
380    /// One `applyWrites` call, unchunked. Private: an oversized call is a
381    /// refused call, so nothing reaches this except through
382    /// [`apply_writes`](Self::apply_writes).
383    async fn apply_writes_once(&self, writes: &[WriteOp]) -> Result<()> {
384        let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
385        let body = json!({ "repo": self.session.sub, "writes": ops });
386        self.send(
387            &self.url("com.atproto.repo.applyWrites"),
388            DpopBody::Json(serde_json::to_vec(&body)?),
389            "com.atproto.repo.applyWrites",
390        )
391        .await
392        .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
393        Ok(())
394    }
395
396    async fn write(&self, nsid: &str, body: Value) -> Result<WriteResult> {
397        let value = self
398            .send(
399                &self.url(nsid),
400                DpopBody::Json(serde_json::to_vec(&body)?),
401                nsid,
402            )
403            .await?;
404        serde_json::from_value(value).with_context(|| format!("{nsid} returned no usable result"))
405    }
406}
407
408/// Format an XRPC failure for an error message.
409fn xrpc_error(body: &[u8], status: u16) -> String {
410    match Repo::error_fields(body) {
411        Some(detail) => format!("status {status} ({detail})"),
412        None => format!("status {status}"),
413    }
414}
415
416/// The reader's typed surface, over [`Repo`].
417///
418/// Deliberately thin: each method is one repo call plus a parse, and the
419/// orderings come from [`crate::lexicon::sort`], SHARED with the sidecar client
420/// so the two cannot disagree across the cutover. A divergence there would not
421/// be subtle — it would reorder the user's feed list the moment the
422/// implementation swapped.
423impl Repo<'_> {
424    /// List a collection and parse each record into `T`, paired with its rkey.
425    ///
426    /// An unparseable record is SKIPPED with a warning rather than failing the
427    /// list. Records are written by other clients and by future versions of this
428    /// one; one record this build cannot read must not black out the whole feed
429    /// list.
430    async fn list_typed<T: serde::de::DeserializeOwned>(
431        &self,
432        collection: &str,
433    ) -> Result<Vec<(String, T)>> {
434        let records = self.list_all_records(collection).await?;
435        let mut out = Vec::with_capacity(records.len());
436        for record in records {
437            let rkey = record.rkey().unwrap_or_default().to_string();
438            match record.parse::<T>() {
439                Ok(value) => out.push((rkey, value)),
440                Err(err) => tracing::warn!(
441                    collection,
442                    uri = %record.uri,
443                    error = %err,
444                    "skipping unparseable record in collection"
445                ),
446            }
447        }
448        Ok(out)
449    }
450
451    // ── subscriptions ────────────────────────────────────────────────────────
452
453    pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
454        self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
455    }
456
457    /// Every subscription, in the reader's deterministic order.
458    pub async fn list_subscriptions_sorted(
459        &self,
460    ) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
461        let mut subs = self.list_subscriptions().await?;
462        subs.sort_by(crate::lexicon::sort::subscriptions);
463        Ok(subs)
464    }
465
466    /// Subscribe to a feed. Returns the new record's rkey so the caller can
467    /// address it (rename, delete) without re-listing.
468    pub async fn add_subscription(
469        &self,
470        sub: &crate::vetted::VettedSubscription,
471    ) -> Result<String> {
472        Ok(self
473            .create_record(crate::lexicon::nsid::SUBSCRIPTION, sub)
474            .await?
475            .into_rkey())
476    }
477
478    pub async fn remove_subscription(&self, rkey: &str) -> Result<()> {
479        self.delete_record(crate::lexicon::nsid::SUBSCRIPTION, rkey)
480            .await
481    }
482
483    /// Replace a subscription in place — retitle, refile, change cadence.
484    /// Returns the [`WriteResult`] rather than discarding it: the caller logs
485    /// the resulting record URI, and a `()` here would have silently dropped
486    /// that field from the log line after the cutover.
487    pub async fn update_subscription(
488        &self,
489        rkey: &str,
490        sub: &crate::vetted::VettedSubscription,
491    ) -> Result<WriteResult> {
492        self.put_record(crate::lexicon::nsid::SUBSCRIPTION, rkey, sub)
493            .await
494    }
495
496    /// Batch-add many subscriptions via `applyWrites` (chunked) — the OPML-import path.
497    ///
498    /// One round trip per 200 feeds rather than one per feed: an import of
499    /// several hundred feeds is the case this exists for. More than one call
500    /// means the import can part-land; on an error,
501    /// [`crate::atproto::ApplyWritesIncomplete::of`] gives `landed`, and the
502    /// first `landed` of `subs` are in the repo.
503    /// Returns the new rkeys, which are assigned HERE rather than by the server:
504    /// client-side TIDs keep the imported feeds in input order and make the
505    /// batch reproducible. The sidecar client does the same, with the same
506    /// generator.
507    pub async fn add_subscriptions_bulk(
508        &self,
509        subs: &[crate::vetted::VettedSubscription],
510    ) -> Result<Vec<String>> {
511        let mut gen = crate::atproto::TidGenerator::new();
512        let mut rkeys = Vec::with_capacity(subs.len());
513        let mut writes = Vec::with_capacity(subs.len());
514        for sub in subs {
515            let rkey = gen.next();
516            writes.push(WriteOp::Create {
517                collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
518                rkey: Some(rkey.clone()),
519                // Propagated, NOT defaulted: `unwrap_or(Value::Null)` here would
520                // write a null record into the user's repo on a serialization
521                // failure rather than failing the import.
522                value: serde_json::to_value(sub)?,
523            });
524            rkeys.push(rkey);
525        }
526        self.apply_writes(&writes).await?;
527        Ok(rkeys)
528    }
529
530    // ── folders ──────────────────────────────────────────────────────────────
531
532    pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
533        self.list_typed(crate::lexicon::nsid::FOLDER).await
534    }
535
536    pub async fn list_folders_sorted(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
537        let mut folders = self.list_folders().await?;
538        folders.sort_by(crate::lexicon::sort::folders);
539        Ok(folders)
540    }
541
542    pub async fn add_folder(&self, folder: &crate::lexicon::Folder) -> Result<String> {
543        Ok(self
544            .create_record(crate::lexicon::nsid::FOLDER, folder)
545            .await?
546            .into_rkey())
547    }
548
549    /// Delete a folder. Subscriptions referencing it are left alone; a dangling
550    /// reference reads as "unfiled", which is the same behaviour the sidecar
551    /// client has.
552    pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
553        self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
554    }
555
556    /// Returns the [`WriteResult`], matching the sidecar client — the caller
557    /// logs the record URI from it.
558    pub async fn rename_folder(
559        &self,
560        rkey: &str,
561        folder: &crate::lexicon::Folder,
562    ) -> Result<WriteResult> {
563        self.put_record(crate::lexicon::nsid::FOLDER, rkey, folder)
564            .await
565    }
566
567    // ── saved ────────────────────────────────────────────────────────────────
568
569    pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
570        self.list_typed(crate::lexicon::nsid::SAVED).await
571    }
572
573    /// Saved entries, newest first.
574    pub async fn list_saved_sorted(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
575        let mut saved = self.list_saved().await?;
576        saved.sort_by(crate::lexicon::sort::saved);
577        Ok(saved)
578    }
579
580    pub async fn add_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<String> {
581        Ok(self
582            .create_record(crate::lexicon::nsid::SAVED, saved)
583            .await?
584            .into_rkey())
585    }
586
587    pub async fn remove_saved(&self, rkey: &str) -> Result<()> {
588        self.delete_record(crate::lexicon::nsid::SAVED, rkey).await
589    }
590
591    // ── read state ───────────────────────────────────────────────────────────
592
593    pub async fn list_read_states(&self) -> Result<Vec<(String, crate::lexicon::ReadState)>> {
594        self.list_typed(crate::lexicon::nsid::READ_STATE).await
595    }
596
597    /// Upsert one read cursor at its feed-derived rkey.
598    pub async fn put_read_state(
599        &self,
600        rkey: &str,
601        state: &crate::lexicon::ReadState,
602    ) -> Result<()> {
603        self.put_record(crate::lexicon::nsid::READ_STATE, rkey, state)
604            .await?;
605        Ok(())
606    }
607
608    /// Flush many dirty read cursors via `applyWrites` (chunked).
609    ///
610    /// Read state changes on nearly every page view, so this is the hottest
611    /// write path in the app; one round trip per flush rather than per feed is
612    /// the whole point. A flush past the op or byte bound is several calls,
613    /// and a failure part-way leaves the first
614    /// [`landed`](crate::atproto::ApplyWritesIncomplete::landed) cursors
615    /// written — see `crate::atproto::apply_writes_chunked`.
616    /// The `bool` is whether the record already exists in the PDS. It is not
617    /// optional bookkeeping: an `#update` on a missing record ERRORS, and
618    /// `applyWrites` is atomic per repo, so one not-yet-created cursor in the
619    /// batch would drop the whole DID's flush. The op builder is shared with the
620    /// sidecar client so the two cannot decide create-vs-update differently.
621    pub async fn flush_read_states(
622        &self,
623        cursors: &[(String, crate::lexicon::ReadState, bool)],
624    ) -> Result<()> {
625        if cursors.is_empty() {
626            return Ok(());
627        }
628        let writes = crate::atproto::read_state_write_ops(cursors)?;
629        self.apply_writes(&writes).await
630    }
631}
632
633#[cfg(test)]
634mod tests {
635    use super::*;
636
637    /// **Both clients must build the SAME write ops.**
638    ///
639    /// `flush_read_states` is the hottest write path in the app, and the
640    /// create-vs-update choice is the one part of it that cannot be got wrong
641    /// quietly: an `#update` on a record that does not exist errors, and
642    /// `applyWrites` is atomic per repo, so a single first-flush cursor in the
643    /// batch takes the whole DID's flush down with it.
644    ///
645    /// The first version of this method here ignored the flag and emitted
646    /// `#create` unconditionally, which would have broken every feed's first
647    /// flush after the cutover. Sharing the builder is what makes that
648    /// impossible rather than merely fixed.
649    #[test]
650    fn read_state_writes_choose_create_or_update_per_cursor() {
651        let state = crate::lexicon::ReadState::new(
652            "https://example.com/feed",
653            Some("2026-01-01T00:00:00Z".to_string()),
654            "2026-01-01T00:00:00Z",
655        );
656        let cursors = vec![
657            ("existing".to_string(), state.clone(), true),
658            ("brand-new".to_string(), state.clone(), false),
659        ];
660
661        let ops = crate::atproto::read_state_write_ops(&cursors).expect("ops build");
662        assert_eq!(ops.len(), 2);
663
664        let rendered: Vec<Value> = ops.iter().map(|op| op.to_json()).collect();
665        assert_eq!(
666            rendered[0]["$type"], "com.atproto.repo.applyWrites#update",
667            "an existing record must be UPDATED, not re-created"
668        );
669        assert_eq!(
670            rendered[1]["$type"], "com.atproto.repo.applyWrites#create",
671            "a first flush must CREATE, or the whole atomic batch fails"
672        );
673    }
674
675    /// **Bulk import, through the real `Repo`, asserted on the bytes.** The
676    /// test this replaces exercised `TidGenerator` directly and never called
677    /// `add_subscriptions_bulk`; the function could stop assigning rkeys and
678    /// return garbage with the suite green.
679    #[tokio::test]
680    async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
681        let (base, log) = crate::net::tests::serve_json_capturing(b"{}".to_vec()).await;
682        let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
683        crate::net::test_host_override(
684            "bulk-pds.test",
685            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
686        );
687        let http = Client::new();
688        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
689        crate::store::init_schema(&pool).await.unwrap();
690        let key = SigningKey::generate("k");
691        let mut s = session();
692        s.aud = format!("http://bulk-pds.test:{port}");
693        let repo = repo(&http, &pool, &s, &key);
694        let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
695            .map(|i| {
696                crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
697                    format!("https://f{i}.example/feed.xml"),
698                    "2026-07-12T00:00:00.000Z",
699                ))
700            })
701            .collect();
702
703        let rkeys = repo
704            .add_subscriptions_bulk(&subs)
705            .await
706            .expect("bulk write failed");
707
708        let sent = log.lock().unwrap().clone();
709        assert_eq!(
710            sent.len(),
711            1,
712            "expected one applyWrites request, got {sent:?}"
713        );
714        let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
715            .expect("request body is JSON");
716        let writes = body["writes"].as_array().expect("writes array");
717        assert_eq!(writes.len(), 3);
718        for (i, w) in writes.iter().enumerate() {
719            assert_eq!(w["collection"], crate::lexicon::nsid::SUBSCRIPTION);
720            assert_eq!(w["rkey"].as_str(), Some(rkeys[i].as_str()));
721        }
722        let mut sorted = rkeys.clone();
723        sorted.sort();
724        assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
725    }
726
727    fn session() -> OAuthSession {
728        OAuthSession {
729            sub: "did:plc:ewvi7nxzyoun6zhxrhs64oiz".into(),
730            issuer: "https://pds.example.com".into(),
731            aud: "https://pds.example.com".into(),
732            dpop_key_jwk: "{}".into(),
733            access_token: "tok".into(),
734            refresh_token: "ref".into(),
735            token_type: "DPoP".into(),
736            granted_scope: "atproto".into(),
737            expires_at: None,
738        }
739    }
740
741    fn repo<'a>(
742        http: &'a Client,
743        pool: &'a SqlitePool,
744        session: &'a OAuthSession,
745        key: &'a SigningKey,
746    ) -> Repo<'a> {
747        Repo {
748            http,
749            pool,
750            session,
751            key,
752        }
753    }
754
755    /// **The PDS comes from the session**, so a call cannot be pointed at
756    /// another host by a caller who passes the wrong base.
757    #[tokio::test]
758    async fn endpoints_are_built_from_the_sessions_audience() {
759        let http = Client::new();
760        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
761        let key = SigningKey::generate("k");
762        let mut s = session();
763        s.aud = "https://pds.example.com/".into();
764        let repo = repo(&http, &pool, &s, &key);
765        assert_eq!(
766            repo.url("com.atproto.repo.listRecords"),
767            "https://pds.example.com/xrpc/com.atproto.repo.listRecords",
768            "a trailing slash on the audience must not double the separator"
769        );
770    }
771
772    /// An XRPC error document is summarised by its `error`/`message` fields
773    /// only — never by echoing the raw body, which on other paths holds tokens.
774    #[test]
775    fn an_xrpc_error_is_summarised_not_echoed() {
776        let body = br#"{"error":"InvalidRequest","message":"unknown collection"}"#;
777        let rendered = xrpc_error(body, 400);
778        assert!(rendered.contains("InvalidRequest"));
779        assert!(rendered.contains("unknown collection"));
780
781        // A body that is not an XRPC error contributes nothing but the status.
782        let opaque = xrpc_error(br#"{"access_token":"SECRET"}"#, 500);
783        assert_eq!(opaque, "status 500");
784        assert!(!opaque.contains("SECRET"));
785        assert_eq!(xrpc_error(b"<html>oops</html>", 502), "status 502");
786    }
787
788    /// Every repo call must fail closed on an internal target, like every other
789    /// outbound path — asserted on the guard's own error.
790    #[tokio::test]
791    async fn repo_calls_fail_closed_on_an_internal_pds() {
792        let http = Client::new();
793        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
794        super::super::store::init_schema(&pool).await.unwrap();
795        let key = SigningKey::generate("k");
796        let mut s = session();
797        s.aud = "http://127.0.0.1:2583".into();
798        let repo = repo(&http, &pool, &s, &key);
799
800        let err = repo
801            .list_records("app.feather.subscription", None, None)
802            .await
803            .expect_err("must refuse a loopback PDS");
804        assert!(
805            format!("{err:#}").contains("forbidden (internal) address"),
806            "failed for the wrong reason: {err:#}"
807        );
808    }
809
810    /// **A 200 carrying an error envelope must not read as an empty repo.**
811    ///
812    /// `records` was taken off the JSON with `unwrap_or(Array([]))`, so a PDS
813    /// answering `200 {"error": …}` produced `Ok(vec![])`. That is not the
814    /// fail-closed branch in `web::resolve_subscriptions`: `sync_sub_refs`
815    /// writes the empty set through and `replace_sub_refs` DELETEs the DID's
816    /// entire `sub_ref` projection — one bad response revokes the reader's
817    /// access to every feed they have. Driven through the real client against
818    /// a real server, because the bug was the missing CALL, not the check.
819    #[tokio::test]
820    async fn a_200_error_envelope_is_not_an_empty_repo() {
821        let base = crate::net::tests::serve_body(
822            br#"{"error":"InvalidRequest","message":"bad cursor"}"#.to_vec(),
823        )
824        .await;
825        let port: u16 = base
826            .trim_end_matches('/')
827            .rsplit(':')
828            .next()
829            .unwrap()
830            .parse()
831            .unwrap();
832        crate::net::test_host_override(
833            "envelope-pds.test",
834            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
835        );
836
837        let http = Client::new();
838        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
839        crate::store::init_schema(&pool).await.unwrap();
840        let key = SigningKey::generate("k");
841        let mut s = session();
842        s.aud = format!("http://envelope-pds.test:{port}");
843        let repo = repo(&http, &pool, &s, &key);
844
845        let err = repo
846            .list_records("app.feather.subscription", None, None)
847            .await
848            .expect_err("an error envelope was read as an empty page");
849        assert!(
850            format!("{err:#}").contains("InvalidRequest"),
851            "failed for the wrong reason: {err:#}"
852        );
853    }
854
855    /// **The live walk spends its budget across pages too.**
856    ///
857    /// This is the `backend=rust` walk whose result reaches `replace_sub_refs`,
858    /// so a bound that silently failed to accumulate here would revoke a
859    /// reader's access to every feed past the cut. Three pages, a two-page
860    /// budget: the walk must refuse, and must say it kept two.
861    #[tokio::test]
862    async fn the_live_walk_spends_its_budget_across_pages() {
863        let (bodies, per_page) = crate::atproto::tests::paged_bodies(3, 4096, false);
864        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
865        let port: u16 = base
866            .trim_end_matches('/')
867            .rsplit(':')
868            .next()
869            .unwrap()
870            .parse()
871            .unwrap();
872        crate::net::test_host_override(
873            "live-budget-pages.test",
874            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
875        );
876
877        let http = Client::new();
878        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
879        crate::store::init_schema(&pool).await.unwrap();
880        let key = SigningKey::generate("k");
881        let mut s = session();
882        s.aud = format!("http://live-budget-pages.test:{port}");
883        let repo = repo(&http, &pool, &s, &key);
884
885        let err = repo
886            .list_all_records_within(
887                "app.feather.subscription",
888                &mut crate::atproto::ByteBudget::new(per_page * 2),
889            )
890            .await
891            .expect_err("three pages cannot fit in a two-page budget");
892        let msg = format!("{err:#}");
893        assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
894        assert!(
895            msg.contains("2 held"),
896            "the live walk did not accumulate across pages: {msg}"
897        );
898    }
899
900    /// **The live walk refuses a short list.** Its result reaches
901    /// `replace_sub_refs`, so returning a truncated list as a complete one
902    /// deletes every subscription past the page budget.
903    #[tokio::test]
904    async fn the_live_walk_that_runs_out_of_pages_refuses() {
905        let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
906            .map(|i| {
907                serde_json::json!({
908                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
909                    "cursor": format!("p{}", i + 1),
910                })
911                .to_string()
912                .into_bytes()
913            })
914            .collect();
915        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
916        let port: u16 = base
917            .trim_end_matches('/')
918            .rsplit(':')
919            .next()
920            .unwrap()
921            .parse()
922            .unwrap();
923        crate::net::test_host_override(
924            "pages-exhausted-live.test",
925            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
926        );
927        let http = Client::new();
928        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
929        crate::store::init_schema(&pool).await.unwrap();
930        let key = SigningKey::generate("k");
931        let mut s = session();
932        s.aud = format!("http://pages-exhausted-live.test:{port}");
933        let repo = repo(&http, &pool, &s, &key);
934
935        let err = repo
936            .list_all_records("app.feather.subscription")
937            .await
938            .expect_err("a truncated list was returned as a complete one");
939        assert!(
940            format!("{err:#}").contains("did not finish"),
941            "failed for the wrong reason: {err:#}"
942        );
943    }
944
945    /// #177 on the live backend: a malformed record in the reader's own repo
946    /// refuses the walk, by type, rather than dropping that subscription.
947    #[tokio::test]
948    async fn the_live_walk_refuses_a_page_with_a_malformed_record() {
949        let body = serde_json::json!({ "records": [
950            { "uri": "at://did:plc:x/c/3labGOOD", "value": {} },
951            { "cid": "bafy", "value": {} },
952        ]})
953        .to_string()
954        .into_bytes();
955        let base = crate::net::tests::serve_bodies_in_sequence(vec![body]).await;
956        let port: u16 = base
957            .trim_end_matches('/')
958            .rsplit(':')
959            .next()
960            .unwrap()
961            .parse()
962            .unwrap();
963        crate::net::test_host_override(
964            "malformed-live.test",
965            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
966        );
967        let http = Client::new();
968        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
969        crate::store::init_schema(&pool).await.unwrap();
970        let key = SigningKey::generate("k");
971        let mut s = session();
972        s.aud = format!("http://malformed-live.test:{port}");
973        let repo = repo(&http, &pool, &s, &key);
974        let err = repo
975            .list_all_records("app.feather.subscription")
976            .await
977            .expect_err("a page with a malformed record was accepted");
978        assert!(
979            err.downcast_ref::<crate::atproto::MalformedRecords>()
980                .is_some(),
981            "refused for the wrong reason: {err:#}"
982        );
983    }
984
985    /// The live walk's clean finish must still be a success — see the sidecar's
986    /// twin for why this direction is the dangerous one.
987    #[tokio::test]
988    async fn the_live_walk_that_finishes_cleanly_returns_the_records() {
989        let mut bodies: Vec<Vec<u8>> = (0..3)
990            .map(|i| {
991                serde_json::json!({
992                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
993                    "cursor": format!("p{}", i + 1),
994                })
995                .to_string()
996                .into_bytes()
997            })
998            .collect();
999        // **The terminator CARRIES a record**, because a real PDS ends on a
1000        // partial page and those are the records an off-by-one loses. Ending on
1001        // an empty page kept "drop the last page's records" alive: the count
1002        // below was right either way.
1003        bodies.push(
1004            serde_json::json!({
1005                "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1006            })
1007            .to_string()
1008            .into_bytes(),
1009        );
1010        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1011        let port: u16 = base
1012            .trim_end_matches('/')
1013            .rsplit(':')
1014            .next()
1015            .unwrap()
1016            .parse()
1017            .unwrap();
1018        crate::net::test_host_override(
1019            "clean-finish-live.test",
1020            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1021        );
1022        let http = Client::new();
1023        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1024        crate::store::init_schema(&pool).await.unwrap();
1025        let key = SigningKey::generate("k");
1026        let mut s = session();
1027        s.aud = format!("http://clean-finish-live.test:{port}");
1028        let repo = repo(&http, &pool, &s, &key);
1029
1030        let records = repo
1031            .list_all_records("app.feather.subscription")
1032            .await
1033            .expect("a walk that ran out of records is not a short list");
1034        assert_eq!(records.len(), 4);
1035        assert!(
1036            records.iter().any(|r| r.uri.ends_with("3labLAST")),
1037            "the LAST page's records were dropped: {:?}",
1038            records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
1039        );
1040    }
1041
1042    /// **The page cap is pinned exactly, not to within one.**
1043    ///
1044    /// `the_live_walk_that_runs_out_of_pages_refuses` serves `MAX_LIST_PAGES + 1`
1045    /// pages, so a budget one page SHORT refuses too and that mutation survives
1046    /// it. A walk whose last allowed request is the terminating one must come
1047    /// back `Ok` — and on this backend an `Err` is `replace_sub_refs` never
1048    /// running, which is the direction that costs a reader their subscriptions.
1049    #[tokio::test]
1050    async fn a_live_walk_that_terminates_on_its_last_allowed_page_succeeds() {
1051        let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
1052            .map(|i| {
1053                serde_json::json!({
1054                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
1055                    "cursor": format!("p{}", i + 1),
1056                })
1057                .to_string()
1058                .into_bytes()
1059            })
1060            .collect();
1061        bodies.push(
1062            serde_json::json!({
1063                "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
1064            })
1065            .to_string()
1066            .into_bytes(),
1067        );
1068        assert_eq!(bodies.len(), MAX_LIST_PAGES);
1069        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
1070        let port: u16 = base
1071            .trim_end_matches('/')
1072            .rsplit(':')
1073            .next()
1074            .unwrap()
1075            .parse()
1076            .unwrap();
1077        crate::net::test_host_override(
1078            "last-allowed-page-live.test",
1079            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1080        );
1081        let http = Client::new();
1082        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1083        crate::store::init_schema(&pool).await.unwrap();
1084        let key = SigningKey::generate("k");
1085        let mut s = session();
1086        s.aud = format!("http://last-allowed-page-live.test:{port}");
1087        let repo = repo(&http, &pool, &s, &key);
1088
1089        let records = repo
1090            .list_all_records("app.feather.subscription")
1091            .await
1092            .expect("a walk that terminated inside its budget is not a short list");
1093        assert_eq!(
1094            records.len(),
1095            MAX_LIST_PAGES,
1096            "a walk that used its whole page budget and finished lost records",
1097        );
1098    }
1099
1100    /// **The ERROR twin of `send`, which a review found unguarded.**
1101    ///
1102    /// `send`'s success branch goes through `PostOutcome::json` and its cap; its
1103    /// failure branch goes to `xrpc_error` → `error_fields`, which deserialised the
1104    /// whole body. So a hostile PDS answering **500** instead of 200 with the same
1105    /// explosion got the full amplification on every write and on every listing
1106    /// failure — the guard bypassed by a status code.
1107    ///
1108    /// Both directions: a small error body still yields its reason, because that
1109    /// string is what a reader's log needs to tell "your PDS said no" from "we
1110    /// broke".
1111    #[test]
1112    fn an_oversized_error_body_is_not_parsed_for_its_reason() {
1113        let small = br#"{"error":"InvalidSwap","message":"record changed"}"#;
1114        assert_eq!(
1115            Repo::error_fields(small).as_deref(),
1116            Some("InvalidSwap: record changed"),
1117            "a real error body must still render its reason",
1118        );
1119
1120        let mut huge = String::from(r#"{"error":"InvalidSwap","pad":["#);
1121        while huge.len() < crate::oauth::MAX_ERROR_BODY + 1_024 {
1122            huge.push_str("{},");
1123        }
1124        huge.push_str("{}]}");
1125        assert!(huge.len() > crate::oauth::MAX_ERROR_BODY);
1126        assert_eq!(
1127            Repo::error_fields(huge.as_bytes()),
1128            None,
1129            "an oversized error body was deserialised to fish out one string",
1130        );
1131    }
1132
1133    /// **A PDS rejection is structured, not only a sentence (#241).**
1134    ///
1135    /// The read-state flusher has to tell "the PDS refused this batch" from "the
1136    /// network broke", and the sidecar client already said so with
1137    /// [`crate::atproto::AtProtoError::Xrpc`]. This client only said it in a
1138    /// string, so the one backend production runs could be matched on nothing
1139    /// sturdier than its wording. The rendered message must not change: it is
1140    /// what every existing log line and assertion reads.
1141    #[tokio::test]
1142    async fn a_rejected_write_carries_the_status_and_error_name() {
1143        use axum::response::IntoResponse as _;
1144        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1145        let addr = listener.local_addr().unwrap();
1146        let host = format!("rejecting-{}.xrpc.test", addr.port());
1147        crate::net::test_host_override(&host, addr);
1148        let app = axum::Router::new().fallback(|| async {
1149            (
1150                axum::http::StatusCode::INTERNAL_SERVER_ERROR,
1151                axum::Json(
1152                    json!({ "error": "InternalServerError", "message": "Internal Server Error" }),
1153                ),
1154            )
1155                .into_response()
1156        });
1157        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1158
1159        let http = Client::new();
1160        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1161        crate::store::init_schema(&pool).await.unwrap();
1162        let key = SigningKey::generate("k");
1163        let mut s = session();
1164        s.aud = format!("http://{host}:{}", addr.port());
1165        let repo = repo(&http, &pool, &s, &key);
1166
1167        let err = repo
1168            .apply_writes(&[WriteOp::Delete {
1169                collection: crate::lexicon::nsid::READ_STATE.into(),
1170                rkey: "rs-0".into(),
1171            }])
1172            .await
1173            .expect_err("a 500 is a failure");
1174        assert_eq!(
1175            err.to_string(),
1176            "com.atproto.repo.applyWrites failed: status 500 \
1177             (InternalServerError: Internal Server Error)",
1178            "the rendered message changed",
1179        );
1180        let xrpc = err
1181            .chain()
1182            .find_map(|cause| cause.downcast_ref::<crate::atproto::AtProtoError>());
1183        match xrpc {
1184            Some(crate::atproto::AtProtoError::Xrpc { status, error, .. }) => {
1185                assert_eq!(status.as_u16(), 500);
1186                assert_eq!(error, "InternalServerError");
1187            }
1188            other => panic!("no structured XRPC error in the chain: {other:?}"),
1189        }
1190    }
1191
1192    /// **The write path parses a PDS body too, and it had no bound before the
1193    /// parse.**
1194    ///
1195    /// `listRecords` is guarded by `parse_list_records`; every OTHER body this
1196    /// client turns into a `Value` goes through `PostOutcome::json` — the repo
1197    /// writers' responses, the PAR response, the token response, the session
1198    /// refresh. `read_capped` bounds the WIRE at 8 MB, which is the *input* to the
1199    /// amplification rather than a limit on it: 8 MB of the cheapest node shape
1200    /// measured 824 MB retained on a 512 MB box.
1201    ///
1202    /// Drives `delete_record`, which is the shortest route from a handler to
1203    /// `send` → `json`.
1204    #[tokio::test]
1205    async fn the_live_write_path_refuses_a_node_explosion() {
1206        let mut body = String::from(r#"{"uri":"at://d/c/r","value":["#);
1207        for _ in 0..1_200_000 {
1208            body.push_str("{},");
1209        }
1210        body.push_str("{}]}");
1211        assert!(
1212            crate::atproto::count_structural_chars(body.as_bytes())
1213                > crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
1214            "the probe body is not over the cap, so this test proves nothing",
1215        );
1216        let base = crate::net::tests::serve_body(body.into_bytes()).await;
1217        let port: u16 = base
1218            .trim_end_matches('/')
1219            .rsplit(':')
1220            .next()
1221            .unwrap()
1222            .parse()
1223            .unwrap();
1224        crate::net::test_host_override(
1225            "write-explosion.test",
1226            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1227        );
1228        let http = Client::new();
1229        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1230        crate::store::init_schema(&pool).await.unwrap();
1231        let key = SigningKey::generate("k");
1232        let mut s = session();
1233        s.aud = format!("http://write-explosion.test:{port}");
1234        let repo = repo(&http, &pool, &s, &key);
1235
1236        let err = repo
1237            .delete_record("c", "r")
1238            .await
1239            .expect_err("a node explosion on the write path was parsed rather than refused");
1240        assert!(
1241            format!("{err:#}").contains("structural characters"),
1242            "failed for the wrong reason: {err:#}"
1243        );
1244    }
1245
1246    /// The other direction: an ordinary write response still parses. Without this
1247    /// a guard that refused every body would pass the test above.
1248    #[tokio::test]
1249    async fn the_live_write_path_accepts_an_ordinary_response() {
1250        let base =
1251            crate::net::tests::serve_body(br#"{"commit":{"cid":"bafy","rev":"3lab"}}"#.to_vec())
1252                .await;
1253        let port: u16 = base
1254            .trim_end_matches('/')
1255            .rsplit(':')
1256            .next()
1257            .unwrap()
1258            .parse()
1259            .unwrap();
1260        crate::net::test_host_override(
1261            "write-ordinary.test",
1262            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1263        );
1264        let http = Client::new();
1265        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1266        crate::store::init_schema(&pool).await.unwrap();
1267        let key = SigningKey::generate("k");
1268        let mut s = session();
1269        s.aud = format!("http://write-ordinary.test:{port}");
1270        let repo = repo(&http, &pool, &s, &key);
1271
1272        repo.delete_record("c", "r")
1273            .await
1274            .expect("an ordinary write response was refused");
1275    }
1276
1277    /// **A duplicated `records` key must not be able to empty a page.**
1278    ///
1279    /// This is the live `backend=rust` walk, so an empty page here reaches
1280    /// `replace_sub_refs` and deletes the reader's subscriptions. serde refuses
1281    /// a repeated field outright; a `serde_json::Value` takes the last one
1282    /// silently, so routing this body through a `Value` first turns a smuggled
1283    /// second key into a successful, empty listing. The test exists to pin
1284    /// which of the two this client uses.
1285    #[tokio::test]
1286    async fn the_live_walk_refuses_a_duplicated_records_key() {
1287        let base = crate::net::tests::serve_body(
1288            br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#.to_vec(),
1289        )
1290        .await;
1291        let port: u16 = base
1292            .trim_end_matches('/')
1293            .rsplit(':')
1294            .next()
1295            .unwrap()
1296            .parse()
1297            .unwrap();
1298        crate::net::test_host_override(
1299            "dup-records.test",
1300            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1301        );
1302
1303        let http = Client::new();
1304        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1305        crate::store::init_schema(&pool).await.unwrap();
1306        let key = SigningKey::generate("k");
1307        let mut s = session();
1308        s.aud = format!("http://dup-records.test:{port}");
1309        let repo = repo(&http, &pool, &s, &key);
1310
1311        let err = repo
1312            .list_records("app.feather.subscription", None, None)
1313            .await
1314            .expect_err("a duplicated records key was read as an empty page");
1315        assert!(
1316            format!("{err:#}").contains("duplicate"),
1317            "failed for the wrong reason: {err:#}"
1318        );
1319    }
1320
1321    /// **An empty 2xx body is not an empty repo either.** `send` maps a
1322    /// zero-length 2xx to `Value::Null` — deliberately, for `deleteRecord` and
1323    /// `applyWrites` — and while `listRecords` still went through it, that
1324    /// slipped past the envelope guard
1325    /// and `records.unwrap_or(Array([]))` produced `Ok(vec![])`: the same
1326    /// `sub_ref` wipe the guard was added to prevent, through the sibling door.
1327    #[tokio::test]
1328    async fn an_empty_200_body_is_not_an_empty_repo() {
1329        let base = crate::net::tests::serve_body(Vec::new()).await;
1330        let port: u16 = base
1331            .trim_end_matches('/')
1332            .rsplit(':')
1333            .next()
1334            .unwrap()
1335            .parse()
1336            .unwrap();
1337        crate::net::test_host_override(
1338            "empty-body.test",
1339            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1340        );
1341        let http = Client::new();
1342        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1343        crate::store::init_schema(&pool).await.unwrap();
1344        let key = SigningKey::generate("k");
1345        let mut s = session();
1346        s.aud = format!("http://empty-body.test:{port}");
1347        let repo = repo(&http, &pool, &s, &key);
1348
1349        let err = repo
1350            .list_records("app.feather.subscription", None, None)
1351            .await
1352            .expect_err("an empty body was read as an empty repo");
1353        assert!(
1354            format!("{err:#}").contains("no records"),
1355            "failed for the wrong reason: {err:#}"
1356        );
1357    }
1358
1359    /// The write paths discarded the body too: `delete_record` and
1360    /// `apply_writes` are `send(..).await?; Ok(())`, so a 200 carrying an
1361    /// error envelope reported success for a delete that did not happen.
1362    #[tokio::test]
1363    async fn a_200_error_envelope_is_not_a_successful_write() {
1364        let base = crate::net::tests::serve_body(
1365            br#"{"error":"InvalidRequest","message":"nope"}"#.to_vec(),
1366        )
1367        .await;
1368        let port: u16 = base
1369            .trim_end_matches('/')
1370            .rsplit(':')
1371            .next()
1372            .unwrap()
1373            .parse()
1374            .unwrap();
1375        crate::net::test_host_override(
1376            "envelope-write.test",
1377            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1378        );
1379        let http = Client::new();
1380        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1381        crate::store::init_schema(&pool).await.unwrap();
1382        let key = SigningKey::generate("k");
1383        let mut s = session();
1384        s.aud = format!("http://envelope-write.test:{port}");
1385        let repo = repo(&http, &pool, &s, &key);
1386
1387        let err = repo
1388            .delete_record("app.feather.subscription", "rk1")
1389            .await
1390            .expect_err("a failed delete was reported as success");
1391        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1392
1393        let err = repo
1394            .apply_writes(&[crate::atproto::WriteOp::Delete {
1395                collection: "app.feather.subscription".to_string(),
1396                rkey: "rk1".to_string(),
1397            }])
1398            .await
1399            .expect_err("a failed batch was reported as success");
1400        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1401    }
1402
1403    /// An empty batch must not produce a request at all — an `applyWrites` with
1404    /// no writes is a round trip that can only fail.
1405    #[tokio::test]
1406    async fn an_empty_batch_is_not_sent() {
1407        let http = Client::new();
1408        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1409        let key = SigningKey::generate("k");
1410        let mut s = session();
1411        // A target that would fail loudly if it were ever contacted.
1412        s.aud = "http://127.0.0.1:2583".into();
1413        let repo = repo(&http, &pool, &s, &key);
1414        assert!(repo.apply_writes(&[]).await.is_ok());
1415    }
1416
1417    // ── applyWrites chunking (#240) ──────────────────────────────────────────
1418    //
1419    // This is the production backend, so the chunking is asserted here on the
1420    // bytes it sends to a fake PDS that refuses what the reference PDS refuses
1421    // (see `crate::atproto::tests::serve_apply_writes`).
1422
1423    /// A session pointed at the strict fake, reached through a per-port host
1424    /// override (the override table is process-wide and tests run in parallel).
1425    async fn strict_pds(
1426        fail_call: Option<usize>,
1427    ) -> (
1428        OAuthSession,
1429        SqlitePool,
1430        crate::atproto::tests::ApplyWritesLog,
1431    ) {
1432        let (base, log) = crate::atproto::tests::serve_apply_writes(fail_call).await;
1433        let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
1434        let host = format!("chunk-oauth-{port}.test");
1435        crate::net::test_host_override(&host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
1436        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1437        crate::store::init_schema(&pool).await.unwrap();
1438        let mut s = session();
1439        s.aud = format!("http://{host}:{port}");
1440        (s, pool, log)
1441    }
1442
1443    fn subs(n: usize) -> Vec<crate::vetted::VettedSubscription> {
1444        (0..n)
1445            .map(|i| {
1446                crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
1447                    format!("https://f{i}.example/feed.xml"),
1448                    "2026-07-12T00:00:00.000Z",
1449                ))
1450            })
1451            .collect()
1452    }
1453
1454    /// **An OPML import of 201 feeds is two calls, 200 then 1, in order** — on
1455    /// the backend production runs. One call of 201 is what the reference PDS
1456    /// refuses with `Too many writes. Max: 200`.
1457    #[tokio::test]
1458    async fn bulk_subscribe_of_201_is_two_calls_in_order() {
1459        let (s, pool, log) = strict_pds(None).await;
1460        let (http, key) = (Client::new(), SigningKey::generate("k"));
1461        let rkeys = repo(&http, &pool, &s, &key)
1462            .add_subscriptions_bulk(&subs(201))
1463            .await
1464            .expect("a 201-feed import must succeed against a PDS that caps at 200");
1465        assert_eq!(crate::atproto::tests::call_sizes(&log), vec![200, 1]);
1466        assert_eq!(
1467            crate::atproto::tests::sent_rkeys(&log),
1468            rkeys,
1469            "every feed, once, in input order"
1470        );
1471    }
1472
1473    /// 500 feeds — the default per-DID cap — are three calls; exactly 200 is one.
1474    #[tokio::test]
1475    async fn bulk_subscribe_splits_at_200_and_not_before() {
1476        for (n, want) in [(500, vec![200, 200, 100]), (200, vec![200])] {
1477            let (s, pool, log) = strict_pds(None).await;
1478            let (http, key) = (Client::new(), SigningKey::generate("k"));
1479            repo(&http, &pool, &s, &key)
1480                .add_subscriptions_bulk(&subs(n))
1481                .await
1482                .expect("bulk write");
1483            assert_eq!(crate::atproto::tests::call_sizes(&log), want, "{n} feeds");
1484        }
1485    }
1486
1487    /// **A failed chunk stops the run**, and chunk 3 is never sent.
1488    #[tokio::test]
1489    async fn bulk_subscribe_stops_at_the_first_failed_chunk() {
1490        let (s, pool, log) = strict_pds(Some(2)).await;
1491        let (http, key) = (Client::new(), SigningKey::generate("k"));
1492        let err = repo(&http, &pool, &s, &key)
1493            .add_subscriptions_bulk(&subs(500))
1494            .await
1495            .expect_err("a failed chunk must fail the call");
1496        assert_eq!(
1497            crate::atproto::tests::call_sizes(&log),
1498            vec![200, 200],
1499            "chunk 3 must NOT be sent"
1500        );
1501        assert!(format!("{err:#}").contains("boom"), "{err:#}");
1502    }
1503
1504    /// **The byte bound splits a read-state flush well under 200 ops.** Ten
1505    /// cursors at the lexicon's 1,000-id cap are ~500 KB — one call of that is
1506    /// refused by every reference PDS older than atproto#4989.
1507    #[tokio::test]
1508    async fn read_state_flush_splits_on_bytes_under_200_ops() {
1509        let (s, pool, log) = strict_pds(None).await;
1510        let (http, key) = (Client::new(), SigningKey::generate("k"));
1511        let cursors: Vec<(String, crate::lexicon::ReadState, bool)> = (0..10)
1512            .map(|i| {
1513                let mut state = crate::lexicon::ReadState::new(
1514                    format!("https://f{i}.example/feed.xml"),
1515                    None,
1516                    "2026-07-12T00:00:00.000Z",
1517                );
1518                state.read_ids = (0..crate::lexicon::ReadState::MAX_IDS)
1519                    .map(|j| format!("https://f{i}.example/posts/{j:04}/an-entry-permalink"))
1520                    .collect();
1521                (format!("rk{i:04}"), state, false)
1522            })
1523            .collect();
1524        repo(&http, &pool, &s, &key)
1525            .flush_read_states(&cursors)
1526            .await
1527            .expect("a byte-heavy flush must succeed in chunks");
1528        let sizes = crate::atproto::tests::call_sizes(&log);
1529        assert!(sizes.len() > 1, "one call for ~500 KB: {sizes:?}");
1530        let want: Vec<String> = cursors.iter().map(|(rkey, _, _)| rkey.clone()).collect();
1531        assert_eq!(crate::atproto::tests::sent_rkeys(&log), want);
1532    }
1533}