Skip to main content

feather_reader/
atproto.rs

1//! The atproto identity + PDS record layer.
2//!
3//! FeatherReader's defining bet is that a user's feed
4//! subscriptions, folders, saved items, and batched read-state live as records
5//! in the user's **own** atproto PDS under the open `community.lexicon.rss.*`
6//! community lexicon — not in the app's database. This module is the client that
7//! reads and writes those records.
8//!
9//! It has three layers:
10//!
11//! 1. **Identity resolution** ([`resolve_handle`], [`resolve_did_to_pds`]) —
12//!    turn a handle (`alice.example.com`) into a DID (`did:plc:…`), then resolve
13//!    the DID document to the PDS service endpoint. Handles resolve via the
14//!    account's PDS `com.atproto.identity.resolveHandle` (or the well-known
15//!    `/.well-known/atproto-did`); DIDs resolve via the PLC directory
16//!    (`did:plc:*`) or the `did:web` well-known document.
17//! 2. **A lightweight [`PdsClient`]** — holds the resolved DID, the PDS base URL,
18//!    and an [`Auth`] token, and exposes typed calls over `com.atproto.repo.*`:
19//!    [`list_records`](PdsClient::list_records),
20//!    `create_record`,
21//!    `put_record`,
22//!    [`delete_record`](PdsClient::delete_record), and
23//!    `apply_writes` (the **batch** call the
24//!    read-state flusher uses to coalesce many per-feed cursor writes into one
25//!    round-trip).
26//! 3. **Typed convenience wrappers** wired to the [`crate::lexicon`] record
27//!    types (list/create [`Subscription`]/[`Folder`]/[`Saved`], put
28//!    [`ReadState`], batch-flush many `ReadState` cursors).
29//!
30//! ## Auth — the OAuth sidecar is the live path
31//!
32//! Auth is a **trait/enum boundary** so the mechanism can vary without touching
33//! call sites. There are three paths:
34//!
35//! * **The live path — the atproto OAuth confidential client, via [`SidecarClient`].**
36//!   atproto OAuth (DPoP, PAR, token refresh) is fiddly and is **not** hand-rolled
37//!   in Rust: it runs in a small, supported `@atproto/oauth-client-node` sidecar.
38//!   The Rust server never holds
39//!   PDS tokens — it POSTs every `com.atproto.repo.*` op to the sidecar's
40//!   `/internal/repo` endpoint (gated by a shared `X-Internal-Secret`), and the
41//!   sidecar restores the DID's OAuth session (transparent DPoP + token refresh)
42//!   and runs the matching XRPC call. [`SidecarClient`] is that client; the typed
43//!   convenience wrappers (list/create/put/delete subscriptions, batch-flush
44//!   read-state) live on it and map 1:1 to the old [`PdsClient`] surface.
45//! * **The interim path — [`Auth::Session`] (app password).** A session obtained
46//!   from `com.atproto.server.createSession`. Kept behind the [`Auth`] seam, but
47//!   it is **no longer the live path**: [`PdsClient`] and
48//!   [`login_with_app_password`] remain for tests, while [`SidecarClient`] is
49//!   what the web layer routes through.
50//!
51//!   ⚠️ **This is no longer a working "local runs without the sidecar" fallback,
52//!   and the docs used to claim otherwise.** Since v0.2.8 every [`PdsClient`]
53//!   request goes through the SSRF guard, which refuses loopback, RFC1918, ULA
54//!   and `100.64/10` (Tailscale). So pointing this at `http://localhost:2583`
55//!   or a tailnet PDS now fails with *"refusing to fetch forbidden (internal)
56//!   address"* rather than returning a session. That is the guard behaving
57//!   correctly — the target host is attacker-influenced in the cases that
58//!   matter, and a dev-only escape hatch is exactly the kind of flag that ends
59//!   up set in production — but it does mean a local-PDS workflow needs the
60//!   PDS reachable on a public address, or a deliberate change here.
61//! * **The public-read path — [`Auth::Anonymous`], via [`PdsClient::anonymous`].**
62//!   `com.atproto.repo.listRecords` is public on a standard PDS, so a stranger's
63//!   `community.lexicon.rss.*` records can be read with no credentials at all.
64//!   An anonymous client sends no `Authorization` header and is **read-only** —
65//!   every write fails closed on [`Auth::bearer`]. Because the target host is
66//!   then chosen by a stranger, the read is routed through
67//!   [`crate::net::guarded_get_no_privacy`] (per-hop SSRF re-validation +
68//!   connect-pinning) and capped by [`crate::net::read_capped`].
69//!
70//! ## Every PDS request goes through the SSRF guard
71//!
72//! A PDS host is *never* a host FeatherReader chose: it comes out of a DID
73//! document, which is attacker-controllable. So identity resolution, the record
74//! **reads**, and the record **writes** all route through [`crate::net`] —
75//! [`crate::net::guarded_get_no_privacy`] and
76//! [`crate::net::guarded_post_json`] — rather than the shared
77//! `reqwest::Client`. [`resolve_did_to_pds`] runs
78//! [`crate::net::assert_public_target`] on the `serviceEndpoint` it returns, but
79//! that check is a *separate DNS resolution* from the later request; only
80//! re-vetting and connect-pinning at request time closes the rebinding window.
81//! The writes matter most: they carry the session bearer, and
82//! [`login_with_app_password`] carries the app password in the request **body**,
83//! where reqwest's cross-origin header sanitisation offers no protection at all
84//! — which is why the guarded POST refuses redirects outright.
85//!
86//! All network I/O is `reqwest` (rustls, no OpenSSL); every fallible path returns
87//! [`anyhow::Result`] or the typed [`AtProtoError`] — nothing panics.
88
89use std::sync::Arc;
90
91use anyhow::{Context, Result};
92use reqwest::header::{HeaderName, HeaderValue, AUTHORIZATION};
93use reqwest::{Client, StatusCode};
94use serde::de::DeserializeOwned;
95use serde::{Deserialize, Serialize};
96use serde_json::{json, Value};
97
98use crate::lexicon::{self, Folder, ReadState, Saved, Subscription};
99
100/// The public PLC directory, used to resolve `did:plc:*` DIDs to their DID
101/// document (and thus their PDS service endpoint).
102pub const DEFAULT_PLC_DIRECTORY: &str = "https://plc.directory";
103
104/// The default appview/entryway used only as a bootstrap host for handle
105/// resolution when the caller has no PDS hint yet. Handle resolution ultimately
106/// works against any atproto host that implements
107/// `com.atproto.identity.resolveHandle`; `bsky.social` is a reliable default.
108pub const DEFAULT_RESOLVER_HOST: &str = "https://bsky.social";
109
110/// Hard cap on cursor pages any `list_all_records` walk will follow.
111///
112/// [`crate::net::read_capped`] bounds each individual response, but nothing
113/// bounded the *accumulation* across pages: a repo host that returns a full page
114/// and a fresh cursor forever walks memory until the (512 MB) box dies. At 100
115/// records per page this admits 20 000 records — far past any real
116/// `community.lexicon.rss.*` collection — while making the loop finite against a
117/// host we do not control. Mirrors [`crate::network::MAX_PAGES`], which bounds
118/// the relay walk for the same reason.
119const MAX_LIST_PAGES: usize = 200;
120
121/// Hard cap on the records a single `list_all_records` walk will accumulate.
122///
123/// [`MAX_LIST_PAGES`] bounds how many REQUESTS a walk makes. It bounds the
124/// accumulated memory only if the server honours `limit=100` — and a repo host
125/// we did not choose has no obligation to. Measured: an 8 MB page (the
126/// [`crate::net::read_capped`] ceiling) holds ~95 000 minimal records and
127/// retains ~23 MB as `Vec<RecordEntry>`, so the page cap alone admits gigabytes
128/// on a 512 MB box.
129///
130/// 20 000 is the number [`MAX_LIST_PAGES`]'s own comment already claimed — this
131/// makes the claim true rather than conditional on the server's cooperation.
132const MAX_LIST_RECORDS: usize = 20_000;
133
134/// The same cap for a collection whose records are **large**.
135///
136/// [`MAX_LIST_RECORDS`]'s figure was measured against *minimal* records
137/// (~1 KB). A `site.standard.document` carries the whole article — ~17 KB
138/// measured across 449 real ones — so 20 000 of them is ~340 MB retained on a
139/// 512 MB box. Sized to the record, not to the protocol.
140pub(crate) const MAX_LARGE_RECORDS: usize = 2_000;
141
142/// The at-URI scheme prefix, **the one Rust spelling**. Every Rust guard that
143/// asks "is this an at-URI" strips or compares this.
144///
145/// **SQL no longer holds a second opinion.** There used to be a matching string
146/// predicate in `store`, and this comment claimed a test pinned the two in
147/// agreement. Both are gone: the predicate was deleted when `feeds.kind` became
148/// a cache of [`crate::feed::FeedKind::of`], re-derived from the URL rather than
149/// re-described in SQL, and no such test survived it. Nothing outside Rust
150/// decides what an at-URI is, so there is nothing left to keep in step.
151pub(crate) const AT_URI_PREFIX: &str = "at://";
152
153/// Strip the at-URI scheme **case-insensitively**, returning the body.
154///
155/// Schemes are case-insensitive per RFC 3986 and `Url::parse` folds them, so
156/// `At://` names the same thing as `at://`. Recognition has to match that, or a
157/// mixed-case row is an at-URI to the fetcher (which refuses it) and an
158/// ordinary URL to every guard — polled forever, failing forever. Whether such
159/// a spelling may be STORED is a separate question, answered no.
160pub(crate) fn strip_at_prefix(url: &str) -> Option<&str> {
161    url.get(..AT_URI_PREFIX.len())
162        .filter(|p| p.eq_ignore_ascii_case(AT_URI_PREFIX))
163        .map(|p| &url[p.len()..])
164}
165
166/// atproto's record-key rules, all of them: charset `[A-Za-z0-9._:~-]`, length
167/// 1..=512, and not `.` or `..`. The repo's TID tests state the same rule; this
168/// is the one place it is enforced on a key that arrives from outside.
169pub(crate) fn is_valid_rkey(rkey: &str) -> bool {
170    !rkey.is_empty()
171        && rkey.len() <= 512
172        && rkey != "."
173        && rkey != ".."
174        && rkey
175            .chars()
176            .all(|c| c.is_ascii_alphanumeric() || matches!(c, '.' | '_' | ':' | '~' | '-'))
177}
178
179/// Accumulate a page for a **reading** walk, keeping what fits and reporting
180/// whether anything was dropped.
181///
182/// **A truncation, never an error** — the opposite of [`extend_bounded`], and
183/// deliberately so. That function's refusal exists because its caller feeds
184/// `replace_sub_refs`, where a short list is revoked access. A walk that only
185/// ADDS entries has no such hazard, and refusing there is strictly worse: a
186/// publication with more documents than the cap would fail on every poll, so
187/// an ordinary long-running blog becomes permanently unreadable instead of
188/// partially read. The records kept are the ones the PDS returned first.
189pub(crate) fn extend_truncating(
190    out: &mut Vec<RecordEntry>,
191    page: Vec<RecordEntry>,
192    max: usize,
193) -> bool {
194    let room = max.saturating_sub(out.len());
195    // **Strictly greater.** `>=` called an exactly-full final page a
196    // truncation, so a collection holding exactly `max` records warned that it
197    // had dropped something on every poll.
198    let dropped = page.len() > room;
199    out.extend(page.into_iter().take(room));
200    dropped
201}
202
203/// The result of a bounded walk: what was read, and whether that is all of it.
204///
205/// **`complete` is a fact the caller cannot recover afterwards.** A short list
206/// from a truncating walk looks exactly like a short collection, and the
207/// difference is the one that matters: "this publication has nine articles" and
208/// "this reader gave up after nine" are the same `Vec` and very different
209/// answers.
210#[derive(Debug)]
211pub struct RecordWalk {
212    /// The records kept, in the order the PDS returned them.
213    pub records: Vec<RecordEntry>,
214    /// True when the collection ran out before any bound did.
215    pub complete: bool,
216    /// Records skipped because their envelope was malformed (#177). Only a
217    /// walk that SKIPS can report this; the walks that feed `replace_sub_refs`
218    /// refuse instead, with [`MalformedRecords`].
219    pub malformed: usize,
220}
221
222impl RecordWalk {
223    fn complete(records: Vec<RecordEntry>) -> Self {
224        Self {
225            records,
226            complete: true,
227            malformed: 0,
228        }
229    }
230    fn partial(records: Vec<RecordEntry>) -> Self {
231        Self {
232            records,
233            complete: false,
234            malformed: 0,
235        }
236    }
237}
238
239/// The memory one walk may retain.
240///
241/// **A record cap bounds memory only if you know what a record costs.** The
242/// caps above are counts, chosen against a measured ~17 KB document, and
243/// `MAX_LIST_PAGES` bounds requests rather than bytes. A PDS whose records are
244/// not that shape satisfies every count and still exhausts the box.
245///
246/// **128 MiB of ACCUMULATION per read — which is not the same as 128 MiB of
247/// memory, and an earlier version of this comment said it was.**
248///
249/// Every charge here is taken after `serde_json` has already built the page, so
250/// the true peak is this ceiling plus one page's tree, and a page's tree is not
251/// small: measured, an 8 MiB response of `{"":0}` objects retains 824 MB, a
252/// wire-to-heap amplification of 98x. A bound consulted after the allocation
253/// cannot prevent that allocation. What it does prevent is the accumulation
254/// across pages and across the walks of one read, which is the part that scales
255/// with how long a walk runs rather than with one response.
256///
257/// Closing the single-page case needs a smaller wire cap for `listRecords` or a
258/// parser that counts as it goes. Neither belongs to this bound; both are filed
259/// as #197 rather than implied here.
260///
261/// A caller passes one [`ByteBudget`] into every walk it makes, so a publication
262/// read — which runs a second walk while still holding the first's records — is
263/// bounded once rather than twice. Two independent ceilings put roughly 384 MB of
264/// accumulation in flight: 128 for the publications, 128 for the documents, and a
265/// further 128 of transient page because the per-page check compared against the
266/// ceiling instead of what was left.
267///
268/// **Only one caller threads it today**: the publication reader, polled by the
269/// scheduler's publication loop since 0.4.0. Every other live read
270/// builds its own ceiling per walk, so the per-request total is still a multiple
271/// of this number — two walks on an OPML export, four on a login — and nothing
272/// bounds concurrent requests at all.
273/// An earlier version of this constant was also 128 MiB while
274/// [`approx_bytes`] charged serialized length — 42x optimistic on hostile
275/// shapes, so the bound was nearly fiction. It was then cut to 64 MiB to
276/// compensate. Charging nodes removed the reason for the cut: the charge is now
277/// at or above what the page really retains, so 128 MiB of budget is at most
278/// 128 MiB of memory, which a 512 MB box carries.
279///
280/// The cut had a cost, measured rather than assumed. Per walk:
281///
282/// - The subscription walks cap at 20 000 records (5 000 on the live one). A
283///   real five-field subscription charges **2 188 bytes** here, and one carrying
284///   a folder and a fetch hint charges **2 764** — so a full repo is 42 to 53 MB,
285///   which was 65 to 82 % of a 64 MiB budget. Their verdict is a hard refusal
286///   that drops the reader into the fail-closed branch, so an account near the
287///   record cap with slightly longer titles would have served a stale projection
288///   on every poll, permanently. At 128 MiB that is 41 % and the count still
289///   binds first. An earlier version of this comment claimed 1.5 KB and 30 MB;
290///   that is a three-field record, not a real one.
291/// - The publication walk caps at [`MAX_LARGE_RECORDS`] (2 000). At the measured
292///   ~17 KB document that is about 37 MB either way. Above roughly 66 KB per
293///   article the budget binds first and the walk truncates early, reporting
294///   `complete: false` as it already does for the record cap.
295///
296/// So the counts bind first on everything measured, and a publication of
297/// extremely long articles truncates sooner than the count would. It remains
298/// untrue that nothing truncates that did not truncate before.
299pub(crate) const MAX_LIST_BYTES: usize = 128 * 1024 * 1024;
300
301/// What one record retains once parsed.
302///
303/// **Nodes, not serialized text.** An earlier version of this charged the
304/// length of the JSON, which is the wrong quantity by up to 42x: a parsed value
305/// is a tree of 32-byte nodes held in vectors that over-allocate, so `[[],[]…]`
306/// costs three bytes on the wire and well over a hundred in memory. Measured
307/// against that estimate, a budget reporting 119 MiB held a process at 5.6 GiB.
308///
309/// Every arm therefore charges at least the node itself, and a container
310/// charges for the slack its backing allocation carries. The result
311/// over-estimates on every adversarial shape and costs honest traffic a couple
312/// of percent, which is the direction a bound has to err in.
313pub(crate) fn approx_bytes(entry: &RecordEntry) -> usize {
314    2 * std::mem::size_of::<RecordEntry>()
315        + entry.uri.len()
316        + entry.cid.as_ref().map_or(0, String::len)
317        + json_bytes(&entry.value)
318}
319
320/// What a parsed JSON value retains, without measuring the heap.
321fn json_bytes(v: &serde_json::Value) -> usize {
322    /// Every value, of every kind, occupies one of these wherever it sits.
323    const NODE: usize = std::mem::size_of::<serde_json::Value>();
324    /// Two nodes per value: the slot it occupies, and the slack the container
325    /// holding it carries — a `Vec` grows by doubling, so up to one spare slot
326    /// per live one.
327    const SLOT: usize = 2 * NODE;
328    /// A map entry is a tree node of its own, with links and a key beside the
329    /// value. Rounded up rather than derived, since the layout is not ours.
330    const MAP_ENTRY: usize = 104;
331    /// A map's backing node, allocated whole.
332    ///
333    /// **Empirical, and not derived from anything the compiler checks.** Unlike
334    /// [`NODE`], which is a `size_of`, this and `MAP_ENTRY` come from measuring
335    /// `std`'s `BTreeMap` layout — B = 6, so eleven pairs to a leaf — under the
336    /// `serde_json` in this lockfile. A toolchain that changes that layout, or a
337    /// `serde_json` that swaps the map type, moves the real cost without moving
338    /// these. The known-answer test below is the tripwire, and it is only as
339    /// good as the day its figures were taken.
340    ///
341    /// `serde_json::Map` is a `BTreeMap` here — no `preserve_order` in the
342    /// lock — and its leaf carries room for eleven pairs whether or not they
343    /// are used, measured at ~632 bytes. So a one-key object costs what an
344    /// eleven-key one does, and a chain of them costs that per level. Charging
345    /// a container's minimum the way an array does under-reports this by about
346    /// half, which is the same failure as the version this replaces, two orders
347    /// of magnitude smaller.
348    const MAP_NODE: usize = 512;
349    match v {
350        // The `4 * NODE` is the container's own minimum allocation; each child
351        // then charges for itself, recursively. Dropping that recursion is what
352        // made an array of empty arrays look free.
353        serde_json::Value::Array(a) => 4 * NODE + a.iter().map(json_bytes).sum::<usize>(),
354        serde_json::Value::Object(o) => {
355            MAP_NODE
356                + o.iter()
357                    .map(|(k, v)| MAP_ENTRY + k.len().max(NODE / 2) + SLOT + json_bytes(v))
358                    .sum::<usize>()
359        }
360        serde_json::Value::String(s) => SLOT + s.len(),
361        // Null, bool and number are all the node and nothing else.
362        _ => SLOT,
363    }
364}
365
366/// Running byte accounting for one walk.
367pub(crate) struct ByteBudget {
368    used: usize,
369    max: usize,
370}
371
372impl ByteBudget {
373    pub(crate) fn new(max: usize) -> Self {
374        Self { used: 0, max }
375    }
376
377    /// Charge a page. `false` when the walk must stop; a refused page is NOT
378    /// charged, so `used` always describes what the caller actually kept.
379    pub(crate) fn admit(&mut self, page: &[RecordEntry]) -> bool {
380        let cost: usize = page.iter().map(approx_bytes).sum();
381        match self.used.checked_add(cost) {
382            Some(total) if total <= self.max => {
383                self.used = total;
384                true
385            }
386            _ => false,
387        }
388    }
389
390    /// Charge `bytes` that no record accounts for. Same contract as
391    /// [`Self::admit`]: `false` means stop, and a refused charge is not taken.
392    pub(crate) fn charge(&mut self, bytes: usize) -> bool {
393        match self.used.checked_add(bytes) {
394            Some(total) if total <= self.max => {
395                self.used = total;
396                true
397            }
398            _ => false,
399        }
400    }
401
402    pub(crate) fn used(&self) -> usize {
403        self.used
404    }
405
406    /// The ceiling this budget was built with.
407    pub(crate) fn max(&self) -> usize {
408        self.max
409    }
410
411    /// What is left. A transient page has to fit in this, not in the ceiling —
412    /// otherwise a walk that has already retained most of its budget can still
413    /// hold a full budget's worth of page on top of it.
414    pub(crate) fn remaining(&self) -> usize {
415        self.max.saturating_sub(self.used)
416    }
417}
418
419/// Append a page, refusing to exceed `max`.
420///
421/// **An error, never a truncation.** The caller of the live walk is
422/// `web::resolve_subscriptions`, whose result reaches `store::replace_sub_refs`
423/// — a `DELETE` followed by reinserting exactly what it was handed. A short
424/// list there is not a short list, it is revoked access to whatever fell off
425/// the end. Returning `Err` lets `resolve_subscriptions` take its documented
426/// fail-closed branch and serve the last-known projection instead.
427///
428/// `out` is left untouched on refusal, so a partial page cannot survive.
429pub(crate) fn extend_bounded(
430    out: &mut Vec<RecordEntry>,
431    page: Vec<RecordEntry>,
432    max: usize,
433    collection: &str,
434) -> Result<()> {
435    if out.len() + page.len() > max {
436        anyhow::bail!(
437            "listRecords for {collection} exceeded the {max}-record cap \
438             ({} held, {} more offered) — refusing to accumulate further",
439            out.len(),
440            page.len(),
441        );
442    }
443    out.extend(page);
444    Ok(())
445}
446
447/// Errors from the atproto identity + PDS layer.
448///
449/// Wraps the transport, the atproto XRPC error envelope (`{"error","message"}`),
450/// and the identity-resolution failure modes so callers can distinguish "the
451/// network broke" from "the PDS said no" from "this handle doesn't resolve".
452#[derive(Debug, thiserror::Error)]
453pub enum AtProtoError {
454    /// The underlying HTTP transport failed (DNS, TLS, timeout, connect).
455    #[error("atproto transport error: {0}")]
456    Transport(#[from] reqwest::Error),
457
458    /// The XRPC endpoint returned a non-2xx status with an atproto error
459    /// envelope (or an opaque body). `error` is the atproto error name (e.g.
460    /// `RecordNotFound`, `AuthMissing`), `message` the human string.
461    #[error("atproto XRPC error {status}: {error}{}", .message.as_deref().map(|m| format!(" — {m}")).unwrap_or_default())]
462    Xrpc {
463        /// The HTTP status code.
464        status: StatusCode,
465        /// The atproto error name (the `error` field), or `"Unknown"`.
466        error: String,
467        /// The optional human-readable `message` field.
468        message: Option<String>,
469    },
470
471    /// A handle could not be resolved to a DID.
472    #[error("could not resolve handle {handle:?} to a DID")]
473    HandleResolution {
474        /// The handle that failed to resolve.
475        handle: String,
476    },
477
478    /// A DID document could not be resolved, or lacks a usable PDS service
479    /// endpoint (`#atproto_pds`).
480    #[error("could not resolve DID {did:?} to a PDS endpoint: {reason}")]
481    DidResolution {
482        /// The DID that failed to resolve.
483        did: String,
484        /// Why resolution failed.
485        reason: String,
486        /// The same, as a value a caller can branch on without reading `reason`.
487        cause: DidResolutionCause,
488    },
489}
490
491/// Why a DID did not resolve to a PDS, structured — so a poller can file a
492/// deleted account under "the server answered" rather than "the network broke".
493#[derive(Debug, Clone, Copy, PartialEq, Eq)]
494#[non_exhaustive]
495pub enum DidResolutionCause {
496    /// A DID method this reader does not resolve.
497    UnsupportedMethod,
498    /// The DID document fetch got an answer, and it was not a success (a
499    /// tombstoned or unknown DID is a 404 from the PLC directory).
500    Status,
501    /// The DID document has no `#atproto_pds` service.
502    NoPdsEndpoint,
503    /// The PDS endpoint it names is refused by the SSRF guard.
504    NotAPublicTarget,
505}
506
507impl AtProtoError {
508    /// True when the XRPC error is a "record not found" — handy for upsert paths
509    /// that treat a missing record as "create instead of update".
510    pub fn is_record_not_found(&self) -> bool {
511        matches!(
512            self,
513            AtProtoError::Xrpc { error, .. } if error == "RecordNotFound"
514        )
515    }
516}
517
518// ---------------------------------------------------------------------------
519// Auth — the direct-PDS path (dev / tests)
520// ---------------------------------------------------------------------------
521
522/// A source of atproto access tokens.
523///
524/// This trait abstracts over token acquisition for the direct [`PdsClient`]
525/// (used by local runs and tests). A [`PdsClient`] can hold a `dyn TokenSource`
526/// instead of a static [`Auth`] without any call-site change, so a token source
527/// that refreshes out of band can be dropped in later.
528///
529/// It is async + `Send + Sync` so a background refresh can live behind it.
530#[allow(async_fn_in_trait)]
531pub trait TokenSource: Send + Sync {
532    /// Return the current bearer access token to send as `Authorization`.
533    async fn access_token(&self) -> Result<String>;
534}
535
536/// The auth material a [`PdsClient`] carries.
537///
538/// A small enum rather than a bare string, so the match stays exhaustive if a
539/// second direct-auth mechanism is added alongside app-password sessions.
540#[derive(Clone)]
541pub enum Auth {
542    /// A bearer access token from a `com.atproto.server.createSession`
543    /// (app-password) session. This is the direct-PDS auth used by local runs
544    /// and tests; the live web path authenticates via the OAuth sidecar instead
545    /// (see [`SidecarClient`]).
546    Session(SessionAuth),
547
548    /// The atproto OAuth confidential-client path is handled entirely by the
549    /// `@atproto/oauth-client` sidecar ([`SidecarClient`]), which mints, DPoP-binds,
550    /// and refreshes tokens. The direct [`PdsClient`] does not carry OAuth tokens;
551    /// this variant is a placeholder so the `Auth` enum documents that the OAuth
552    /// path lives elsewhere.
553    Oauth(OauthPlaceholder),
554
555    /// **No credentials at all** — an unauthenticated public read of a repo the
556    /// caller does not own. `com.atproto.repo.listRecords` is public on a
557    /// standard PDS, so a stranger's `community.lexicon.rss.*` records can be
558    /// read with no session; this variant makes that expressible without
559    /// inventing a fake token.
560    ///
561    /// A client holding it is **read-only**: [`Auth::bearer`] returns an error,
562    /// so every write path (`create_record` / `put_record` / `delete_record` /
563    /// `apply_writes`, all of which go through
564    /// `authed_headers`) fails closed. Construct one
565    /// via [`PdsClient::anonymous`].
566    Anonymous,
567}
568
569impl Auth {
570    /// The bearer access token to present on `com.atproto.repo.*` calls.
571    ///
572    /// Only [`Auth::Session`] carries a token (the session's `accessJwt`).
573    /// [`Auth::Oauth`] carries none — the sidecar owns the OAuth path — so it
574    /// returns an error pointing callers at [`SidecarClient`]. [`Auth::Anonymous`]
575    /// carries none by construction, which is what makes an anonymous client
576    /// read-only.
577    pub fn bearer(&self) -> Result<&str> {
578        match self {
579            Auth::Session(s) => Ok(&s.access_jwt),
580            Auth::Oauth(_) => anyhow::bail!(
581                "the direct PdsClient does not carry OAuth tokens — atproto OAuth is \
582                 handled by the @atproto/oauth-client sidecar (SidecarClient); \
583                 use Auth::Session (app-password) for the direct-PDS path"
584            ),
585            Auth::Anonymous => anyhow::bail!(
586                "this PdsClient is anonymous (unauthenticated public read) and carries no \
587                 bearer token — authenticated repo writes require Auth::Session or the \
588                 SidecarClient"
589            ),
590        }
591    }
592}
593
594/// A session obtained from `com.atproto.server.createSession` (interim
595/// app-password auth). Holds the DID + tokens + handle the server returned.
596#[derive(Clone, Debug, Deserialize)]
597pub struct SessionAuth {
598    /// The account DID this session authenticates.
599    pub did: String,
600    /// The account handle at session-creation time.
601    #[serde(default)]
602    pub handle: Option<String>,
603    /// The bearer access token presented on authed XRPC calls.
604    #[serde(rename = "accessJwt")]
605    pub access_jwt: String,
606    /// The refresh token, exchanged via `com.atproto.server.refreshSession`.
607    /// The direct-PDS refresh flow is not implemented here; the live web path
608    /// refreshes via the OAuth sidecar instead.
609    #[serde(rename = "refreshJwt", default)]
610    pub refresh_jwt: Option<String>,
611}
612
613/// Placeholder for the OAuth variant of [`Auth`].
614///
615/// Intentionally empty: the OAuth session material (DPoP key handle, token
616/// references) is held entirely by the sidecar, not by the direct [`PdsClient`].
617/// This type exists only so [`Auth::Oauth`] is a real variant and the split is
618/// visible in the type system.
619#[derive(Clone, Debug, Default)]
620#[non_exhaustive]
621pub struct OauthPlaceholder {}
622
623// ---------------------------------------------------------------------------
624// Identity resolution
625// ---------------------------------------------------------------------------
626
627/// Resolve an atproto handle to its DID.
628///
629/// Uses `com.atproto.identity.resolveHandle` against `resolver_base` (any host
630/// that implements it; [`DEFAULT_RESOLVER_HOST`] is a safe bootstrap). A fuller
631/// implementation would also try the DNS `_atproto` TXT record and the
632/// `https://<handle>/.well-known/atproto-did` fallback; the XRPC path is the
633/// common case and the one implemented here.
634pub async fn resolve_handle(client: &Client, resolver_base: &str, handle: &str) -> Result<String> {
635    // Build the query manually rather than via reqwest's `.query()` so we don't
636    // depend on the optional `query`/`url` reqwest feature (the declared feature
637    // set is rustls + gzip + json only).
638    let url = format!(
639        "{}/xrpc/com.atproto.identity.resolveHandle?handle={}",
640        resolver_base.trim_end_matches('/'),
641        urlencode(handle)
642    );
643
644    #[derive(Deserialize)]
645    struct ResolveHandleOut {
646        did: String,
647    }
648
649    // Route through the SSRF guard: `resolver_base` can be a user-influenced PDS
650    // host (from a prior DID-doc resolution), so a hostile endpoint must not be
651    // able to target loopback / link-local / metadata. Feed-privacy is NOT
652    // applied here (this is a legitimate atproto XRPC call, not a feed fetch).
653    let resp = crate::net::guarded_get_no_privacy(client, &url, &[]).await?;
654    if !resp.status().is_success() {
655        // Surface the XRPC envelope but map the common "not found" to the typed
656        // handle-resolution error so callers get a clean signal.
657        let err = xrpc_error_from(resp).await;
658        if let AtProtoError::Xrpc { status, .. } = &err {
659            if *status == StatusCode::BAD_REQUEST || *status == StatusCode::NOT_FOUND {
660                return Err(AtProtoError::HandleResolution {
661                    handle: handle.to_string(),
662                }
663                .into());
664            }
665        }
666        return Err(err.into());
667    }
668
669    // Capped: `resolver_base` can be a user-influenced PDS host, as the comment
670    // above this function's guard already says.
671    let raw = crate::net::read_capped(resp).await?;
672    let out: ResolveHandleOut =
673        serde_json::from_slice(&raw).context("parsing resolveHandle response")?;
674    Ok(out.did)
675}
676
677/// Resolve a DID to its PDS service endpoint by fetching + parsing its DID
678/// document.
679///
680/// * `did:plc:*` → the PLC directory (`{plc_directory}/{did}`).
681/// * `did:web:host` → `https://host/.well-known/did.json`.
682///
683/// The PDS endpoint is the service in the DID doc whose `id` ends with
684/// `#atproto_pds` (type `AtprotoPersonalDataServer`); its `serviceEndpoint` is
685/// the base URL for all `com.atproto.repo.*` calls.
686pub async fn resolve_did_to_pds(client: &Client, plc_directory: &str, did: &str) -> Result<String> {
687    let doc_url = if let Some(rest) = did.strip_prefix("did:web:") {
688        // did:web host may itself be percent-encoded / contain a path; the
689        // common case is a bare host.
690        let host = rest.replace(':', "/");
691        format!("https://{host}/.well-known/did.json")
692    } else if did.starts_with("did:plc:") {
693        format!("{}/{}", plc_directory.trim_end_matches('/'), did)
694    } else {
695        return Err(AtProtoError::DidResolution {
696            did: did.to_string(),
697            reason: "unsupported DID method (only did:plc and did:web are handled)".to_string(),
698            cause: DidResolutionCause::UnsupportedMethod,
699        }
700        .into());
701    };
702
703    // SSRF guard: `doc_url` is attacker-controllable for `did:web:<host>` (the
704    // host comes straight from the DID) — a hostile `did:web:169.254.169.254`
705    // or `did:web:localhost` would otherwise make the server fetch an internal
706    // target and reflect its body. Route through the IP/scheme guard (no
707    // feed-privacy layer — this is a DID document, not a feed).
708    let resp = crate::net::guarded_get_no_privacy(client, &doc_url, &[]).await?;
709    if !resp.status().is_success() {
710        return Err(AtProtoError::DidResolution {
711            did: did.to_string(),
712            reason: format!("DID document fetch returned {}", resp.status()),
713            cause: DidResolutionCause::Status,
714        }
715        .into());
716    }
717
718    // **Capped, and this is the most remote-controlled body of the lot.** For a
719    // `did:web:` the host is taken straight out of the DID, so whoever supplies
720    // the DID chooses the server — and the SSRF guard only proves the address is
721    // public, not that the body is finite.
722    let raw = crate::net::read_capped(resp).await?;
723    let doc: DidDocument = serde_json::from_slice(&raw).context("parsing DID document")?;
724    let endpoint = doc
725        .pds_endpoint()
726        .ok_or_else(|| AtProtoError::DidResolution {
727            did: did.to_string(),
728            reason: "DID document has no #atproto_pds service endpoint".to_string(),
729            cause: DidResolutionCause::NoPdsEndpoint,
730        })?;
731
732    // SSRF guard on the RESOLVED endpoint: the `serviceEndpoint` is fully
733    // attacker-controlled (it's whatever the DID document says) and is handed to
734    // XRPC clients that fetch it directly. Reject a private/loopback/metadata
735    // target here so a hostile DID doc can't point the PDS at an internal host.
736    crate::net::assert_public_target(&endpoint)
737        .await
738        .map_err(|e| AtProtoError::DidResolution {
739            did: did.to_string(),
740            reason: format!("PDS serviceEndpoint is not a public target: {e}"),
741            cause: DidResolutionCause::NotAPublicTarget,
742        })?;
743    Ok(endpoint)
744}
745
746/// The subset of a DID document FeatherReader needs: its services, so it can
747/// find the `#atproto_pds` endpoint.
748#[derive(Debug, Clone, Deserialize)]
749pub struct DidDocument {
750    /// The document subject (the DID itself).
751    #[serde(default)]
752    pub id: String,
753    /// The declared services; the PDS is the one whose `id` ends `#atproto_pds`.
754    #[serde(default)]
755    pub service: Vec<DidService>,
756}
757
758/// One service entry in a [`DidDocument`].
759#[derive(Debug, Clone, Deserialize)]
760pub struct DidService {
761    /// The service id fragment (e.g. `#atproto_pds`).
762    pub id: String,
763    /// The service type (e.g. `AtprotoPersonalDataServer`).
764    #[serde(rename = "type", default)]
765    pub r#type: String,
766    /// The service base URL.
767    #[serde(rename = "serviceEndpoint")]
768    pub service_endpoint: String,
769}
770
771impl DidDocument {
772    /// The `#atproto_pds` service endpoint, if present.
773    pub fn pds_endpoint(&self) -> Option<String> {
774        self.service
775            .iter()
776            .find(|s| s.id.ends_with("#atproto_pds"))
777            .map(|s| s.service_endpoint.trim_end_matches('/').to_string())
778    }
779}
780
781// ---------------------------------------------------------------------------
782// Direct-PDS auth: app-password session
783// ---------------------------------------------------------------------------
784
785/// Create a session with an **app password** via
786/// `com.atproto.server.createSession`.
787///
788/// This is the direct-PDS path that makes [`PdsClient`] usable without the OAuth
789/// sidecar (local runs and tests). `pds_base` is the account's PDS (resolve it
790/// first with [`resolve_handle`] + [`resolve_did_to_pds`], or pass the entryway
791/// like `https://bsky.social`, which will service-proxy). `identifier` is a
792/// handle or DID; `app_password` is an app-password (never the main password).
793///
794/// The POST goes through [`crate::net::guarded_post_json`]. This is the single
795/// most credential-dense request in the crate — the app password travels in the
796/// JSON **body**, where reqwest's cross-origin header sanitisation cannot help
797/// it — so it gets the scheme/IP allow-list, the connect pin (no second DNS
798/// resolution to rebind), and a hard refusal to follow a redirect that would
799/// re-send that body to another host.
800pub async fn login_with_app_password(
801    client: &Client,
802    pds_base: &str,
803    identifier: &str,
804    app_password: &str,
805) -> Result<SessionAuth> {
806    let url = format!(
807        "{}/xrpc/com.atproto.server.createSession",
808        pds_base.trim_end_matches('/')
809    );
810    let body = serde_json::to_vec(&json!({ "identifier": identifier, "password": app_password }))
811        .context("serializing createSession request")?;
812    let resp = crate::net::guarded_post_json(client, &url, &[], body).await?;
813    if !resp.status().is_success() {
814        return Err(xrpc_error_from(resp).await.into());
815    }
816    let raw = crate::net::read_capped(resp).await?;
817    serde_json::from_slice(&raw).context("parsing createSession response")
818}
819
820// ---------------------------------------------------------------------------
821// The PDS client
822// ---------------------------------------------------------------------------
823
824/// A lightweight client for one user's PDS repo.
825///
826/// Holds the user's DID (the repo to read/write), the PDS base URL (resolved
827/// from the DID doc), the shared `reqwest::Client`, and the [`Auth`] token.
828/// All the `com.atproto.repo.*` methods below act on `self.did`'s repo.
829///
830/// The client may also be **anonymous** ([`PdsClient::anonymous`]), in which case
831/// it is read-only: it sends no `Authorization` header and every write path
832/// errors out of [`Auth::bearer`].
833///
834/// Cheap to clone (`Arc` internals); one is held per logged-in session.
835#[derive(Clone)]
836pub struct PdsClient {
837    http: Client,
838    /// The PDS base URL, e.g. `https://pds.example.com` (no trailing slash).
839    pds_base: Arc<str>,
840    /// The repo DID all calls target.
841    did: Arc<str>,
842    /// The auth material (an app-password session bearer for the direct path, or
843    /// [`Auth::Anonymous`] for a read-only public read of a stranger's repo).
844    auth: Auth,
845}
846
847/// A single record as returned in a `listRecords` / `getRecord` response.
848///
849/// `value` is the raw record body (with its `$type`); typed wrappers
850/// deserialize it into the matching [`crate::lexicon`] struct.
851#[derive(Debug, Clone, Deserialize)]
852pub struct RecordEntry {
853    /// The `at://did/collection/rkey` strong ref to this record.
854    pub uri: String,
855    /// The record CID (content hash).
856    #[serde(default)]
857    pub cid: Option<String>,
858    /// The raw record body.
859    pub value: Value,
860}
861
862impl RecordEntry {
863    /// The record key (the last `/`-segment of the `at://` URI).
864    pub fn rkey(&self) -> Option<&str> {
865        self.uri.rsplit('/').next()
866    }
867
868    /// Deserialize this record's `value` into a typed lexicon record.
869    pub fn parse<T: DeserializeOwned>(&self) -> Result<T> {
870        serde_json::from_value(self.value.clone())
871            .with_context(|| format!("deserializing record {}", self.uri))
872    }
873}
874
875/// The `com.atproto.repo.listRecords` response envelope.
876#[derive(Debug, Clone, Deserialize)]
877pub struct ListRecordsResponse {
878    /// The page of records.
879    #[serde(default)]
880    pub records: Vec<RecordEntry>,
881    /// The opaque pagination cursor for the next page, if any.
882    #[serde(default)]
883    pub cursor: Option<String>,
884    /// Records on this page whose envelope was malformed and were left out of
885    /// `records` (#177): one bad record no longer fails the page it is on.
886    #[serde(skip)]
887    pub malformed: usize,
888    /// The page's size on the wire, where the client knows it (0 otherwise).
889    /// A walk that SKIPS malformed records charges this to its budget, since
890    /// the skipped records are invisible to the per-record accounting.
891    #[serde(skip)]
892    pub wire_bytes: usize,
893}
894
895/// **A reader's own repo holds records this server cannot read** (#177).
896///
897/// Returned by every walk whose result is written through
898/// `store::replace_sub_refs`. Skipping a record there would silently drop it
899/// from the reader's subscriptions, so the walk refuses instead, and the web
900/// layer recognises this error by type to tell the reader why.
901#[derive(Debug, Clone, PartialEq, Eq)]
902pub struct MalformedRecords {
903    /// The collection being listed.
904    pub collection: String,
905    /// How many records on the refused page were malformed.
906    pub count: usize,
907}
908
909impl std::fmt::Display for MalformedRecords {
910    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
911        write!(
912            f,
913            "{} record(s) in {} have a malformed envelope; refusing the listing rather than \
914             dropping them",
915            self.count, self.collection
916        )
917    }
918}
919
920impl std::error::Error for MalformedRecords {}
921
922/// The `com.atproto.repo.createRecord` / `putRecord` response (a strong ref to
923/// the written record).
924#[derive(Debug, Clone, Deserialize)]
925pub struct WriteResult {
926    /// The `at://` URI of the written record.
927    pub uri: String,
928    /// The record CID after the write.
929    #[serde(default)]
930    pub cid: Option<String>,
931}
932
933impl WriteResult {
934    /// The record key — the last `/`-segment of the `at://` URI.
935    ///
936    /// The reader-facing `add_*` wrappers return this so the web layer can
937    /// address the freshly-created record (delete/rename) without a re-list.
938    pub fn rkey(&self) -> Option<&str> {
939        self.uri.rsplit('/').next()
940    }
941
942    /// The record key as an owned `String`, or the empty string if the URI is
943    /// somehow segment-less (never in practice — a PDS always returns an
944    /// `at://did/collection/rkey`). Convenience for the `-> rkey` wrappers.
945    pub fn into_rkey(self) -> String {
946        self.rkey().unwrap_or_default().to_string()
947    }
948}
949
950impl PdsClient {
951    /// Construct a client against an already-resolved PDS base + DID + auth.
952    pub fn new(
953        http: Client,
954        pds_base: impl Into<String>,
955        did: impl Into<String>,
956        auth: Auth,
957    ) -> Self {
958        Self {
959            http,
960            pds_base: Arc::from(pds_base.into().trim_end_matches('/')),
961            did: Arc::from(did.into()),
962            auth,
963        }
964    }
965
966    /// Construct a **read-only, unauthenticated** client for a public repo the
967    /// caller does not own — `com.atproto.repo.listRecords` is public on a
968    /// standard PDS, so a stranger's records need no credentials.
969    ///
970    /// Every `com.atproto.repo.*` **write** returns an error (there is no bearer;
971    /// see [`Auth::Anonymous`]). Callers are expected to have obtained `pds_base`
972    /// from [`resolve_did_to_pds`], which already runs
973    /// [`crate::net::assert_public_target`] on the resolved `serviceEndpoint` —
974    /// but that is not what makes the fetch safe: every read is **re-vetted at
975    /// fetch time** by [`crate::net::guarded_get_no_privacy`], which closes the
976    /// DNS-rebinding window between resolve and connect. This constructor is
977    /// deliberately synchronous and does no validation of its own, so the
978    /// authoritative check is not duplicated (or, worse, mistaken for sufficient).
979    pub fn anonymous(http: Client, pds_base: impl Into<String>, did: impl Into<String>) -> Self {
980        Self::new(http, pds_base, did, Auth::Anonymous)
981    }
982
983    /// Resolve `handle` → DID → PDS, obtain an app-password session, and build a
984    /// ready-to-use client. A convenience constructor for the direct-PDS path
985    /// that exercises the whole stack end-to-end.
986    ///
987    /// `resolver_base` / `plc_directory` default to [`DEFAULT_RESOLVER_HOST`] /
988    /// [`DEFAULT_PLC_DIRECTORY`] when passed `None`.
989    pub async fn login(
990        http: Client,
991        handle: &str,
992        app_password: &str,
993        resolver_base: Option<&str>,
994        plc_directory: Option<&str>,
995    ) -> Result<Self> {
996        let resolver = resolver_base.unwrap_or(DEFAULT_RESOLVER_HOST);
997        let plc = plc_directory.unwrap_or(DEFAULT_PLC_DIRECTORY);
998
999        let did = resolve_handle(&http, resolver, handle).await?;
1000        let pds_base = resolve_did_to_pds(&http, plc, &did).await?;
1001        let session = login_with_app_password(&http, &pds_base, &did, app_password).await?;
1002
1003        Ok(Self::new(
1004            http,
1005            pds_base,
1006            session.did.clone(),
1007            Auth::Session(session),
1008        ))
1009    }
1010
1011    /// The repo DID this client targets.
1012    pub fn did(&self) -> &str {
1013        &self.did
1014    }
1015
1016    /// The PDS base URL this client talks to.
1017    pub fn pds_base(&self) -> &str {
1018        &self.pds_base
1019    }
1020
1021    /// Build the `Authorization: Bearer …` header pair for an authed write.
1022    ///
1023    /// A `Vec` of pairs rather than a [`HeaderMap`] because every write now goes
1024    /// through [`crate::net::guarded_post_json`], which takes header pairs and
1025    /// sets `Content-Type: application/json` itself. Fails closed on
1026    /// [`Auth::Anonymous`] (there is no bearer), which is what makes an anonymous
1027    /// client read-only.
1028    fn authed_headers(&self) -> Result<Vec<(HeaderName, HeaderValue)>> {
1029        let bearer = self.auth.bearer()?;
1030        let mut value = HeaderValue::from_str(&format!("Bearer {bearer}"))
1031            .context("building Authorization header")?;
1032        value.set_sensitive(true);
1033        Ok(vec![(AUTHORIZATION, value)])
1034    }
1035
1036    fn xrpc_url(&self, method: &str) -> String {
1037        format!("{}/xrpc/{}", self.pds_base, method)
1038    }
1039
1040    // -- com.atproto.repo.* --------------------------------------------------
1041
1042    /// `com.atproto.repo.listRecords` — one page of a collection's records.
1043    ///
1044    /// `cursor` continues a previous page; `limit` caps the page (atproto's max
1045    /// is 100). Use [`list_all_records`](Self::list_all_records) to page fully.
1046    ///
1047    /// The fetch is routed through [`crate::net::guarded_get_no_privacy`] — the
1048    /// same per-hop scheme/IP allow-list and connect-pinning the feed poller and
1049    /// the identity-resolution paths use. `pds_base` was vetted by
1050    /// [`crate::net::assert_public_target`] at resolve time, but that is a
1051    /// *separate* DNS resolution from this fetch; routing the request through the
1052    /// guard closes the rebinding window, which matters as soon as the repo (and
1053    /// therefore the host) is chosen by a stranger. The response body is read via
1054    /// [`crate::net::read_capped`] so a hostile PDS cannot stream an unbounded
1055    /// body at a 512 MB box.
1056    pub async fn list_records(
1057        &self,
1058        collection: &str,
1059        limit: Option<u32>,
1060        cursor: Option<&str>,
1061    ) -> Result<ListRecordsResponse> {
1062        // Build the query manually (see `resolve_handle`): no reqwest `query`
1063        // feature dependency.
1064        let mut url = format!(
1065            "{}?repo={}&collection={}",
1066            self.xrpc_url("com.atproto.repo.listRecords"),
1067            urlencode(&self.did),
1068            urlencode(collection),
1069        );
1070        if let Some(limit) = limit {
1071            url.push_str(&format!("&limit={limit}"));
1072        }
1073        if let Some(cursor) = cursor {
1074            url.push_str(&format!("&cursor={}", urlencode(cursor)));
1075        }
1076
1077        // listRecords is public/unauthenticated on most PDSes, but we send the
1078        // bearer when we have a session one so private repos work too. An
1079        // Auth::Oauth / Auth::Anonymous client sends no Authorization header at
1080        // all. The guard drops the header if a redirect leaves this PDS's origin.
1081        let mut headers: Vec<(reqwest::header::HeaderName, HeaderValue)> = Vec::new();
1082        if let Auth::Session(s) = &self.auth {
1083            let mut value = HeaderValue::from_str(&format!("Bearer {}", s.access_jwt))
1084                .context("building Authorization header")?;
1085            value.set_sensitive(true);
1086            headers.push((AUTHORIZATION, value));
1087        }
1088        let resp = crate::net::guarded_get_no_privacy(&self.http, &url, &headers).await?;
1089        if !resp.status().is_success() {
1090            return Err(xrpc_error_from(resp).await.into());
1091        }
1092        let body = crate::net::read_capped(resp).await?;
1093        parse_list_records(&body)
1094    }
1095
1096    /// Page through **all** records in a collection, following the cursor until
1097    /// exhausted. Convenience over [`list_records`](Self::list_records) for the
1098    /// login-time "load the whole follow-list" read.
1099    ///
1100    /// Bounded by `MAX_LIST_PAGES` and by cursor-repetition detection, because
1101    /// `pds_base` may be a host we did not choose (see [`PdsClient::anonymous`]).
1102    pub async fn list_all_records(&self, collection: &str) -> Result<Vec<RecordEntry>> {
1103        self.list_all_records_within(collection, &mut ByteBudget::new(MAX_LIST_BYTES))
1104            .await
1105    }
1106
1107    /// [`list_all_records`](Self::list_all_records) against a caller's budget.
1108    ///
1109    /// **Shared, not per-walk.** Two walks that nest — a publication read runs a
1110    /// second walk while still holding the first's records — each had their own
1111    /// ceiling, so the process could hold twice it. Passing one budget in makes
1112    /// the bound a property of the caller's whole read, which is the thing that
1113    /// has to fit in the box, and the type enforces it where a comment would not.
1114    pub(crate) async fn list_all_records_within(
1115        &self,
1116        collection: &str,
1117        budget: &mut ByteBudget,
1118    ) -> Result<Vec<RecordEntry>> {
1119        self.walk_all_within(collection, budget, OnMalformed::Refuse)
1120            .await
1121            .map(|(records, _)| records)
1122    }
1123
1124    /// [`list_all_records_within`](Self::list_all_records_within) for a
1125    /// **stranger's** repo: a malformed record is skipped and counted instead of
1126    /// refusing the walk (#177). Never for a walk that reaches
1127    /// `replace_sub_refs`, where a skipped record is a dropped subscription.
1128    pub(crate) async fn list_all_records_skipping_within(
1129        &self,
1130        collection: &str,
1131        budget: &mut ByteBudget,
1132    ) -> Result<(Vec<RecordEntry>, usize)> {
1133        self.walk_all_within(collection, budget, OnMalformed::Skip)
1134            .await
1135    }
1136
1137    async fn walk_all_within(
1138        &self,
1139        collection: &str,
1140        budget: &mut ByteBudget,
1141        on_malformed: OnMalformed,
1142    ) -> Result<(Vec<RecordEntry>, usize)> {
1143        let mut out = Vec::new();
1144        let max_bytes = budget.max();
1145        let mut cursor: Option<String> = None;
1146        let mut more_offered = false;
1147        let mut malformed = 0usize;
1148        for _ in 0..MAX_LIST_PAGES {
1149            let page = self
1150                .list_records(collection, Some(100), cursor.as_deref())
1151                .await?;
1152            if on_malformed == OnMalformed::Refuse {
1153                refuse_malformed(&page, collection)?;
1154            }
1155            // **Skipped records still cost what they cost.** The per-record
1156            // accounting below never sees them, so without this a repo serving
1157            // pages of junk walks every page MAX_LIST_PAGES allows — gigabytes —
1158            // under a budget meant to stop it (found in review). The whole page
1159            // is charged: conservative, and only on pages that skipped something.
1160            // Charged ONCE, at the larger of what crossed the wire and what the
1161            // kept records retain: a parsed record can hold up to 42x its wire
1162            // size, so the wire alone under-charges (third review of #224).
1163            if page.malformed > 0
1164                && !budget.charge(
1165                    page.wire_bytes
1166                        .max(page.records.iter().map(approx_bytes).sum()),
1167                )
1168            {
1169                anyhow::bail!(
1170                    "listRecords for {collection} exceeded the {max_bytes}-byte cap on pages \
1171                     of malformed records — refusing to read further"
1172                );
1173            }
1174            malformed += page.malformed;
1175            // Skipped records count toward "the page had something", so a page
1176            // of nothing BUT malformed records does not end the walk early.
1177            let got = page.records.len() + page.malformed;
1178            // **Refused, not truncated**, for the reason `extend_bounded`
1179            // gives: this walk feeds `replace_sub_refs`, where a short list is
1180            // revoked access.
1181            // A page that skipped records was already charged its whole wire
1182            // size above, which covers its good records too; charging them again
1183            // failed walks that fit (found in review).
1184            if page.malformed == 0 && !budget.admit(&page.records) {
1185                anyhow::bail!(
1186                    "listRecords for {collection} exceeded the {max_bytes}-byte cap \
1187                     ({} held, {} bytes charged) — refusing to accumulate further",
1188                    out.len(),
1189                    budget.used(),
1190                );
1191            }
1192            extend_bounded(&mut out, page.records, MAX_LIST_RECORDS, collection)?;
1193            match page.cursor {
1194                // Guard against a PDS that echoes a cursor with an empty page,
1195                // or that hands back the SAME cursor forever (an infinite walk
1196                // that would otherwise re-count the same page every pass).
1197                Some(next) if got > 0 && Some(&next) != cursor.as_ref() => {
1198                    cursor = Some(next);
1199                    more_offered = true;
1200                }
1201                _ => {
1202                    more_offered = false;
1203                    break;
1204                }
1205            }
1206        }
1207        // **Running out of pages is a refusal, not a short answer.** Falling out
1208        // of the loop used to return `Ok(out)`, so a repo bigger than the page
1209        // budget produced a truncated list indistinguishable from a complete
1210        // one — and `resolve_subscriptions` needs an `Err` for its fail-closed
1211        // branch. Given `Ok`, it hands the short list to `replace_sub_refs`,
1212        // which DELETEs the reader's whole `sub_ref` projection and reinserts
1213        // only what it was given. `extend_bounded` cannot catch this either:
1214        // `MAX_LIST_PAGES` x the 100 we request is `MAX_LIST_RECORDS`, so the
1215        // page budget runs out first.
1216        //
1217        // **The cap is on REQUESTS, so where it bites in RECORDS is the server's
1218        // choice and not ours.** We ask for 100 a page; a PDS MAY answer with
1219        // fewer, and only one that honours the limit puts the boundary anywhere
1220        // near `MAX_LIST_PAGES` x 100. Halve the page size and the same budget
1221        // reaches half as many records; a server that returns MORE than asked
1222        // trips `extend_bounded` first, which is the case the sentence above does
1223        // not cover. Said this way because an earlier version of this comment
1224        // named a fixed record window as though our own constants decided it.
1225        //
1226        // **And at the boundary the refusal is a FALSE one.** Terminating costs
1227        // one extra request, because a short page can still carry a cursor — this
1228        // project's own PDS does exactly that — so a walk that fills its last
1229        // allowed page is holding every record it was ever going to hold and
1230        // refuses anyway, on the strength of a cursor it never followed. With
1231        // `limit=100` honoured that window is a repo of roughly 19 901 to 20 000
1232        // records. The direction is safe and the alternative is deleting feeds,
1233        // but it is a false refusal and not a clean boundary.
1234        if more_offered {
1235            anyhow::bail!(
1236                "listRecords for {collection} did not finish within {MAX_LIST_PAGES} pages \
1237                 ({} held, and the PDS still offered more) — refusing a short list",
1238                out.len(),
1239            );
1240        }
1241        Ok((out, malformed))
1242    }
1243
1244    /// See [`RecordWalk`].
1245    /// The most recent records of a collection that the caller **keeps**,
1246    /// truncating rather than refusing.
1247    ///
1248    /// **The cap counts kept records, not walked ones.** Applying it to the
1249    /// raw collection starves a caller whose filter is selective: a quiet
1250    /// standard.site publication in a repo whose busy sibling fills the
1251    /// window returns nothing at all, permanently, and worse with every post
1252    /// the sibling makes. `MAX_LIST_PAGES` still bounds the request count, so
1253    /// a filter that matches nothing costs a fixed number of round trips.
1254    ///
1255    /// Truncating, not refusing, because this is an additive read: see
1256    /// `extend_truncating` for why the `extend_bounded` refusal would be
1257    /// strictly worse here.
1258    ///
1259    /// `page_size` is the caller's, because the right page depends on how big
1260    /// the records are: [`crate::net::read_capped`] bounds a response at 8 MB,
1261    /// so 100 long-form articles per page can exceed it and fail the whole
1262    /// walk.
1263    ///
1264    /// **Ordering is the PDS's**: `listRecords` is descending by *rkey*, which
1265    /// is newest-first only when rkeys are TIDs. For a publisher using slug
1266    /// rkeys the truncation keeps a lexicographic subset rather than a recent
1267    /// one — acceptable because the cap is now per-publication rather than
1268    /// per-repo, so reaching it at all means an archive larger than this
1269    /// reader stores.
1270    pub async fn list_recent_matching(
1271        &self,
1272        collection: &str,
1273        max_records: usize,
1274        page_size: u32,
1275        keep: impl FnMut(&RecordEntry) -> bool,
1276    ) -> Result<RecordWalk> {
1277        self.list_recent_matching_within(
1278            collection,
1279            max_records,
1280            &mut ByteBudget::new(MAX_LIST_BYTES),
1281            page_size,
1282            keep,
1283        )
1284        .await
1285    }
1286
1287    /// [`list_recent_matching`](Self::list_recent_matching) against a caller's
1288    /// budget. See [`list_all_records_within`](Self::list_all_records_within) for
1289    /// why it is the caller's and not the walk's.
1290    pub(crate) async fn list_recent_matching_within(
1291        &self,
1292        collection: &str,
1293        max_records: usize,
1294        budget: &mut ByteBudget,
1295        page_size: u32,
1296        mut keep: impl FnMut(&RecordEntry) -> bool,
1297    ) -> Result<RecordWalk> {
1298        let mut out = Vec::new();
1299        let mut cursor: Option<String> = None;
1300        let mut malformed = 0usize;
1301        // Every exit carries the skipped count, so none can forget it.
1302        let walk = |mut w: RecordWalk, malformed: usize| {
1303            w.malformed = malformed;
1304            w
1305        };
1306        for _ in 0..MAX_LIST_PAGES {
1307            let page = self
1308                .list_records(collection, Some(page_size), cursor.as_deref())
1309                .await?;
1310            // **Skipped, not refused**: this reads a stranger's collection, and
1311            // one bad record must not stall everything beside it (#177). Counted
1312            // in `got` too, so a page whose only records were malformed is not
1313            // mistaken for the end of the collection.
1314            // As in `walk_all_within`: skipped records are invisible to the
1315            // per-record charge, so their page pays for them here.
1316            let wire_charged = page.malformed > 0;
1317            malformed += page.malformed;
1318            let got = page.records.len() + page.malformed;
1319            // Is there a next page that is actually new? (A PDS may echo a
1320            // cursor with an empty page, or hand back the same one forever.)
1321            let more =
1322                matches!(&page.cursor, Some(next) if got > 0 && Some(next) != cursor.as_ref());
1323            // **The transient page needs its own bound.** The running total
1324            // charges what is KEPT, because that is what the walk retains and a
1325            // page is dropped after the filter. But transient is not free, and
1326            // `net::read_capped`'s 8 MB bounds the WIRE — the whole point of
1327            // this budget is that wire size and retained size are not the same
1328            // number. A page of records the filter rejects entirely charges
1329            // nothing against the total and can still hold hundreds of
1330            // megabytes, so no single page may exceed the walk's budget alone.
1331            // On a page that skipped records, the transient is the larger of
1332            // its wire size and its parsed records: the skipped ones are
1333            // invisible to the per-record sum.
1334            let mut page_cost: usize = page.records.iter().map(approx_bytes).sum();
1335            if wire_charged {
1336                page_cost = page_cost.max(page.wire_bytes);
1337            }
1338            if page_cost > budget.remaining() {
1339                return Ok(walk(RecordWalk::partial(out), malformed));
1340            }
1341            let kept: Vec<RecordEntry> = page.records.into_iter().filter(|r| keep(r)).collect();
1342            // Charged on what is KEPT, which is what this walk retains. **This
1343            // cannot refuse**, and saying so matters: the page check above
1344            // already proved the whole page fits in the remainder, and `kept` is
1345            // a subset of it. Written as `if !admit(…) { return }` it reads as a
1346            // second stopping rule, and a reader would look for the case that
1347            // trips it. There isn't one — the call is the bookkeeping.
1348            // A page that skipped records also pays for the skipped bytes,
1349            // once: the larger of its wire size and what it keeps.
1350            let charged = if wire_charged {
1351                budget.charge(page.wire_bytes.max(kept.iter().map(approx_bytes).sum()))
1352            } else {
1353                budget.admit(&kept)
1354            };
1355            debug_assert!(charged, "the page charge already proved this fits");
1356            if extend_truncating(&mut out, kept, max_records) {
1357                return Ok(walk(RecordWalk::partial(out), malformed));
1358            }
1359            if out.len() >= max_records {
1360                // Landing exactly on the cap is only a truncation if the
1361                // collection had more to give — `extend_truncating` cannot see
1362                // that, so the caller's "incomplete" signal is decided here.
1363                return Ok(RecordWalk {
1364                    complete: !more,
1365                    records: out,
1366                    malformed,
1367                });
1368            }
1369            if !more {
1370                return Ok(walk(RecordWalk::complete(out), malformed));
1371            }
1372            cursor = page.cursor;
1373        }
1374        // **The page budget ran out with the collection still going.** Silence
1375        // here reintroduces the starvation this function exists to prevent, one
1376        // order of magnitude further out: a quiet publication in a repo whose
1377        // busy sibling has more records than MAX_LIST_PAGES × page_size can
1378        // reach returns nothing at all, forever, having spent every round trip
1379        // to find out. The caller is told so it can say which feed.
1380        Ok(walk(RecordWalk::partial(out), malformed))
1381    }
1382
1383    /// `com.atproto.repo.createRecord` — create a new record (server assigns the
1384    /// rkey, `key: tid`). Returns the written record's strong ref.
1385    /// **Private, not `pub` — and not `pub(crate)`.** This is generic over
1386    /// `T: Serialize`, so it will happily write a raw `lexicon::Subscription`:
1387    /// the general case of the hole `create_subscriptions_batch` was one
1388    /// instance of. The vetted wrappers in this `impl` are the sanctioned entry
1389    /// points. `pub(crate)` was tried first and stops nothing that matters — a
1390    /// handler in `web.rs` is in this crate. Private is what makes the wrappers
1391    /// a fact rather than a convention, and it costs nothing: nothing outside
1392    /// this module ever called it.
1393    async fn create_record<T: Serialize>(
1394        &self,
1395        collection: &str,
1396        record: &T,
1397    ) -> Result<WriteResult> {
1398        let body = json!({
1399            "repo": self.did.as_ref(),
1400            "collection": collection,
1401            "record": record,
1402        });
1403        self.repo_write("com.atproto.repo.createRecord", body).await
1404    }
1405
1406    /// `com.atproto.repo.putRecord` — upsert a record at a **known** rkey
1407    /// (`key: any`). This is the `readState` upsert primitive: a feed-derived
1408    /// rkey makes the write idempotent (one record per feed).
1409    /// **Private, not `pub` — and not `pub(crate)`.** This is generic over
1410    /// `T: Serialize`, so it will happily write a raw `lexicon::Subscription`:
1411    /// the general case of the hole `create_subscriptions_batch` was one
1412    /// instance of. The vetted wrappers in this `impl` are the sanctioned entry
1413    /// points. `pub(crate)` was tried first and stops nothing that matters — a
1414    /// handler in `web.rs` is in this crate. Private is what makes the wrappers
1415    /// a fact rather than a convention, and it costs nothing: nothing outside
1416    /// this module ever called it.
1417    async fn put_record<T: Serialize>(
1418        &self,
1419        collection: &str,
1420        rkey: &str,
1421        record: &T,
1422    ) -> Result<WriteResult> {
1423        let body = json!({
1424            "repo": self.did.as_ref(),
1425            "collection": collection,
1426            "rkey": rkey,
1427            "record": record,
1428        });
1429        self.repo_write("com.atproto.repo.putRecord", body).await
1430    }
1431
1432    /// The single outbound path for every authenticated `com.atproto.repo.*`
1433    /// **write**, routed through [`crate::net::guarded_post_json`].
1434    ///
1435    /// Reads were hardened first (see [`list_records`](Self::list_records)), but
1436    /// the argument applies with more force here: `pds_base` is vetted by
1437    /// [`crate::net::assert_public_target`] at *resolve* time, and the write is a
1438    /// *separate* DNS resolution — the rebinding window `net.rs` exists to close.
1439    /// A write also carries the session bearer and, in
1440    /// [`login_with_app_password`], the app password itself, so the guard's
1441    /// refusal to follow redirects (a `307` re-sends the body verbatim to the new
1442    /// host) is doing real work and not just symmetry.
1443    async fn guarded_post(&self, url: &str, body: &Value) -> Result<reqwest::Response> {
1444        let headers = self.authed_headers()?;
1445        let payload = serde_json::to_vec(body).context("serializing XRPC request body")?;
1446        crate::net::guarded_post_json(&self.http, url, &headers, payload).await
1447    }
1448
1449    /// `com.atproto.repo.deleteRecord` — delete a record by collection + rkey
1450    /// (e.g. unsubscribe → delete the subscription record).
1451    pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
1452        let url = self.xrpc_url("com.atproto.repo.deleteRecord");
1453        let body = json!({
1454            "repo": self.did.as_ref(),
1455            "collection": collection,
1456            "rkey": rkey,
1457        });
1458        let resp = self.guarded_post(&url, &body).await?;
1459        if !resp.status().is_success() {
1460            return Err(xrpc_error_from(resp).await.into());
1461        }
1462        Ok(())
1463    }
1464
1465    /// `com.atproto.repo.applyWrites` — a **batch** of create/update/delete
1466    /// operations in one atomic-per-repo round-trip.
1467    ///
1468    /// This is the read-state flusher's workhorse: dozens of dirty per-feed
1469    /// [`ReadState`] cursors coalesce into one call rather than one `putRecord`
1470    /// each. See [`flush_read_states`](Self::flush_read_states).
1471    /// **Private, not `pub` — and not `pub(crate)`.** This is generic over
1472    /// `T: Serialize`, so it will happily write a raw `lexicon::Subscription`:
1473    /// the general case of the hole `create_subscriptions_batch` was one
1474    /// instance of. The vetted wrappers in this `impl` are the sanctioned entry
1475    /// points. `pub(crate)` was tried first and stops nothing that matters — a
1476    /// handler in `web.rs` is in this crate. Private is what makes the wrappers
1477    /// a fact rather than a convention, and it costs nothing: nothing outside
1478    /// this module ever called it.
1479    ///
1480    /// Sent in chunks within the PDS's limits — see [`apply_writes_chunked`]
1481    /// for what a failure part-way means. This used to send any batch as one
1482    /// call, including an empty one; an empty batch now sends nothing, as the
1483    /// other two clients already did.
1484    async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
1485        apply_writes_chunked(writes, |chunk| self.apply_writes_once(chunk)).await
1486    }
1487
1488    /// One `applyWrites` call, unchunked. Reached only through
1489    /// [`apply_writes`](Self::apply_writes).
1490    async fn apply_writes_once(&self, writes: &[WriteOp]) -> Result<()> {
1491        let url = self.xrpc_url("com.atproto.repo.applyWrites");
1492        let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
1493        let body = json!({
1494            "repo": self.did.as_ref(),
1495            "writes": ops,
1496        });
1497        let resp = self.guarded_post(&url, &body).await?;
1498        if !resp.status().is_success() {
1499            return Err(xrpc_error_from(resp).await.into());
1500        }
1501        Ok(())
1502    }
1503
1504    /// Shared create/put path (both return a `{uri,cid}` strong ref).
1505    async fn repo_write(&self, method: &str, body: Value) -> Result<WriteResult> {
1506        let url = self.xrpc_url(method);
1507        let resp = self.guarded_post(&url, &body).await?;
1508        if !resp.status().is_success() {
1509            return Err(xrpc_error_from(resp).await.into());
1510        }
1511        // `read_capped` rather than `resp.json()`: a hostile PDS must not be able
1512        // to stream an unbounded body at a 512 MB box (same rule as the reads).
1513        let raw = crate::net::read_capped(resp).await?;
1514        serde_json::from_slice(&raw).with_context(|| format!("parsing {method} response"))
1515    }
1516
1517    // -- typed lexicon wrappers ---------------------------------------------
1518
1519    /// List every [`Subscription`] record in the user's repo (paged fully). The
1520    /// login-time "what does this user follow?" read.
1521    pub async fn list_subscriptions(&self) -> Result<Vec<(String, Subscription)>> {
1522        self.list_typed(lexicon::nsid::SUBSCRIPTION).await
1523    }
1524
1525    /// Create a [`Subscription`] record (subscribe to a feed).
1526    pub async fn create_subscription(
1527        &self,
1528        sub: &crate::vetted::VettedSubscription,
1529    ) -> Result<WriteResult> {
1530        self.create_record(lexicon::nsid::SUBSCRIPTION, sub).await
1531    }
1532
1533    /// List every [`Folder`] record in the user's repo.
1534    pub async fn list_folders(&self) -> Result<Vec<(String, Folder)>> {
1535        self.list_typed(lexicon::nsid::FOLDER).await
1536    }
1537
1538    /// Create a [`Folder`] record.
1539    pub async fn create_folder(&self, folder: &Folder) -> Result<WriteResult> {
1540        self.create_record(lexicon::nsid::FOLDER, folder).await
1541    }
1542
1543    /// List every [`Saved`] (starred) record in the user's repo.
1544    pub async fn list_saved(&self) -> Result<Vec<(String, Saved)>> {
1545        self.list_typed(lexicon::nsid::SAVED).await
1546    }
1547
1548    /// Create a [`Saved`] record (star an article).
1549    pub async fn create_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<WriteResult> {
1550        self.create_record(lexicon::nsid::SAVED, saved).await
1551    }
1552
1553    /// List every [`ReadState`] cursor in the user's repo (the read side a
1554    /// login-time read-state merge would consume).
1555    pub async fn list_read_states(&self) -> Result<Vec<(String, ReadState)>> {
1556        self.list_typed(lexicon::nsid::READ_STATE).await
1557    }
1558
1559    /// Upsert a single [`ReadState`] cursor at its feed-derived rkey. For a
1560    /// batch of dirty cursors prefer [`flush_read_states`](Self::flush_read_states).
1561    pub async fn put_read_state(&self, rkey: &str, state: &ReadState) -> Result<WriteResult> {
1562        self.put_record(lexicon::nsid::READ_STATE, rkey, state)
1563            .await
1564    }
1565
1566    /// Batch-flush many dirty [`ReadState`] cursors via `applyWrites` (chunked) —
1567    /// the debounced read-state flusher's coalesced write.
1568    ///
1569    /// Each `(rkey, state, pds_created)` becomes a `create` op at the feed-derived
1570    /// rkey when the record does not yet exist, and an `update` when it does — so a
1571    /// feed's FIRST flush succeeds (an `#update` on a missing record errors, and
1572    /// `applyWrites` is atomic per-repo). Both kinds ride the same batch.
1573    pub async fn flush_read_states(&self, cursors: &[(String, ReadState, bool)]) -> Result<()> {
1574        if cursors.is_empty() {
1575            return Ok(());
1576        }
1577        let writes = read_state_write_ops(cursors)?;
1578        self.apply_writes(&writes).await
1579    }
1580
1581    /// List a collection and parse each record's value into `T`, pairing it with
1582    /// its rkey. Records that fail to deserialize are skipped with a warning
1583    /// (forward-compat: a future writer's extra fields shouldn't break login).
1584    async fn list_typed<T: DeserializeOwned>(&self, collection: &str) -> Result<Vec<(String, T)>> {
1585        let records = self.list_all_records(collection).await?;
1586        let mut out = Vec::with_capacity(records.len());
1587        for rec in records {
1588            let rkey = rec.rkey().unwrap_or_default().to_string();
1589            match rec.parse::<T>() {
1590                Ok(value) => out.push((rkey, value)),
1591                Err(e) => tracing::warn!(
1592                    collection,
1593                    uri = %rec.uri,
1594                    error = %e,
1595                    "skipping unparseable record in collection"
1596                ),
1597            }
1598        }
1599        Ok(out)
1600    }
1601}
1602
1603// ---------------------------------------------------------------------------
1604// The OAuth sidecar client — the LIVE com.atproto.repo.* path
1605// ---------------------------------------------------------------------------
1606
1607/// A client for the atproto OAuth sidecar's **internal** API.
1608///
1609/// This is the live path for every authed repo operation. Rather than the Rust
1610/// server holding PDS tokens, it POSTs `{did, action, …}` to the sidecar's
1611/// `/internal/repo` endpoint (gated by the shared `X-Internal-Secret`); the
1612/// sidecar `restore(did)`s the OAuth session — transparent DPoP + token refresh —
1613/// and runs the matching XRPC call via `@atproto/api`. The `did` (plus the shared
1614/// secret) is what authorizes the call; there is no bearer token on the Rust side.
1615///
1616/// It also fronts `/internal/session/:id`, the one-shot handoff the Rust callback
1617/// uses to turn a `session_id` (from the sidecar's browser redirect) into the
1618/// `{did, handle}` it keys its own signed cookie by.
1619///
1620/// Cheap to clone (shared `reqwest::Client` + `Arc`'d config).
1621#[derive(Clone)]
1622pub struct SidecarClient {
1623    http: Client,
1624    public_url: Arc<str>,
1625    internal_url: Arc<str>,
1626    internal_secret: Arc<str>,
1627}
1628
1629/// The `{did, handle}` a session-id resolves to (the sidecar's
1630/// `/internal/session/:id` body).
1631#[derive(Debug, Clone, Deserialize)]
1632pub struct SidecarSession {
1633    /// The account DID that logged in.
1634    pub did: String,
1635    /// The account handle at login time.
1636    #[serde(default)]
1637    pub handle: Option<String>,
1638}
1639
1640/// The sidecar's `/internal/revoke` response body:
1641/// `{ ok:true, did, revoked, hadSession }`.
1642#[derive(Debug, Clone, Deserialize)]
1643pub struct RevokeResult {
1644    /// The DID that was revoked.
1645    #[serde(default)]
1646    pub did: String,
1647    /// Whether the OAuth token revocation at the PDS succeeded. `false` means
1648    /// the local rows were still purged (best-effort), but the PDS-side tokens
1649    /// may not have been invalidated (network failure).
1650    #[serde(default)]
1651    pub revoked: bool,
1652    /// Whether the sidecar actually had a stored session for the DID.
1653    #[serde(default, rename = "hadSession")]
1654    pub had_session: bool,
1655}
1656
1657/// The action verbs the sidecar's `/internal/repo` endpoint dispatches on.
1658#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1659pub enum RepoAction {
1660    /// `com.atproto.repo.listRecords`.
1661    List,
1662    /// `com.atproto.repo.createRecord`.
1663    Create,
1664    /// `com.atproto.repo.putRecord`.
1665    Put,
1666    /// `com.atproto.repo.deleteRecord`.
1667    Delete,
1668    /// `com.atproto.repo.applyWrites` (batch).
1669    ApplyWrites,
1670}
1671
1672impl RepoAction {
1673    fn as_str(self) -> &'static str {
1674        match self {
1675            RepoAction::List => "list",
1676            RepoAction::Create => "create",
1677            RepoAction::Put => "put",
1678            RepoAction::Delete => "delete",
1679            RepoAction::ApplyWrites => "applyWrites",
1680        }
1681    }
1682}
1683
1684/// The `/internal/repo` success envelope: `{ ok:true, data:<raw XRPC JSON> }`.
1685#[derive(Debug, Deserialize)]
1686struct RepoOk {
1687    #[serde(default)]
1688    data: Value,
1689}
1690
1691/// `/internal/repo`'s ok envelope for a **listing**, typed all the way down.
1692///
1693/// `data` absent is not `data` empty, the same distinction `records` carries: it
1694/// is what a proxy makes of an unexpected upstream body.
1695#[derive(Debug, Deserialize)]
1696struct RepoOkList {
1697    /// The sidecar's own `ok`, which it can set false on a 200.
1698    #[serde(default)]
1699    ok: Option<bool>,
1700    /// And its own `error` — a different envelope from the PDS's, one layer out.
1701    #[serde(default)]
1702    error: Option<Value>,
1703    #[serde(default)]
1704    data: Option<ListRecordsBody>,
1705}
1706
1707/// The `/internal/repo` error envelope: `{ ok:false, error, message, status? }`.
1708#[derive(Debug, Deserialize)]
1709struct RepoErr {
1710    #[serde(default)]
1711    error: Option<String>,
1712    #[serde(default)]
1713    message: Option<String>,
1714    #[serde(default)]
1715    status: Option<u16>,
1716}
1717
1718impl SidecarClient {
1719    /// Build a sidecar client from the shared `reqwest::Client` and the
1720    /// resolved public + internal base URLs + internal secret (from
1721    /// [`crate::config::SidecarConfig`]). `public_url` anchors the browser
1722    /// `/login` redirect; `internal_url` is the loopback base for the `/internal/*`
1723    /// API (they collapse to the same value in single-URL local dev).
1724    pub fn new(
1725        http: Client,
1726        public_url: impl Into<String>,
1727        internal_url: impl Into<String>,
1728        internal_secret: impl Into<String>,
1729    ) -> Self {
1730        Self {
1731            http,
1732            public_url: Arc::from(public_url.into().trim_end_matches('/')),
1733            internal_url: Arc::from(internal_url.into().trim_end_matches('/')),
1734            internal_secret: Arc::from(internal_secret.into()),
1735        }
1736    }
1737
1738    /// The sidecar's public `/login` URL for a handle, round-tripping an opaque
1739    /// `return` value through OAuth state (used to bounce the browser back to a
1740    /// specific place after login). The browser is redirected here.
1741    pub fn login_url(&self, handle: &str, return_to: Option<&str>) -> String {
1742        let mut url = format!("{}/login?handle={}", self.public_url, urlencode(handle));
1743        if let Some(r) = return_to {
1744            url.push_str(&format!("&return={}", urlencode(r)));
1745        }
1746        url
1747    }
1748
1749    /// Resolve a one-shot `session_id` (from the sidecar's post-OAuth redirect)
1750    /// to the `{did, handle}` that logged in. `Ok(None)` on `404 SessionNotFound`.
1751    pub async fn resolve_session(&self, session_id: &str) -> Result<Option<SidecarSession>> {
1752        let url = format!(
1753            "{}/internal/session/{}",
1754            self.internal_url,
1755            urlencode(session_id)
1756        );
1757        let resp = self
1758            .http
1759            .get(&url)
1760            .header("X-Internal-Secret", self.internal_secret.as_ref())
1761            .send()
1762            .await?;
1763        if resp.status() == StatusCode::NOT_FOUND {
1764            return Ok(None);
1765        }
1766        if !resp.status().is_success() {
1767            return Err(xrpc_error_from(resp).await.into());
1768        }
1769        let raw = crate::net::read_capped(resp).await?;
1770        let session: SidecarSession =
1771            serde_json::from_slice(&raw).context("parsing /internal/session response")?;
1772        Ok(Some(session))
1773    }
1774
1775    /// Revoke a DID's OAuth session at the sidecar: `POST /internal/revoke`.
1776    ///
1777    /// This revokes the refresh + access tokens at the PDS **and** purges the
1778    /// sidecar's stored `oauth_session` + `app_session` rows for the DID. It is
1779    /// idempotent — revoking a DID with no live session returns
1780    /// `had_session: false`. Called on `/logout` (so the cookie clear isn't the
1781    /// only thing that ends the session) and on `/account/delete`.
1782    pub async fn revoke_session(&self, did: &str) -> Result<RevokeResult> {
1783        let url = format!("{}/internal/revoke", self.internal_url);
1784        let resp = self
1785            .http
1786            .post(&url)
1787            .header("X-Internal-Secret", self.internal_secret.as_ref())
1788            .json(&json!({ "did": did }))
1789            .send()
1790            .await?;
1791        if !resp.status().is_success() {
1792            return Err(xrpc_error_from(resp).await.into());
1793        }
1794        let raw = crate::net::read_capped(resp).await?;
1795        let result: RevokeResult =
1796            serde_json::from_slice(&raw).context("parsing /internal/revoke response")?;
1797        Ok(result)
1798    }
1799
1800    /// POST one op to `/internal/repo` and return the raw XRPC `data` payload.
1801    ///
1802    /// `body` must already carry `did` + `action` + the action's required fields
1803    /// (the typed wrappers below build these). Maps the sidecar's error envelope
1804    /// to [`AtProtoError`]: `404 SessionNotFound` → `Xrpc{error:"SessionNotFound"}`
1805    /// so callers can treat it as "re-login required".
1806    async fn repo(&self, body: Value) -> Result<Value> {
1807        let raw = self.repo_bytes(body).await?;
1808        // **The same guard as the listing path, twelve lines below.** This is
1809        // the response path for every write — create, put, delete, applyWrites —
1810        // and `RepoOk.data` is an unbounded `Value`. Measured without it: the
1811        // identical 8 MB attack retains 786 MB. The listing had the guard and
1812        // this did not, which is the drift a shared helper exists to prevent.
1813        refuse_a_structure_explosion(&raw, "the /internal/repo body")?;
1814        let ok: RepoOk = serde_json::from_slice(&raw).context("parsing /internal/repo ok body")?;
1815        Ok(ok.data)
1816    }
1817
1818    /// [`repo`](Self::repo) without the `Value`.
1819    ///
1820    /// The listing path needs the bytes: a `serde_json::Map` resolves a repeated
1821    /// key last-wins, so a body carrying a second, empty `records` array read as
1822    /// a successful empty page — and an empty page here is `replace_sub_refs`
1823    /// deleting every `sub_ref` the reader has. serde refuses a duplicated field
1824    /// outright, but only if it sees the bytes.
1825    async fn repo_bytes(&self, body: Value) -> Result<Vec<u8>> {
1826        let url = format!("{}/internal/repo", self.internal_url);
1827        let resp = self
1828            .http
1829            .post(&url)
1830            .header("X-Internal-Secret", self.internal_secret.as_ref())
1831            .json(&body)
1832            .send()
1833            .await?;
1834        let status = resp.status();
1835        // **Capped, like every other body this codebase reads.** `resp.json()`
1836        // buffers whatever arrives; `/internal/repo` proxies the account's PDS,
1837        // so that length is chosen by a host the reader picked and we did not.
1838        // The 8 MB ceiling that bounds the direct client did not exist here,
1839        // and the sidecar is the default backend — so the one path with no byte
1840        // bound at all was the one most deployments run.
1841        let raw = crate::net::read_capped(resp).await;
1842        if status.is_success() {
1843            return raw;
1844        }
1845        // Error path: parse the sidecar's `{ok:false,error,message,status}` shape.
1846        //
1847        // **A body we could not read must not cost us the status.** Reading
1848        // before the branch was the obvious shape and it swallowed the HTTP
1849        // status on an over-cap or truncated error body, turning a `404
1850        // SessionNotFound` into a bare "body exceeded the cap". `xrpc_error_from`
1851        // already makes the opposite choice deliberately, for the same reason.
1852        let err: RepoErr = raw
1853            .ok()
1854            .and_then(|body| serde_json::from_slice(&body).ok())
1855            .unwrap_or(RepoErr {
1856                error: None,
1857                message: None,
1858                status: None,
1859            });
1860        let mapped = err
1861            .status
1862            .and_then(|s| StatusCode::from_u16(s).ok())
1863            .unwrap_or(status);
1864        Err(AtProtoError::Xrpc {
1865            status: mapped,
1866            error: err.error.unwrap_or_else(|| "Unknown".to_string()),
1867            message: err.message,
1868        }
1869        .into())
1870    }
1871
1872    // -- raw com.atproto.repo.* over the sidecar -----------------------------
1873
1874    /// `list` — one page of a collection's records for `did`.
1875    pub async fn list_records(
1876        &self,
1877        did: &str,
1878        collection: &str,
1879        limit: Option<u32>,
1880        cursor: Option<&str>,
1881    ) -> Result<ListRecordsResponse> {
1882        let mut body = json!({
1883            "did": did,
1884            "action": RepoAction::List.as_str(),
1885            "collection": collection,
1886        });
1887        if let Some(limit) = limit {
1888            body["limit"] = json!(limit);
1889        }
1890        if let Some(cursor) = cursor {
1891            body["cursor"] = json!(cursor);
1892        }
1893        // Straight into the shared wire struct, one parse, no `Value` between —
1894        // so this client gets the same guards as the other two, including the
1895        // duplicated-key refusal that only serde can make.
1896        let raw = self.repo_bytes(body).await?;
1897        refuse_a_structure_explosion(&raw, "the sidecar listRecords body")?;
1898        let envelope: RepoOkList =
1899            serde_json::from_slice(&raw).context("parsing sidecar listRecords data")?;
1900        // **Two envelope layers here, not one.** `page_from_body` guards the
1901        // PDS's, which arrives inside `data`; this is the sidecar's own, and it
1902        // can say `{"ok":false,"error":"ExpiredToken"}` on a 200 while still
1903        // carrying a `data` that reads as a perfectly good empty page.
1904        if envelope.ok == Some(false) {
1905            let name = envelope
1906                .error
1907                .as_ref()
1908                .and_then(envelope_error_name)
1909                .unwrap_or_else(|| "unspecified".to_string());
1910            anyhow::bail!("the sidecar answered 2xx with ok:false ({name})");
1911        }
1912        if let Some(name) = envelope.error.as_ref().and_then(envelope_error_name) {
1913            anyhow::bail!("the sidecar answered 2xx with an error envelope: {name}");
1914        }
1915        let Some(data) = envelope.data else {
1916            anyhow::bail!("listRecords returned no records field (empty or unexpected body)");
1917        };
1918        page_from_body(data).context("parsing sidecar listRecords data")
1919    }
1920
1921    /// Page through **all** records in a collection for `did`.
1922    ///
1923    /// Bounded by `MAX_LIST_PAGES` and cursor-repetition detection, same as
1924    /// [`PdsClient::list_all_records`] — the sidecar proxies to the account's
1925    /// PDS, so the page count is ultimately remote-controlled here too.
1926    pub async fn list_all_records(&self, did: &str, collection: &str) -> Result<Vec<RecordEntry>> {
1927        self.list_all_records_within(did, collection, &mut ByteBudget::new(MAX_LIST_BYTES))
1928            .await
1929    }
1930
1931    /// [`list_all_records`](Self::list_all_records) against a caller's budget.
1932    pub(crate) async fn list_all_records_within(
1933        &self,
1934        did: &str,
1935        collection: &str,
1936        budget: &mut ByteBudget,
1937    ) -> Result<Vec<RecordEntry>> {
1938        let mut out = Vec::new();
1939        let max_bytes = budget.max();
1940        let mut cursor: Option<String> = None;
1941        let mut more_offered = false;
1942        for _ in 0..MAX_LIST_PAGES {
1943            let page = self
1944                .list_records(did, collection, Some(100), cursor.as_deref())
1945                .await?;
1946            refuse_malformed(&page, collection)?;
1947            let got = page.records.len();
1948            // The sidecar proxies the account's PDS, so this walk's size is as
1949            // remote-controlled as the direct client's. It carried no budget at
1950            // all until a review noticed it was the default backend.
1951            if !budget.admit(&page.records) {
1952                anyhow::bail!(
1953                    "listRecords for {collection} exceeded the {max_bytes}-byte cap \
1954                     ({} held, {} bytes charged) — refusing to accumulate further",
1955                    out.len(),
1956                    budget.used(),
1957                );
1958            }
1959            extend_bounded(&mut out, page.records, MAX_LIST_RECORDS, collection)?;
1960            match page.cursor {
1961                Some(next) if got > 0 && Some(&next) != cursor.as_ref() => {
1962                    cursor = Some(next);
1963                    more_offered = true;
1964                }
1965                _ => {
1966                    more_offered = false;
1967                    break;
1968                }
1969            }
1970        }
1971        // **Running out of pages is a refusal, not a short answer.** Falling out
1972        // of the loop used to return `Ok(out)`, so a repo bigger than the page
1973        // budget produced a truncated list indistinguishable from a complete
1974        // one — and `resolve_subscriptions` needs an `Err` for its fail-closed
1975        // branch. Given `Ok`, it hands the short list to `replace_sub_refs`,
1976        // which DELETEs the reader's whole `sub_ref` projection and reinserts
1977        // only what it was given. `extend_bounded` cannot catch this either:
1978        // `MAX_LIST_PAGES` x the 100 we request is `MAX_LIST_RECORDS`, so the
1979        // page budget runs out first.
1980        //
1981        // **The cap is on REQUESTS, so where it bites in RECORDS is the server's
1982        // choice and not ours.** We ask for 100 a page; a PDS MAY answer with
1983        // fewer, and only one that honours the limit puts the boundary anywhere
1984        // near `MAX_LIST_PAGES` x 100. Halve the page size and the same budget
1985        // reaches half as many records; a server that returns MORE than asked
1986        // trips `extend_bounded` first, which is the case the sentence above does
1987        // not cover. Said this way because an earlier version of this comment
1988        // named a fixed record window as though our own constants decided it.
1989        //
1990        // **And at the boundary the refusal is a FALSE one.** Terminating costs
1991        // one extra request, because a short page can still carry a cursor — this
1992        // project's own PDS does exactly that — so a walk that fills its last
1993        // allowed page is holding every record it was ever going to hold and
1994        // refuses anyway, on the strength of a cursor it never followed. With
1995        // `limit=100` honoured that window is a repo of roughly 19 901 to 20 000
1996        // records. The direction is safe and the alternative is deleting feeds,
1997        // but it is a false refusal and not a clean boundary.
1998        if more_offered {
1999            anyhow::bail!(
2000                "listRecords for {collection} did not finish within {MAX_LIST_PAGES} pages \
2001                 ({} held, and the PDS still offered more) — refusing a short list",
2002                out.len(),
2003            );
2004        }
2005        Ok(out)
2006    }
2007
2008    /// `create` — create a record (server-assigned rkey). Returns its strong ref.
2009    /// **Private, not `pub` — and not `pub(crate)`.** This is generic over
2010    /// `T: Serialize`, so it will happily write a raw `lexicon::Subscription`:
2011    /// the general case of the hole `create_subscriptions_batch` was one
2012    /// instance of. The vetted wrappers in this `impl` are the sanctioned entry
2013    /// points. `pub(crate)` was tried first and stops nothing that matters — a
2014    /// handler in `web.rs` is in this crate. Private is what makes the wrappers
2015    /// a fact rather than a convention, and it costs nothing: nothing outside
2016    /// this module ever called it.
2017    async fn create_record<T: Serialize>(
2018        &self,
2019        did: &str,
2020        collection: &str,
2021        record: &T,
2022    ) -> Result<WriteResult> {
2023        let body = json!({
2024            "did": did,
2025            "action": RepoAction::Create.as_str(),
2026            "collection": collection,
2027            "record": record,
2028        });
2029        let data = self.repo(body).await?;
2030        serde_json::from_value(data).context("parsing sidecar createRecord data")
2031    }
2032
2033    /// `put` — upsert a record at a known rkey. Returns its strong ref.
2034    /// **Private, not `pub` — and not `pub(crate)`.** This is generic over
2035    /// `T: Serialize`, so it will happily write a raw `lexicon::Subscription`:
2036    /// the general case of the hole `create_subscriptions_batch` was one
2037    /// instance of. The vetted wrappers in this `impl` are the sanctioned entry
2038    /// points. `pub(crate)` was tried first and stops nothing that matters — a
2039    /// handler in `web.rs` is in this crate. Private is what makes the wrappers
2040    /// a fact rather than a convention, and it costs nothing: nothing outside
2041    /// this module ever called it.
2042    async fn put_record<T: Serialize>(
2043        &self,
2044        did: &str,
2045        collection: &str,
2046        rkey: &str,
2047        record: &T,
2048    ) -> Result<WriteResult> {
2049        let body = json!({
2050            "did": did,
2051            "action": RepoAction::Put.as_str(),
2052            "collection": collection,
2053            "rkey": rkey,
2054            "record": record,
2055        });
2056        let data = self.repo(body).await?;
2057        serde_json::from_value(data).context("parsing sidecar putRecord data")
2058    }
2059
2060    /// `delete` — delete a record by collection + rkey.
2061    pub async fn delete_record(&self, did: &str, collection: &str, rkey: &str) -> Result<()> {
2062        let body = json!({
2063            "did": did,
2064            "action": RepoAction::Delete.as_str(),
2065            "collection": collection,
2066            "rkey": rkey,
2067        });
2068        // A 200 carrying an error envelope is not a delete: this reported
2069        // success while the record stayed in the reader's repo, and the UI
2070        // showed them unsubscribed from a feed they still had.
2071        self.repo(body)
2072            .await
2073            .and_then(|data| reject_error_envelope(&data))?;
2074        Ok(())
2075    }
2076
2077    /// `applyWrites` — a batch of create/update/delete ops in one round-trip.
2078    /// **Private, not `pub` — and not `pub(crate)`.** This is generic over
2079    /// `T: Serialize`, so it will happily write a raw `lexicon::Subscription`:
2080    /// the general case of the hole `create_subscriptions_batch` was one
2081    /// instance of. The vetted wrappers in this `impl` are the sanctioned entry
2082    /// points. `pub(crate)` was tried first and stops nothing that matters — a
2083    /// handler in `web.rs` is in this crate. Private is what makes the wrappers
2084    /// a fact rather than a convention, and it costs nothing: nothing outside
2085    /// this module ever called it.
2086    ///
2087    /// Sent in chunks within the PDS's limits — see [`apply_writes_chunked`]
2088    /// for what a failure part-way means. **Chunked here, not in the sidecar**,
2089    /// although the sidecar is what calls the PDS: it answers one
2090    /// `/internal/repo` request with one result, which has no way to say that
2091    /// half a batch landed, and this client is its only caller. Chunking here
2092    /// also keeps the hop itself under the sidecar's 1 MiB Fastify body limit.
2093    async fn apply_writes(&self, did: &str, writes: &[WriteOp]) -> Result<()> {
2094        apply_writes_chunked(writes, |chunk| self.apply_writes_once(did, chunk)).await
2095    }
2096
2097    /// One `/internal/repo` `applyWrites`, unchunked. Reached only through
2098    /// [`apply_writes`](Self::apply_writes).
2099    async fn apply_writes_once(&self, did: &str, writes: &[WriteOp]) -> Result<()> {
2100        let ops: Vec<Value> = writes.iter().map(WriteOp::to_sidecar_json).collect();
2101        let body = json!({
2102            "did": did,
2103            "action": RepoAction::ApplyWrites.as_str(),
2104            "writes": ops,
2105        });
2106        self.repo(body)
2107            .await
2108            .and_then(|data| reject_error_envelope(&data))?;
2109        Ok(())
2110    }
2111
2112    // -- typed lexicon wrappers (mirror the old PdsClient surface) ------------
2113
2114    /// List every [`Subscription`] record in `did`'s repo (paged fully).
2115    pub async fn list_subscriptions(&self, did: &str) -> Result<Vec<(String, Subscription)>> {
2116        self.list_typed(did, lexicon::nsid::SUBSCRIPTION).await
2117    }
2118
2119    /// Create a [`Subscription`] record (subscribe to a feed).
2120    pub async fn create_subscription(
2121        &self,
2122        did: &str,
2123        sub: &crate::vetted::VettedSubscription,
2124    ) -> Result<WriteResult> {
2125        self.create_record(did, lexicon::nsid::SUBSCRIPTION, sub)
2126            .await
2127    }
2128
2129    /// Delete a [`Subscription`] record by rkey (unsubscribe).
2130    pub async fn delete_subscription(&self, did: &str, rkey: &str) -> Result<()> {
2131        self.delete_record(did, lexicon::nsid::SUBSCRIPTION, rkey)
2132            .await
2133    }
2134
2135    /// List every [`Folder`] record in `did`'s repo.
2136    pub async fn list_folders(&self, did: &str) -> Result<Vec<(String, Folder)>> {
2137        self.list_typed(did, lexicon::nsid::FOLDER).await
2138    }
2139
2140    /// List every [`Saved`] record in `did`'s repo.
2141    pub async fn list_saved(&self, did: &str) -> Result<Vec<(String, Saved)>> {
2142        self.list_typed(did, lexicon::nsid::SAVED).await
2143    }
2144
2145    /// List every [`ReadState`] cursor in `did`'s repo (the read side a
2146    /// login-time read-state merge would consume).
2147    pub async fn list_read_states(&self, did: &str) -> Result<Vec<(String, ReadState)>> {
2148        self.list_typed(did, lexicon::nsid::READ_STATE).await
2149    }
2150
2151    /// Upsert a single [`ReadState`] cursor at its feed-derived rkey.
2152    pub async fn put_read_state(
2153        &self,
2154        did: &str,
2155        rkey: &str,
2156        state: &ReadState,
2157    ) -> Result<WriteResult> {
2158        self.put_record(did, lexicon::nsid::READ_STATE, rkey, state)
2159            .await
2160    }
2161
2162    /// Batch-flush many dirty [`ReadState`] cursors via `applyWrites` (chunked).
2163    ///
2164    /// Each `(rkey, state, pds_created)` becomes a `create` op at the feed-derived
2165    /// rkey when the record does NOT yet exist (`pds_created == false`), and an
2166    /// `update` op when it does. This is what makes the FIRST flush of a feed
2167    /// succeed: `applyWrites#update` errors on a record that does not pre-exist,
2168    /// and `applyWrites` is atomic per-repo, so a single not-yet-created cursor
2169    /// would otherwise drop the whole DID batch. Both kinds ride the SAME
2170    /// `applyWrites` batch so batching is preserved.
2171    pub async fn flush_read_states(
2172        &self,
2173        did: &str,
2174        cursors: &[(String, ReadState, bool)],
2175    ) -> Result<()> {
2176        if cursors.is_empty() {
2177            return Ok(());
2178        }
2179        let writes = read_state_write_ops(cursors)?;
2180        self.apply_writes(did, &writes).await
2181    }
2182
2183    // -- reader-facing record CRUD (the surface the web layer calls) ----------
2184    //
2185    // These are the typed convenience methods `web.rs` uses to manage a user's
2186    // feeds/folders/saved items *as records in their PDS*. They mirror the
2187    // create/list surface above but use the reader vocabulary
2188    // (add/remove/rename) and, for the `add_*` verbs, return the server-assigned
2189    // rkey so the caller can address the new record without a re-list. Ordering
2190    // is made deterministic where it matters (see [`list_subscriptions_sorted`]
2191    // etc.) so the server-rendered HTML is stable between reads.
2192
2193    // -- subscriptions -------------------------------------------------------
2194
2195    /// Add a subscription (subscribe to a feed) — `createRecord`, server-assigned
2196    /// `tid` rkey. Returns the new record's **rkey** so the web layer can offer
2197    /// unsubscribe/rename immediately.
2198    pub async fn add_subscription(
2199        &self,
2200        did: &str,
2201        sub: &crate::vetted::VettedSubscription,
2202    ) -> Result<String> {
2203        Ok(self.create_subscription(did, sub).await?.into_rkey())
2204    }
2205
2206    /// Remove a subscription (unsubscribe) by rkey — `deleteRecord`. Alias of
2207    /// [`delete_subscription`](Self::delete_subscription) in the reader vocabulary.
2208    pub async fn remove_subscription(&self, did: &str, rkey: &str) -> Result<()> {
2209        self.delete_subscription(did, rkey).await
2210    }
2211
2212    /// Update / rename a subscription in place at a known rkey — `putRecord`.
2213    ///
2214    /// The whole record is replaced (retitle, move to a folder, change the
2215    /// fetch hint …). Upsert semantics: it also creates the record if the rkey
2216    /// is somehow absent, so it is safe as a general "write this exact record".
2217    pub async fn update_subscription(
2218        &self,
2219        did: &str,
2220        rkey: &str,
2221        sub: &crate::vetted::VettedSubscription,
2222    ) -> Result<WriteResult> {
2223        self.put_record(did, lexicon::nsid::SUBSCRIPTION, rkey, sub)
2224            .await
2225    }
2226
2227    /// List every subscription, **sorted deterministically** — by display title
2228    /// (case-insensitive), then feed URL, then rkey as the final tiebreaker — so
2229    /// the rendered feed list is stable across reads regardless of PDS return
2230    /// order. Untitled feeds sort by their URL.
2231    pub async fn list_subscriptions_sorted(
2232        &self,
2233        did: &str,
2234    ) -> Result<Vec<(String, Subscription)>> {
2235        let mut subs = self.list_subscriptions(did).await?;
2236        // The comparator is SHARED with the Rust-native client so the two
2237        // cannot order the list differently across the cutover.
2238        subs.sort_by(lexicon::sort::subscriptions);
2239        Ok(subs)
2240    }
2241
2242    /// Batch-add many subscriptions via `applyWrites` (chunked) — the OPML-import path.
2243    ///
2244    /// Each feed becomes one `create` op. Client-side monotonic `tid`
2245    /// rkeys are assigned so the batch is deterministic and the imported feeds
2246    /// keep OPML order (server-assigned tids would also be monotonic, but pinning
2247    /// them here makes the whole import reproducible and testable offline).
2248    /// Returns the assigned rkeys in input order.
2249    ///
2250    /// More than [`APPLY_WRITES_MAX_OPS`] feeds is more than one call, so the
2251    /// import can part-land: on an error, [`ApplyWritesIncomplete::of`] gives
2252    /// `landed`, and the first `landed` of `subs` are in the repo.
2253    pub async fn add_subscriptions_bulk(
2254        &self,
2255        did: &str,
2256        subs: &[crate::vetted::VettedSubscription],
2257    ) -> Result<Vec<String>> {
2258        let mut gen = TidGenerator::new();
2259        let mut rkeys = Vec::with_capacity(subs.len());
2260        let mut writes = Vec::with_capacity(subs.len());
2261        for sub in subs {
2262            let rkey = gen.next();
2263            writes.push(WriteOp::Create {
2264                collection: lexicon::nsid::SUBSCRIPTION.to_string(),
2265                rkey: Some(rkey.clone()),
2266                value: serde_json::to_value(sub)?,
2267            });
2268            rkeys.push(rkey);
2269        }
2270        self.apply_writes(did, &writes).await?;
2271        Ok(rkeys)
2272    }
2273
2274    // -- folders -------------------------------------------------------------
2275
2276    /// Add a folder — `createRecord`, server-assigned `tid` rkey. Returns the
2277    /// new folder's rkey (subscriptions reference it by its `at://` URI).
2278    pub async fn add_folder(&self, did: &str, folder: &Folder) -> Result<String> {
2279        Ok(self
2280            .create_record(did, lexicon::nsid::FOLDER, folder)
2281            .await?
2282            .into_rkey())
2283    }
2284
2285    /// Remove a folder by rkey — `deleteRecord`. (Subscriptions referencing it
2286    /// are left untouched; a dangling `folder` ref reads as "unfiled".)
2287    pub async fn remove_folder(&self, did: &str, rkey: &str) -> Result<()> {
2288        self.delete_record(did, lexicon::nsid::FOLDER, rkey).await
2289    }
2290
2291    /// Rename / update a folder in place at a known rkey — `putRecord`
2292    /// (rename, or change its `position` sort hint).
2293    pub async fn rename_folder(
2294        &self,
2295        did: &str,
2296        rkey: &str,
2297        folder: &Folder,
2298    ) -> Result<WriteResult> {
2299        self.put_record(did, lexicon::nsid::FOLDER, rkey, folder)
2300            .await
2301    }
2302
2303    /// List every folder, **sorted deterministically** — by `position` (the
2304    /// lexicon's sort hint; unset sorts last), then name (case-insensitive),
2305    /// then rkey — so the sidebar order is stable.
2306    pub async fn list_folders_sorted(&self, did: &str) -> Result<Vec<(String, Folder)>> {
2307        let mut folders = self.list_folders(did).await?;
2308        folders.sort_by(lexicon::sort::folders);
2309        Ok(folders)
2310    }
2311
2312    // -- saved / starred -----------------------------------------------------
2313
2314    /// Add a saved (starred / save-for-later) entry — `createRecord`,
2315    /// server-assigned `tid` rkey. Returns the new record's rkey.
2316    pub async fn add_saved(&self, did: &str, saved: &crate::vetted::VettedSaved) -> Result<String> {
2317        Ok(self
2318            .create_record(did, lexicon::nsid::SAVED, saved)
2319            .await?
2320            .into_rkey())
2321    }
2322
2323    /// Remove a saved entry by rkey — `deleteRecord` (un-star).
2324    pub async fn remove_saved(&self, did: &str, rkey: &str) -> Result<()> {
2325        self.delete_record(did, lexicon::nsid::SAVED, rkey).await
2326    }
2327
2328    /// List every saved entry, **sorted deterministically** — newest first by
2329    /// `createdAt` (RFC-3339 sorts lexicographically), then rkey — so the
2330    /// "saved for later" list reads most-recent-first and is stable.
2331    pub async fn list_saved_sorted(&self, did: &str) -> Result<Vec<(String, Saved)>> {
2332        let mut saved = self.list_saved(did).await?;
2333        saved.sort_by(lexicon::sort::saved);
2334        Ok(saved)
2335    }
2336
2337    /// List a collection for `did` and parse each record's value into `T`,
2338    /// pairing it with its rkey. Unparseable records are skipped with a warning
2339    /// (forward-compat).
2340    async fn list_typed<T: DeserializeOwned>(
2341        &self,
2342        did: &str,
2343        collection: &str,
2344    ) -> Result<Vec<(String, T)>> {
2345        let records = self.list_all_records(did, collection).await?;
2346        let mut out = Vec::with_capacity(records.len());
2347        for rec in records {
2348            let rkey = rec.rkey().unwrap_or_default().to_string();
2349            match rec.parse::<T>() {
2350                Ok(value) => out.push((rkey, value)),
2351                Err(e) => tracing::warn!(
2352                    collection,
2353                    uri = %rec.uri,
2354                    error = %e,
2355                    "skipping unparseable record in collection"
2356                ),
2357            }
2358        }
2359        Ok(out)
2360    }
2361}
2362
2363// ---------------------------------------------------------------------------
2364// applyWrites operations
2365// ---------------------------------------------------------------------------
2366
2367/// Build the `applyWrites` ops for a batch of dirty read-state cursors.
2368///
2369/// Each `(rkey, state, pds_created)` becomes a `#create` op (at the stable
2370/// feed-derived rkey) when the PDS record does NOT yet exist, and a `#update`
2371/// when it does. This is the crux of the first-flush fix: an `#update` on a
2372/// missing record errors, and `applyWrites` is atomic per-repo, so a single
2373/// not-yet-created cursor in the batch would drop the whole DID's flush. Emitting
2374/// a `create` for those makes a feed's first flush succeed while keeping every
2375/// op in ONE batch. Shared by both the sidecar and direct-PDS flush paths.
2376pub(crate) fn read_state_write_ops(cursors: &[(String, ReadState, bool)]) -> Result<Vec<WriteOp>> {
2377    cursors
2378        .iter()
2379        .map(|(rkey, state, pds_created)| {
2380            let value = serde_json::to_value(state)?;
2381            Ok(if *pds_created {
2382                WriteOp::Update {
2383                    collection: lexicon::nsid::READ_STATE.to_string(),
2384                    rkey: rkey.clone(),
2385                    value,
2386                }
2387            } else {
2388                WriteOp::Create {
2389                    collection: lexicon::nsid::READ_STATE.to_string(),
2390                    rkey: Some(rkey.clone()),
2391                    value,
2392                }
2393            })
2394        })
2395        .collect()
2396}
2397
2398/// Most writes one `com.atproto.repo.applyWrites` call may carry.
2399///
2400/// **The limit is the reference PDS's, not the lexicon's.** The lexicon's
2401/// `writes` array has no `maxLength` (checked against
2402/// `lexicons/com/atproto/repo/applyWrites.json` on bluesky-social/atproto
2403/// `main`, and against its history back to 2024-02); the cap is enforced by the
2404/// handler, `packages/pds/src/api/com/atproto/repo/applyWrites.ts`
2405/// (`if (writes.length > 200) throw new InvalidRequestError('Too many writes.
2406/// Max: 200')`, unchanged at a0c49d9). vlpds documents the same figure for its
2407/// commit coalescing. A PDS that allows more loses nothing by being sent 200.
2408pub const APPLY_WRITES_MAX_OPS: usize = 200;
2409
2410/// Most bytes of serialized writes one `applyWrites` call may carry.
2411///
2412/// **Sized to the smallest limit a deployed PDS is known to apply, not the
2413/// largest.** The reference PDS took `applyWrites` bodies up to its server-wide
2414/// `jsonLimit` of 150 KiB (`150 * 1024` in `packages/pds/src/index.ts`) until
2415/// atproto#4989 (2026-05-21) raised the record methods to `1_000_000` bytes.
2416/// Self-hosted PDSes run older releases for months, so the 150 KiB figure is
2417/// the live one for some readers. The other limits on this path are all
2418/// larger: the newer reference PDS's 1,000,000 bytes, the sidecar hop's Fastify
2419/// default `bodyLimit` of 1 MiB, and vlpds's coalesced commit of "up to 200
2420/// operations or 1 MB of record bytes".
2421///
2422/// 128 KiB leaves 22 KiB under 150 KiB for what this does not count — the
2423/// `{"repo": …, "writes": [ ]}` envelope, around a hundred bytes — and is
2424/// measured on the PDS-shaped op (`WriteOp::to_json`), which is also what the
2425/// sidecar forwards and is larger than the sidecar's own request shape.
2426///
2427/// What it costs, measured: 200 subscriptions with a title and site URL each
2428/// are ~76 KB, so an ordinary OPML import is still one call per 200 feeds. A
2429/// read-state cursor with a full 1,000-id set is ~10 KB with the store's
2430/// integer entry ids (~20 KB with both sets full), so a flush crosses the
2431/// bound at around a dozen full cursors — rare, since a cursor is compacted at
2432/// half the cap — and a refused body is worse than an extra round trip.
2433///
2434/// A single op larger than this is still sent, alone: it cannot be split, and
2435/// the PDS is the one to judge it.
2436pub const APPLY_WRITES_MAX_BYTES: usize = 128 * 1024;
2437
2438/// Split a batch into the consecutive ranges [`apply_writes_chunked`] sends,
2439/// each within [`APPLY_WRITES_MAX_OPS`] and [`APPLY_WRITES_MAX_BYTES`].
2440///
2441/// Ranges rather than slices so a caller can say WHICH writes a chunk held.
2442/// Order is preserved, every op is in exactly one range, and no range is empty;
2443/// an empty batch yields no ranges.
2444pub(crate) fn chunk_writes(writes: &[WriteOp]) -> Vec<std::ops::Range<usize>> {
2445    let mut chunks = Vec::new();
2446    let mut start = 0;
2447    let mut bytes = 0;
2448    for (i, op) in writes.iter().enumerate() {
2449        // +1 for the comma between array elements.
2450        let size = op.to_json().to_string().len() + 1;
2451        let held = i - start;
2452        if held > 0 && (held == APPLY_WRITES_MAX_OPS || bytes + size > APPLY_WRITES_MAX_BYTES) {
2453            chunks.push(start..i);
2454            start = i;
2455            bytes = 0;
2456        }
2457        bytes += size;
2458    }
2459    if start < writes.len() {
2460        chunks.push(start..writes.len());
2461    }
2462    chunks
2463}
2464
2465/// Send `writes` as consecutive `applyWrites` calls within the PDS's limits,
2466/// stopping at the first that fails.
2467///
2468/// **Every client's `apply_writes` goes through this**, so no caller can send
2469/// an oversized call: the OAuth client ([`crate::oauth::xrpc::Repo`]), the
2470/// sidecar client ([`SidecarClient`]) and the direct client ([`PdsClient`]).
2471/// `send` is that client's single-call primitive.
2472///
2473/// **The batch is no longer atomic.** `applyWrites` is atomic per CALL, so a
2474/// split batch can half-land. The contract a caller gets instead:
2475///
2476/// * chunks go in input order, one at a time, and a failure stops the run —
2477///   so what landed is always a PREFIX of `writes`;
2478/// * on failure the error carries an [`ApplyWritesIncomplete`] saying how long
2479///   that prefix is, how many writes after it are in doubt (the failed chunk:
2480///   atomic, so all or none, but a timeout cannot say which), and that the
2481///   rest were never sent.
2482///
2483/// Stopping rather than carrying on is what keeps that a prefix: sending chunk
2484/// 3 after chunk 2 failed would leave a gap no caller could describe in one
2485/// number, and a gap in an OPML import is feeds silently missing from the
2486/// middle of the list.
2487pub(crate) async fn apply_writes_chunked<'a, F, Fut>(
2488    writes: &'a [WriteOp],
2489    mut send: F,
2490) -> Result<()>
2491where
2492    F: FnMut(&'a [WriteOp]) -> Fut,
2493    Fut: std::future::Future<Output = Result<()>>,
2494{
2495    let chunks = chunk_writes(writes);
2496    let count = chunks.len();
2497    for (i, range) in chunks.into_iter().enumerate() {
2498        let (landed, in_doubt) = (range.start, range.len());
2499        if let Err(cause) = send(&writes[range]).await {
2500            return Err(ApplyWritesIncomplete {
2501                landed,
2502                in_doubt,
2503                total: writes.len(),
2504                chunk: i + 1,
2505                chunks: count,
2506                cause,
2507            }
2508            .into());
2509        }
2510    }
2511    Ok(())
2512}
2513
2514/// How far a chunked `applyWrites` got before a chunk failed — carried by the
2515/// error from every client's `apply_writes`, and so from `flush_read_states`
2516/// and `add_subscriptions_bulk`.
2517///
2518/// Read it with [`ApplyWritesIncomplete::of`]. In terms of the caller's own
2519/// input, in order:
2520///
2521/// * `writes[..landed]` were committed (each chunk was acknowledged);
2522/// * `writes[landed..landed + in_doubt]` were the failed call — usually not
2523///   committed, but a timeout or a lost response cannot rule it out;
2524/// * everything after was never sent.
2525///
2526/// Display is the underlying failure's message, with the progress appended
2527/// when the batch had more than one chunk, so a log line still names the
2528/// PDS's reason. The failure's own causes stay reachable through
2529/// [`anyhow::Error::chain`].
2530#[derive(Debug)]
2531pub struct ApplyWritesIncomplete {
2532    /// Writes committed, counted from the start of the input.
2533    pub landed: usize,
2534    /// Writes in the failed call, starting at `landed`.
2535    pub in_doubt: usize,
2536    /// Writes in the whole batch.
2537    pub total: usize,
2538    chunk: usize,
2539    chunks: usize,
2540    cause: anyhow::Error,
2541}
2542
2543impl ApplyWritesIncomplete {
2544    /// The progress record an `apply_writes` error carries, if it has one.
2545    pub fn of(err: &anyhow::Error) -> Option<&Self> {
2546        err.chain().find_map(|e| e.downcast_ref::<Self>())
2547    }
2548
2549    /// The failure that stopped the run, as the client reported it.
2550    pub fn cause(&self) -> &anyhow::Error {
2551        &self.cause
2552    }
2553}
2554
2555impl std::fmt::Display for ApplyWritesIncomplete {
2556    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2557        write!(f, "{}", self.cause)?;
2558        if self.chunks > 1 {
2559            write!(
2560                f,
2561                " (applyWrites call {} of {}; {} of {} writes had landed)",
2562                self.chunk, self.chunks, self.landed, self.total
2563            )?;
2564        }
2565        Ok(())
2566    }
2567}
2568
2569impl std::error::Error for ApplyWritesIncomplete {
2570    // The cause's own Display is already in ours, so the chain continues from
2571    // ITS source — `{:#}` would otherwise print the PDS's message twice.
2572    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
2573        self.cause.chain().nth(1)
2574    }
2575}
2576
2577/// One operation in a `PdsClient::apply_writes` batch.
2578///
2579/// Maps to the `com.atproto.repo.applyWrites` union of
2580/// `#create` / `#update` / `#delete`.
2581#[derive(Debug, Clone)]
2582pub enum WriteOp {
2583    /// Create a record (server-assigned rkey unless `rkey` is given).
2584    Create {
2585        /// The collection NSID.
2586        collection: String,
2587        /// Optional explicit rkey (`None` → server assigns a tid).
2588        rkey: Option<String>,
2589        /// The record body.
2590        value: Value,
2591    },
2592    /// Upsert a record at a known rkey (the read-state cursor case).
2593    Update {
2594        /// The collection NSID.
2595        collection: String,
2596        /// The rkey to write at.
2597        rkey: String,
2598        /// The record body.
2599        value: Value,
2600    },
2601    /// Delete a record by collection + rkey.
2602    Delete {
2603        /// The collection NSID.
2604        collection: String,
2605        /// The rkey to delete.
2606        rkey: String,
2607    },
2608}
2609
2610impl WriteOp {
2611    /// Render this op as the tagged JSON `com.atproto.repo.applyWrites` expects.
2612    ///
2613    /// `pub(crate)` so [`crate::oauth::xrpc`] can build the same batch body.
2614    /// Sharing the rendering rather than reimplementing it is what keeps the two
2615    /// clients wire-identical across the cutover.
2616    pub(crate) fn to_json(&self) -> Value {
2617        match self {
2618            WriteOp::Create {
2619                collection,
2620                rkey,
2621                value,
2622            } => {
2623                let mut op = json!({
2624                    "$type": "com.atproto.repo.applyWrites#create",
2625                    "collection": collection,
2626                    "value": value,
2627                });
2628                if let Some(rkey) = rkey {
2629                    op["rkey"] = json!(rkey);
2630                }
2631                op
2632            }
2633            WriteOp::Update {
2634                collection,
2635                rkey,
2636                value,
2637            } => json!({
2638                "$type": "com.atproto.repo.applyWrites#update",
2639                "collection": collection,
2640                "rkey": rkey,
2641                "value": value,
2642            }),
2643            WriteOp::Delete { collection, rkey } => json!({
2644                "$type": "com.atproto.repo.applyWrites#delete",
2645                "collection": collection,
2646                "rkey": rkey,
2647            }),
2648        }
2649    }
2650
2651    /// Render this op in the shape the OAuth sidecar's `/internal/repo`
2652    /// `applyWrites` expects: `{action, collection, rkey?, value?}` (the sidecar
2653    /// maps `action` → the `com.atproto.repo.applyWrites#<kind>` union member).
2654    fn to_sidecar_json(&self) -> Value {
2655        match self {
2656            WriteOp::Create {
2657                collection,
2658                rkey,
2659                value,
2660            } => {
2661                let mut op = json!({
2662                    "action": "create",
2663                    "collection": collection,
2664                    "value": value,
2665                });
2666                if let Some(rkey) = rkey {
2667                    op["rkey"] = json!(rkey);
2668                }
2669                op
2670            }
2671            WriteOp::Update {
2672                collection,
2673                rkey,
2674                value,
2675            } => json!({
2676                "action": "update",
2677                "collection": collection,
2678                "rkey": rkey,
2679                "value": value,
2680            }),
2681            WriteOp::Delete { collection, rkey } => json!({
2682                "action": "delete",
2683                "collection": collection,
2684                "rkey": rkey,
2685            }),
2686        }
2687    }
2688}
2689
2690// ---------------------------------------------------------------------------
2691// TID rkeys (client-assigned, sortable, deterministic within a batch)
2692// ---------------------------------------------------------------------------
2693
2694/// The atproto base32-sortable alphabet (`s32`) — the digits/letters, minus the
2695/// ambiguous set, in **ascending** order so a bytewise string compare of two
2696/// TIDs matches their timestamp order.
2697const S32_ALPHABET: &[u8; 32] = b"234567abcdefghijklmnopqrstuvwxyz";
2698
2699/// A monotonic generator of atproto **TID** record keys.
2700///
2701/// A TID is a 13-char `s32`-encoded 64-bit integer: a 53-bit microsecond
2702/// timestamp in the high bits and a 10-bit "clock id" in the low bits (the top
2703/// bit is always 0). Encoded in the ascending `s32` alphabet, TIDs sort
2704/// lexicographically in creation order — which is exactly what we want for a
2705/// batched OPML import: assigning the rkeys ourselves keeps the imported feeds
2706/// in input order and makes [`add_subscriptions_bulk`](SidecarClient::add_subscriptions_bulk)
2707/// fully reproducible/testable without a live PDS.
2708///
2709/// Monotonicity within one generator is guaranteed by tracking the last value
2710/// and bumping to `last + 1` if the clock hasn't advanced — so a burst of
2711/// same-microsecond calls still yields strictly increasing, ordered rkeys.
2712pub(crate) struct TidGenerator {
2713    /// The last raw 64-bit TID value emitted (0 = none yet).
2714    last: u64,
2715    /// The low-10-bit clock id, randomized once per generator to avoid
2716    /// cross-instance collisions on the same microsecond.
2717    clock_id: u64,
2718}
2719
2720impl TidGenerator {
2721    /// A fresh generator with a per-instance clock id derived from the current
2722    /// nanosecond clock (no extra deps; uniqueness only needs to hold within a
2723    /// single import batch, and the timestamp bits carry the ordering).
2724    pub(crate) fn new() -> Self {
2725        let nanos = std::time::SystemTime::now()
2726            .duration_since(std::time::UNIX_EPOCH)
2727            .map(|d| d.subsec_nanos() as u64)
2728            .unwrap_or(0);
2729        Self {
2730            last: 0,
2731            clock_id: nanos & 0x3ff,
2732        }
2733    }
2734
2735    /// The next monotonic TID rkey (13 `s32` chars).
2736    pub(crate) fn next(&mut self) -> String {
2737        let micros = std::time::SystemTime::now()
2738            .duration_since(std::time::UNIX_EPOCH)
2739            .map(|d| d.as_micros() as u64)
2740            .unwrap_or(0);
2741        // Timestamp in bits 63..10 (top bit stays 0), clock id in bits 9..0.
2742        let mut raw = ((micros & 0x001f_ffff_ffff_ffff) << 10) | self.clock_id;
2743        if raw <= self.last {
2744            raw = self.last + 1;
2745        }
2746        self.last = raw;
2747        encode_s32_tid(raw)
2748    }
2749}
2750
2751/// Encode a 64-bit TID value as a 13-char big-endian `s32` string.
2752fn encode_s32_tid(mut v: u64) -> String {
2753    let mut buf = [0u8; 13];
2754    for slot in buf.iter_mut().rev() {
2755        *slot = S32_ALPHABET[(v & 0x1f) as usize];
2756        v >>= 5;
2757    }
2758    // 13 * 5 = 65 bits cover the 64-bit value; the leading char carries bits
2759    // 64..60, and bit 64 does not exist in a `u64` while bit 63 is always 0 in
2760    // a real TID, so the leading char is always one of the alphabet's first
2761    // eight symbols. Between 2005-09-05 and 2041-05-10 it is the second one,
2762    // which is why real TIDs all begin with `3`.
2763    String::from_utf8(buf.to_vec()).unwrap_or_default()
2764}
2765
2766/// The earliest instant a real TID can encode: 2020-01-01T00:00:00Z, in
2767/// microseconds.
2768///
2769/// atproto did not exist before this, so a "TID" decoding to earlier is a record
2770/// key that merely *looks* like one.
2771///
2772/// **This bound catches only the slugs that fall outside the window, and that
2773/// is a minority of them.** 13 lowercase alphanumerics is an ordinary slug
2774/// shape and also a valid `s32` value, and one beginning `3` decodes into the
2775/// last few years as readily as a real record key does: `3hoursinparis` reads
2776/// as 2020-11-24, `3ideasforjune` as 2021-08-12. Nothing in the string
2777/// distinguishes them — telling a slug from a TID would mean asking the PDS
2778/// when the record was written, which the listing does not report.
2779///
2780/// What the window does buy is that a mis-read date is always an ordinary past
2781/// instant rather than an unsweepable future one. That is worth having and it
2782/// is *not* harmless: a slug reading as 2020 is older than any realistic
2783/// retention window, so the row is swept, re-listed on the next poll, and
2784/// arrives unread again — the cycle this dating work narrows but does not
2785/// close. Refusing to insert what is already past the floor is what closes it,
2786/// for a mis-read slug and a genuine archive alike, and that belongs with the
2787/// retention floor rather than here.
2788const TID_FLOOR_MICROS: i64 = 1_577_836_800_000_000;
2789
2790/// How far ahead of our own clock a timestamp someone else authored may be and
2791/// still be believed.
2792///
2793/// A PDS a second or two fast would otherwise leave a brand-new document
2794/// undated until the following poll, and an undated row is the least visible
2795/// one in the reading list. Well under any interval that matters to retention
2796/// or the per-feed cap.
2797///
2798/// **Both date sources use it.** It began as a TID-only allowance, which left a
2799/// stated `publishedAt` judged against a bare `now` while the record key two
2800/// lines below got five minutes — the same clock, two different answers, for no
2801/// reason either comment could give.
2802pub(crate) const CLOCK_SKEW_GRACE_SECS: i64 = 300;
2803
2804/// Decode a 13-char `s32` TID rkey back to its raw 64-bit value.
2805///
2806/// The exact inverse of [`encode_s32_tid`] over the values a TID can hold.
2807///
2808/// `None` for anything that is not a 13-character `s32` value: wrong length, a
2809/// character outside the alphabet, or a value whose top bit is set. That last
2810/// rejection is stricter than the TID syntax regex, which admits leading `c`
2811/// through `j`; the spec's separate rule that the high bit is always 0 is the
2812/// one enforced here, and it keeps every decoded value inside the range
2813/// [`tid_timestamp`] can shift without loss.
2814///
2815/// **This does not decide whether the string is a TID**, only whether it is a
2816/// number. Thirteen lowercase alphanumerics is also an ordinary slug, and a
2817/// slug decodes as readily as a record key does. Refusing an implausible
2818/// instant is [`tid_timestamp`]'s job, and it is where that case is caught.
2819pub(crate) fn decode_s32_tid(rkey: &str) -> Option<u64> {
2820    if rkey.len() != 13 {
2821        return None;
2822    }
2823    let mut v: u64 = 0;
2824    for b in rkey.bytes() {
2825        let digit = S32_ALPHABET.iter().position(|c| *c == b)? as u64;
2826        // `checked_*` rather than shifting: 13 chars carry 65 bits, so the
2827        // largest 13-char string overflows a `u64` and must read as "not a
2828        // TID" instead of wrapping to a plausible-looking value.
2829        v = v.checked_mul(32)?.checked_add(digit)?;
2830    }
2831    (v >> 63 == 0).then_some(v)
2832}
2833
2834/// The instant a TID rkey encodes, or `None` if the rkey is not a plausible
2835/// TID.
2836///
2837/// **Bounded at both ends on purpose.** A TID's timestamp is minted from the
2838/// writer's clock, so one decoding far into the future is either a broken clock
2839/// or a slug that happens to be 13 `s32` characters; one decoding to before
2840/// [`TID_FLOOR_MICROS`] predates atproto. Neither is a date worth trusting, and
2841/// the caller's fallback for "no date" is safer than a wrong one.
2842///
2843/// The bounds are not a slug detector — see [`TID_FLOOR_MICROS`] for why they
2844/// cannot be, and for what they do guarantee instead.
2845pub(crate) fn tid_timestamp(rkey: &str) -> Option<chrono::DateTime<chrono::Utc>> {
2846    // The low 10 bits are the clock id; the rest is microseconds since the
2847    // epoch, and clearing bit 63 above bounds it well inside `i64`.
2848    let micros = i64::try_from(decode_s32_tid(rkey)? >> 10).ok()?;
2849    if micros < TID_FLOOR_MICROS {
2850        return None;
2851    }
2852    let at = chrono::DateTime::from_timestamp_micros(micros)?;
2853    let ceiling = chrono::Utc::now() + chrono::Duration::seconds(CLOCK_SKEW_GRACE_SECS);
2854    (at <= ceiling).then_some(at)
2855}
2856
2857// ---------------------------------------------------------------------------
2858// XRPC error helper
2859// ---------------------------------------------------------------------------
2860
2861/// Minimal percent-encoding for a query-string component.
2862///
2863/// Encodes everything outside the RFC 3986 unreserved set, which covers the
2864/// values FeatherReader passes (DIDs like `did:plc:…`, NSIDs, opaque cursors,
2865/// handles) without pulling in the optional reqwest `url`/`query` feature.
2866///
2867/// `pub(crate)` so [`crate::network`] builds its relay query strings the same
2868/// way rather than keeping a second copy of the escape table.
2869pub(crate) fn urlencode(s: &str) -> String {
2870    let mut out = String::with_capacity(s.len());
2871    for b in s.bytes() {
2872        match b {
2873            b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'_' | b'.' | b'~' => {
2874                out.push(b as char)
2875            }
2876            _ => out.push_str(&format!("%{b:02X}")),
2877        }
2878    }
2879    out
2880}
2881
2882/// The most nodes a `listRecords` body may ask us to build.
2883///
2884/// **A bound on the parse, checked before the parse.** Every other limit here is
2885/// consulted after `serde_json` has already materialised the page, which cannot
2886/// prevent the allocation it exists to prevent: one 8 MB response of `{"":0}`
2887/// objects was measured retaining 824 MB, a 98x wire-to-heap amplification, on a
2888/// 512 MB box. A record cap does not see it — the page holds one record. A page
2889/// cap does not see it — there is one request. A byte budget does not see it
2890/// until the memory is already spent.
2891///
2892/// **640 000, measured — not two million, which this file's own arithmetic
2893/// already contradicted.** An earlier version reasoned "32 bytes a node plus
2894/// slack, so two million is about 128 MB". That is the model `json_bytes` two
2895/// hundred lines above explicitly rejects: it charges `MAP_NODE` + `MAP_ENTRY` +
2896/// `SLOT` = 680 bytes for a single-entry object, which costs three counted
2897/// characters. Measured against a counting allocator, the worst shape reaches
2898/// **210 bytes per counted character**, so two million admitted **400 MB** — a
2899/// bound that let through more than the attack it was written to stop, and that
2900/// the walk's own byte budget then refused a step later.
2901///
2902/// 640 000 x 210 B is about 128 MiB, which is the figure the walk budget uses
2903/// and the one this claims.
2904///
2905/// **Re-measured when a review put the worst shape at 221 B; it does not
2906/// reproduce.** Sweeping nesting depths 10, 50, 100 and 120 against the same
2907/// counting allocator, the worst is 210.6 B per counted character (at depth 120)
2908/// and peak equals retained — `serde_json` overshoots by nothing measurable
2909/// while it builds. Depth cannot be pushed further to raise the ratio, either:
2910/// `serde_json`'s own recursion limit of 128 refuses a deeper body outright,
2911/// before this guard would even matter.
2912///
2913/// **The floor is real traffic, not comfort.** The densest legitimate page is a
2914/// full `readState` listing — 100 records each carrying two arrays of
2915/// [`crate::lexicon::ReadState::MAX_IDS`] ids — which counts 403 003. So the cap
2916/// sits above the densest page the lexicons permit, and a test holds it there.
2917///
2918/// Ordinary traffic is nowhere near either number: a page of 100 real-sized
2919/// standard.site documents (seven fields, a 15 kB `textContent`, 1.5 MB on the
2920/// wire) counts **4 003**. An earlier version of this line said 40 000, which was
2921/// wrong by an order of magnitude in the direction that makes the cap look tighter
2922/// than it is; the test that was supposed to hold it served a single record and
2923/// would have passed with the cap set to 1 000.
2924pub(crate) const MAX_LIST_STRUCTURAL_CHARS: usize = 640_000;
2925
2926/// How many nodes `body` would parse into, to within one, without parsing it.
2927///
2928/// **A lower bound, despite what an earlier name said.** `[1,2,3]` counts three
2929/// — one `[` and two commas — and builds four values. The deficit is never more
2930/// than one (verified exhaustively over every body of length 1-5 from a JSON
2931/// alphabet), because every node but the outermost is introduced by one of the
2932/// characters counted here. That is the direction a guard needs: it can
2933/// under-count by one and still refuse everything it must.
2934///
2935/// Counts the structural characters that introduce a value — `{`, `[`, `,`, `:`
2936/// — **outside strings**, which is what makes this sound: a node cannot appear
2937/// without one, and a string's contents cannot invent one. Skipping strings is
2938/// the whole difficulty; counting naively would refuse a legitimate article that
2939/// happens to contain a million commas.
2940pub(crate) fn count_structural_chars(body: &[u8]) -> usize {
2941    let mut nodes = 0usize;
2942    let mut in_string = false;
2943    let mut escaped = false;
2944    for &b in body {
2945        if in_string {
2946            // `\"` stays inside the string; `\\` does not escape the quote that
2947            // follows it. Getting this pair wrong makes the scan count a whole
2948            // document as structure, or none of it.
2949            if escaped {
2950                escaped = false;
2951            } else if b == b'\\' {
2952                escaped = true;
2953            } else if b == b'"' {
2954                in_string = false;
2955            }
2956            continue;
2957        }
2958        match b {
2959            // A string is a node, and everything inside it is not.
2960            b'"' => {
2961                in_string = true;
2962                nodes += 1;
2963            }
2964            b'{' | b'[' | b',' | b':' => nodes += 1,
2965            _ => {}
2966        }
2967    }
2968    nodes
2969}
2970
2971/// Refuse a body carrying more JSON structure than
2972/// [`MAX_LIST_STRUCTURAL_CHARS`].
2973pub(crate) fn refuse_a_structure_explosion(body: &[u8], what: &str) -> Result<()> {
2974    let counted = count_structural_chars(body);
2975    anyhow::ensure!(
2976        counted <= MAX_LIST_STRUCTURAL_CHARS,
2977        "{what} counts at least {counted} structural characters, over the \
2978         {MAX_LIST_STRUCTURAL_CHARS} cap — refusing before parsing it"
2979    );
2980    Ok(())
2981}
2982
2983/// Parse a `listRecords` body, refusing an error envelope that arrived on a 2xx.
2984///
2985/// Some PDS implementations answer 200 for application failures, and the status
2986/// check in the caller cannot see those. Without the guard, `{"error","message"}`
2987/// deserialises as a page with no records — so a walk over a stranger's
2988/// collection returns a healthy, empty result in place of an error, and for the
2989/// walk that feeds `replace_sub_refs` that is revoked access rather than an empty
2990/// repo. The guard itself lives in [`page_from_body`], which every client shares.
2991pub(crate) fn parse_list_records(body: &[u8]) -> Result<ListRecordsResponse> {
2992    // **An empty body is the "unexpected body" case, not a parse error.** Reading
2993    // bytes reaches it as "EOF while parsing", where the OAuth client used to
2994    // reach it as "no records field" (its `send` mapped an empty 2xx to
2995    // `Value::Null`) and the direct client reached it as "EOF" too. Refused
2996    // either way, so this is a unification rather than a preservation — nothing
2997    // outside the tests matches on the text, and `resolve_subscriptions` fails
2998    // closed on any `Err`. It is for whoever reads the log.
2999    if body.is_empty() {
3000        anyhow::bail!("listRecords returned no records field (empty or unexpected body)");
3001    }
3002    refuse_a_structure_explosion(body, "the listRecords body")?;
3003    let parsed: ListRecordsBody =
3004        serde_json::from_slice(body).context("parsing listRecords response")?;
3005    let mut page = page_from_body(parsed)?;
3006    page.wire_bytes = body.len();
3007    Ok(page)
3008}
3009
3010/// Apply both invariants to an already-deserialised body.
3011///
3012/// **The one place the guards live, for all three clients.** They were added a
3013/// client at a time twice over, which is the whole reason a shared function
3014/// exists; splitting the sidecar onto a different route would have started that
3015/// again, so it deserialises into this same struct.
3016fn page_from_body(parsed: ListRecordsBody) -> Result<ListRecordsResponse> {
3017    if let Some(error) = parsed.error.as_ref().and_then(envelope_error_name) {
3018        let message = parsed
3019            .message
3020            .as_ref()
3021            .and_then(Value::as_str)
3022            .map(|m| format!(" — {}", truncate_for_message(m)))
3023            .unwrap_or_default();
3024        anyhow::bail!("PDS answered 2xx with an error envelope: {error}{message}");
3025    }
3026    let entries = parsed.records.ok_or_else(|| {
3027        anyhow::anyhow!("listRecords returned no records field (empty or unexpected body)")
3028    })?;
3029    let mut records = Vec::with_capacity(entries.len());
3030    let mut malformed = 0;
3031    for entry in entries {
3032        match entry {
3033            MaybeRecord::Record(r) => records.push(r),
3034            MaybeRecord::Malformed(_) => malformed += 1,
3035        }
3036    }
3037    Ok(ListRecordsResponse {
3038        records,
3039        cursor: parsed.cursor,
3040        malformed,
3041        wire_bytes: 0,
3042    })
3043}
3044
3045/// The wire shape of a `listRecords` body, read in **one** pass.
3046///
3047/// **Parsing to `Value` and then into the struct materialises the page twice.**
3048/// `serde_json::from_value` rebuilds rather than moves, so an 8 MB response was
3049/// measured holding both copies at once — a peak of roughly double the retained
3050/// size, reached before any accounting the caller does, which is why no budget
3051/// charged after the parse can cover it.
3052///
3053/// The two invariants that used to live on a `Value` are
3054/// expressed here as fields instead of lookups, and mean exactly what they did:
3055/// an `error` present on a 2xx is a failure, not an empty page, and `records`
3056/// ABSENT is not `records` empty.
3057#[derive(Debug, Default)]
3058struct ListRecordsBody {
3059    /// A `Value`, not a `String`. Typing it as a string made
3060    /// `{"error":404,"records":[]}` fail as "invalid type: integer" rather than
3061    /// as an envelope — the wrong reason for the exact shape the guard exists
3062    /// for, and the guard's whole point is that this distinction is load-bearing.
3063    error: Option<Value>,
3064    /// Likewise, and for a duller reason: `message` carries no security role,
3065    /// and typing it as a string made a PDS that stamps a non-string one onto an
3066    /// otherwise good page of a thousand records fail the entire listing.
3067    message: Option<Value>,
3068    /// `None` means the field was absent — what a proxy makes of an empty or
3069    /// unexpected upstream body. `Some(vec![])` is a genuine empty page.
3070    records: Option<Vec<MaybeRecord>>,
3071    cursor: Option<String>,
3072}
3073
3074/// One element of a `listRecords` page: a record, or something that is not one.
3075///
3076/// **Parsed per record, so one malformed envelope costs that record and not
3077/// the page** (#177). `RecordEntry.uri` is required, and the page used to be
3078/// parsed in one `from_value`, so a single `{"cid":…,"value":{}}` failed every
3079/// record beside it. Whoever reads the page decides what a skipped record
3080/// means: a stranger's publication skips it, a reader's own repo refuses.
3081#[derive(Debug, Deserialize)]
3082#[serde(untagged)]
3083enum MaybeRecord {
3084    Record(RecordEntry),
3085    Malformed(serde::de::IgnoredAny),
3086}
3087
3088/// Refuse a page that skipped records, for a walk that must not drop any.
3089///
3090/// Every walk whose result reaches `store::replace_sub_refs` calls this: a
3091/// record left out there is a subscription silently removed.
3092/// What a walk does with a record whose envelope is malformed (#177).
3093#[derive(Debug, Clone, Copy, PartialEq, Eq)]
3094enum OnMalformed {
3095    /// Refuse the walk with [`MalformedRecords`]: the result reaches
3096    /// `replace_sub_refs`, where a skipped record is a dropped subscription.
3097    Refuse,
3098    /// Skip and count it: a stranger's repo, where one bad record must not
3099    /// stall everything beside it.
3100    Skip,
3101}
3102
3103pub(crate) fn refuse_malformed(page: &ListRecordsResponse, collection: &str) -> Result<()> {
3104    if page.malformed > 0 {
3105        return Err(MalformedRecords {
3106            collection: collection.to_string(),
3107            count: page.malformed,
3108        }
3109        .into());
3110    }
3111    Ok(())
3112}
3113
3114/// **Hand-written, because the derive accepts a listing that is not an object.**
3115///
3116/// serde's derived `Deserialize` takes a struct in POSITIONAL form as well as
3117/// map form, so with every field defaulted the fourteen bytes `[null,null,[]]`
3118/// bound `records` to an empty vector and read as a healthy page — on all three
3119/// clients, and the `Value` route this replaced refused it, because
3120/// `Value::Array::get("records")` is always `None`. Neither the envelope guard
3121/// nor the duplicated-key refusal can fire on a body with no keys at all, so one
3122/// short array defeated every protection here at once and reached
3123/// `replace_sub_refs`, which deletes the reader's whole subscription projection.
3124///
3125/// **Unknown fields are read, not skipped.** `IgnoredAny` does not validate what
3126/// it skips, so `{"records":[],"x":"<invalid utf-8>"}` — not valid JSON at all —
3127/// also read as a healthy empty page where the `Value` route refused it. Reading
3128/// the value into a `Value` and dropping it costs an allocation on a field nobody
3129/// wants, and buys back the validation.
3130impl<'de> Deserialize<'de> for ListRecordsBody {
3131    fn deserialize<D: serde::Deserializer<'de>>(d: D) -> std::result::Result<Self, D::Error> {
3132        struct AsMap;
3133        impl<'de> serde::de::Visitor<'de> for AsMap {
3134            type Value = ListRecordsBody;
3135            fn expecting(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
3136                f.write_str("a listRecords object")
3137            }
3138            fn visit_map<M: serde::de::MapAccess<'de>>(
3139                self,
3140                mut map: M,
3141            ) -> std::result::Result<ListRecordsBody, M::Error> {
3142                use serde::de::Error;
3143                let mut out = ListRecordsBody::default();
3144                let (mut error, mut message, mut records, mut cursor) =
3145                    (false, false, false, false);
3146                while let Some(key) = map.next_key::<String>()? {
3147                    let seen = match key.as_str() {
3148                        "error" => std::mem::replace(&mut error, true),
3149                        "message" => std::mem::replace(&mut message, true),
3150                        "records" => std::mem::replace(&mut records, true),
3151                        "cursor" => std::mem::replace(&mut cursor, true),
3152                        _ => false,
3153                    };
3154                    if seen {
3155                        // A repeated key is last-wins in a `Value`, which is how
3156                        // a smuggled second, empty `records` array read as a
3157                        // successful page. Refused here.
3158                        return Err(M::Error::duplicate_field(match key.as_str() {
3159                            "error" => "error",
3160                            "message" => "message",
3161                            "records" => "records",
3162                            _ => "cursor",
3163                        }));
3164                    }
3165                    match key.as_str() {
3166                        "error" => out.error = Some(map.next_value()?),
3167                        "message" => out.message = Some(map.next_value()?),
3168                        "records" => out.records = Some(map.next_value()?),
3169                        "cursor" => out.cursor = map.next_value()?,
3170                        _ => {
3171                            let _validated: Value = map.next_value()?;
3172                        }
3173                    }
3174                }
3175                Ok(out)
3176            }
3177        }
3178        d.deserialize_map(AsMap)
3179    }
3180}
3181
3182/// Refuse an atproto error envelope that arrived on a 2xx.
3183///
3184/// **No listing reaches this any more — [`page_from_body`] is the one every
3185/// client shares.** It survives for the WRITE paths, where the response is still
3186/// a `Value`: `deleteRecord` and `applyWrites` on both live clients.
3187///
3188/// The history is worth keeping, because it is why a shared function exists at
3189/// all. Each client used to take `records` off the JSON its own way — the live
3190/// one with `unwrap_or(Array([]))`, the sidecar through a defaulted `Value` — and
3191/// each turned `200 {"error": …}` into `Ok(empty)`. That is not the fail-closed
3192/// branch in `web::resolve_subscriptions`: `sync_sub_refs` wrote the empty set
3193/// and `replace_sub_refs` DELETEd the DID's entire `sub_ref` projection. One bad
3194/// response revoked a reader's access to every feed they had. The guard was added
3195/// to one client at a time, twice, which is the drift a single function prevents
3196/// — and why the listing guard now lives in exactly one place rather than here.
3197pub(crate) fn reject_error_envelope(value: &Value) -> Result<()> {
3198    let Some(error) = value.get("error").and_then(envelope_error_name) else {
3199        return Ok(());
3200    };
3201    let message = value
3202        .get("message")
3203        .and_then(Value::as_str)
3204        .map(|m| format!(" — {}", truncate_for_message(m)))
3205        .unwrap_or_default();
3206    anyhow::bail!("PDS answered 2xx with an error envelope: {error}{message}")
3207}
3208
3209/// The name in an `error` field, or `None` when the field does not denote one.
3210///
3211/// **Whatever its type.** Keying on `as_str` meant a PDS answering
3212/// `{"error":404,"records":[]}` — or `{}`, or `[]` — passed the guard and read as
3213/// a healthy empty page, the shape that makes `replace_sub_refs` delete every
3214/// `sub_ref` a reader has. A non-string `error` is not a well-formed envelope,
3215/// but it is certainly not a successful listing either.
3216///
3217/// **Except the four spellings of "no error".** Absent and `null` are what an
3218/// ordinary listing carries; `false` and `0` are a convention proxies use, and
3219/// treating those as envelopes turns a good page of a thousand records into a
3220/// hard refusal, which on these walks means the reader's sidebar degrades to a
3221/// stale projection on every request.
3222///
3223/// **Bounded.** The name reaches a `warn!` that also logs the DID, and the value
3224/// is attacker-chosen: a PDS answering with hundreds of kilobytes under `error`
3225/// would otherwise put all of it in the log and allocate another copy, in code
3226/// whose purpose is cutting peak allocation.
3227fn envelope_error_name(error: &Value) -> Option<String> {
3228    match error {
3229        Value::Null | Value::Bool(false) => None,
3230        // Integer zero only. `as_f64() == Some(0.0)` also matched `-0`, `0.0`
3231        // and anything that underflows, so `1e-400` was "no error".
3232        Value::Number(n) if n.as_i64() == Some(0) || n.as_u64() == Some(0) => None,
3233
3234        Value::String(s) => Some(truncate_for_message(s)),
3235        // **The type, not the value.** `to_string()` would serialise the whole
3236        // attacker-chosen subtree before truncating it, allocating a full extra
3237        // copy of up to the body cap — in code whose purpose is cutting peak
3238        // allocation. A non-string `error` is malformed, so its contents tell a
3239        // reader nothing its shape does not.
3240        Value::Bool(_) => Some("<non-string error: bool>".to_string()),
3241        Value::Number(_) => Some("<non-string error: number>".to_string()),
3242        Value::Array(_) => Some("<non-string error: array>".to_string()),
3243        Value::Object(_) => Some("<non-string error: object>".to_string()),
3244    }
3245}
3246
3247/// Cap a string destined for an error message at a readable length.
3248fn truncate_for_message(s: &str) -> String {
3249    /// **Bytes, not characters.** A log line is bytes, and counting characters
3250    /// let astral-plane code points render four times the intended bound.
3251    const MAX_BYTES: usize = 120;
3252    if s.len() <= MAX_BYTES {
3253        return s.to_string();
3254    }
3255    let cut = s
3256        .char_indices()
3257        .map(|(i, _)| i)
3258        .take_while(|i| *i <= MAX_BYTES)
3259        .last()
3260        .unwrap_or(0);
3261    format!("{}… ({} bytes)", &s[..cut], s.len())
3262}
3263
3264/// The atproto XRPC error envelope body: `{"error": "...", "message": "..."}`.
3265#[derive(Debug, Deserialize)]
3266struct XrpcErrorBody {
3267    #[serde(default)]
3268    error: Option<String>,
3269    #[serde(default)]
3270    message: Option<String>,
3271}
3272
3273/// Consume a non-2xx response into a typed [`AtProtoError::Xrpc`], parsing the
3274/// atproto error envelope when present (falling back to `"Unknown"`).
3275///
3276/// The body is read through [`crate::net::read_capped`], **not** `resp.json()`.
3277/// Every guarded call caps its success body; routing the error body through
3278/// `resp.json()` would have left a hole exactly where the hostile-PDS threat
3279/// model points — reqwest decompresses gzip before deserialising, so a `400`
3280/// carrying a decompression bomb was an unbounded allocation on a 512 MB box.
3281/// A body we cannot read (over-cap, transport error) degrades to `"Unknown"`,
3282/// which is the same fallback an unparseable envelope already took.
3283async fn xrpc_error_from(resp: reqwest::Response) -> AtProtoError {
3284    let status = resp.status();
3285    let (error, message) = match crate::net::read_capped(resp).await {
3286        Ok(raw) => match serde_json::from_slice::<XrpcErrorBody>(&raw) {
3287            Ok(body) => (
3288                body.error.unwrap_or_else(|| "Unknown".to_string()),
3289                body.message,
3290            ),
3291            Err(_) => ("Unknown".to_string(), None),
3292        },
3293        Err(_) => ("Unknown".to_string(), None),
3294    };
3295    AtProtoError::Xrpc {
3296        status,
3297        error,
3298        message,
3299    }
3300}
3301
3302// ---------------------------------------------------------------------------
3303// Tests — record (de)serialization against a repo listRecords response shape.
3304// No network.
3305// ---------------------------------------------------------------------------
3306
3307#[cfg(test)]
3308pub(crate) mod tests {
3309    use super::*;
3310
3311    /// **Regression (v0.2.8 review).** Every guarded call caps its *success*
3312    /// body via `read_capped`, but the non-2xx branch went through
3313    /// `resp.json::<XrpcErrorBody>()` — unbounded, and with reqwest's gzip
3314    /// decompression in front of it. That left a hole precisely where the
3315    /// module's own threat model points: a hostile or DNS-rebound PDS answers
3316    /// `400` with a decompression bomb and gets an unbounded allocation on a
3317    /// 512 MB box. Both this PR's review passes checked the success path and
3318    /// walked past the error path, so the cap is asserted here explicitly.
3319    ///
3320    /// Fetched directly rather than through the guard, which rightly refuses
3321    /// loopback — the same reason `net::tests::read_capped_rejects_over_cap_body`
3322    /// bypasses it. The stub answers 200; `xrpc_error_from` reads the status only
3323    /// to record it, so the body handling under test is identical.
3324    #[tokio::test]
3325    async fn xrpc_error_body_is_capped() {
3326        // A syntactically VALID envelope, one byte past the cap. If the body were
3327        // parsed unbounded this would deserialize and yield "TooBig"; capped, it
3328        // is refused unread and degrades to the "Unknown" fallback.
3329        let filler = "x".repeat(crate::net::MAX_BODY_BYTES);
3330        let big = format!(r#"{{"error":"TooBig","message":"{filler}"}}"#).into_bytes();
3331        assert!(big.len() > crate::net::MAX_BODY_BYTES);
3332
3333        let base = crate::net::tests::serve_body(big).await;
3334        let resp = reqwest::Client::builder()
3335            .build()
3336            .unwrap()
3337            .get(&base)
3338            .send()
3339            .await
3340            .unwrap();
3341
3342        match xrpc_error_from(resp).await {
3343            AtProtoError::Xrpc { error, message, .. } => {
3344                assert_eq!(error, "Unknown", "an over-cap error body must not parse");
3345                assert!(message.is_none());
3346            }
3347            other => panic!("expected Xrpc, got {other:?}"),
3348        }
3349    }
3350
3351    /// The other half: a normal-sized envelope still parses, so capping the
3352    /// error path did not cost the diagnostics it exists to provide.
3353    #[tokio::test]
3354    async fn xrpc_error_body_within_the_cap_still_parses() {
3355        let base = crate::net::tests::serve_body(
3356            br#"{"error":"InvalidRequest","message":"bad rkey"}"#.to_vec(),
3357        )
3358        .await;
3359        let resp = reqwest::Client::builder()
3360            .build()
3361            .unwrap()
3362            .get(&base)
3363            .send()
3364            .await
3365            .unwrap();
3366
3367        match xrpc_error_from(resp).await {
3368            AtProtoError::Xrpc { error, message, .. } => {
3369                assert_eq!(error, "InvalidRequest");
3370                assert_eq!(message.as_deref(), Some("bad rkey"));
3371            }
3372            other => panic!("expected Xrpc, got {other:?}"),
3373        }
3374    }
3375
3376    /// A realistic `com.atproto.repo.listRecords` response for the subscription
3377    /// collection, as a PDS returns it — the envelope wraps each record in
3378    /// `{uri, cid, value}` and the record `value` carries its `$type`.
3379    fn subscription_list_json() -> Value {
3380        json!({
3381            "records": [
3382                {
3383                    "uri": "at://did:plc:abc123/community.lexicon.rss.subscription/3ksub0001",
3384                    "cid": "bafyreisubone",
3385                    "value": {
3386                        "$type": "community.lexicon.rss.subscription",
3387                        "url": "https://example.com/feed.xml",
3388                        "title": "Example Blog",
3389                        "siteUrl": "https://example.com/",
3390                        "fetchHint": "hourly",
3391                        "createdAt": "2026-07-12T00:00:00.000Z"
3392                    }
3393                },
3394                {
3395                    "uri": "at://did:plc:abc123/community.lexicon.rss.subscription/3ksub0002",
3396                    "cid": "bafyreisubtwo",
3397                    "value": {
3398                        "$type": "community.lexicon.rss.subscription",
3399                        "url": "https://blog.example.org/atom.xml",
3400                        "createdAt": "2026-07-11T12:00:00.000Z"
3401                    }
3402                }
3403            ],
3404            "cursor": "3ksub0002"
3405        })
3406    }
3407
3408    /// **A big archive is truncated, not refused.** `extend_bounded` bails on
3409    /// its cap, which is right for the `sub_ref` walk (a short list there is
3410    /// revoked access) and wrong for an additive read: a publication with more
3411    /// documents than the cap would return `Err` on every poll — permanently
3412    /// unreadable rather than partially read. 2 000 posts is an ordinary
3413    /// figure for a long-running blog.
3414    #[test]
3415    fn a_reading_walk_truncates_where_the_sub_ref_walk_refuses() {
3416        let page = |n: usize| -> Vec<RecordEntry> {
3417            (0..n)
3418                .map(|i| RecordEntry {
3419                    uri: format!("at://did:plc:x/c/{i}"),
3420                    cid: None,
3421                    value: Value::Null,
3422                })
3423                .collect()
3424        };
3425        let mut out = Vec::new();
3426        assert!(!extend_truncating(&mut out, page(2), 3), "not full yet");
3427        assert_eq!(out.len(), 2);
3428        // The page that overshoots contributes what fits, and says "stop".
3429        assert!(extend_truncating(&mut out, page(5), 3), "must report full");
3430        assert_eq!(out.len(), 3, "a reading walk must keep what fits");
3431        // The same overshoot is a hard error on the fail-closed path.
3432        let mut refused = Vec::new();
3433        assert!(extend_bounded(&mut refused, page(5), 3, "c").is_err());
3434        assert!(refused.is_empty(), "a refusal must leave nothing behind");
3435    }
3436
3437    /// **The cap counts the records the caller KEEPS, not the ones the repo
3438    /// holds.** A repo-wide cap applied before the caller's filter starves a
3439    /// quiet publication whose busy sibling fills the window: poll it, walk
3440    /// the newest 2 000 documents, discard all of them as the sibling's,
3441    /// return nothing — permanently, and worse with every post the sibling
3442    /// makes. The walk pages on until it has `max` MATCHING records (still
3443    /// bounded by `MAX_LIST_PAGES` requests).
3444    #[tokio::test]
3445    async fn the_cap_counts_matching_records_not_walked_ones() {
3446        // Every page: 4 records, only the last of which the caller wants.
3447        let records: Vec<Value> = (0..4)
3448            .map(|i| {
3449                serde_json::json!({
3450                    "uri": format!("at://did:plc:x/c/{i}"),
3451                    "value": {"mine": i == 3}
3452                })
3453            })
3454            .collect();
3455        let body = serde_json::json!({ "records": records, "cursor": serde_json::Value::Null })
3456            .to_string();
3457        let base = crate::net::tests::serve_body(body.into_bytes()).await;
3458        let port: u16 = base
3459            .trim_end_matches('/')
3460            .rsplit(':')
3461            .next()
3462            .unwrap()
3463            .parse()
3464            .unwrap();
3465        crate::net::test_host_override(
3466            "matching-pds.test",
3467            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
3468        );
3469        let client = PdsClient::anonymous(
3470            ssrf_test_client(),
3471            format!("http://matching-pds.test:{port}"),
3472            "did:plc:x",
3473        );
3474
3475        let kept = client
3476            .list_recent_matching("site.standard.document", 3, 100, |r| {
3477                r.value
3478                    .get("mine")
3479                    .and_then(Value::as_bool)
3480                    .unwrap_or(false)
3481            })
3482            .await
3483            .expect("walk failed")
3484            .records;
3485        // One page, no cursor: one match survives. The point is that the three
3486        // non-matching records did NOT consume the cap.
3487        assert_eq!(kept.len(), 1, "the filter ran after the cap, not before it");
3488    }
3489
3490    /// **A walk that stopped early says so.** Landing exactly on the cap, or
3491    /// running out of page budget, returns the same short `Vec` as a small
3492    /// collection — and the caller cannot tell them apart afterwards. That
3493    /// silence is how the starvation this walk exists to prevent came back one
3494    /// order of magnitude further out: a quiet publication whose busy sibling
3495    /// fills every page returns nothing, forever, looking healthy.
3496    #[tokio::test]
3497    async fn a_walk_that_stops_early_reports_itself_incomplete() {
3498        let records: Vec<Value> = (0..2)
3499            .map(|i| serde_json::json!({"uri": format!("at://did:plc:x/c/{i}"), "value": {}}))
3500            .collect();
3501        // Every page is full AND advertises another — the shape that lands on
3502        // the cap with the collection still going.
3503        let body = serde_json::json!({ "records": records, "cursor": "next" }).to_string();
3504        let base = crate::net::tests::serve_body(body.into_bytes()).await;
3505        let port: u16 = base
3506            .trim_end_matches('/')
3507            .rsplit(':')
3508            .next()
3509            .unwrap()
3510            .parse()
3511            .unwrap();
3512        crate::net::test_host_override(
3513            "incomplete-pds.test",
3514            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
3515        );
3516        let client = PdsClient::anonymous(
3517            ssrf_test_client(),
3518            format!("http://incomplete-pds.test:{port}"),
3519            "did:plc:x",
3520        );
3521
3522        let walk = client
3523            .list_recent_matching("c", 2, 100, |_| true)
3524            .await
3525            .expect("walk failed");
3526        assert_eq!(walk.records.len(), 2);
3527        assert!(
3528            !walk.complete,
3529            "a walk that filled its cap with pages still to come called itself complete"
3530        );
3531    }
3532
3533    /// The other side: a collection that runs out IS complete, so the caller
3534    /// does not warn about every ordinary small publication.
3535    #[tokio::test]
3536    async fn a_walk_that_exhausts_the_collection_reports_itself_complete() {
3537        let body = serde_json::json!({
3538            "records": [{"uri": "at://did:plc:x/c/1", "value": {}}]
3539        })
3540        .to_string();
3541        let base = crate::net::tests::serve_body(body.into_bytes()).await;
3542        let port: u16 = base
3543            .trim_end_matches('/')
3544            .rsplit(':')
3545            .next()
3546            .unwrap()
3547            .parse()
3548            .unwrap();
3549        crate::net::test_host_override(
3550            "complete-pds.test",
3551            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
3552        );
3553        let client = PdsClient::anonymous(
3554            ssrf_test_client(),
3555            format!("http://complete-pds.test:{port}"),
3556            "did:plc:x",
3557        );
3558
3559        let walk = client
3560            .list_recent_matching("c", 100, 100, |_| true)
3561            .await
3562            .expect("walk failed");
3563        assert_eq!(walk.records.len(), 1);
3564        assert!(walk.complete, "an exhausted collection is a complete read");
3565    }
3566
3567    /// **`{}` is not a page of zero records.** The records-presence guard
3568    /// landed on the OAuth client first; a proxy answering
3569    /// `{"ok":true,"data":{}}` kept the same `sub_ref`-wipe open on the
3570    /// sidecar path, and `{}` from a stranger's PDS made an empty publication
3571    /// look healthy.
3572    #[test]
3573    fn a_body_without_a_records_field_is_not_an_empty_page() {
3574        let err =
3575            parse_list_records(br#"{}"#).expect_err("`{}` was read as a page of zero records");
3576        assert!(format!("{err:#}").contains("no records field"), "{err:#}");
3577        let err = parse_list_records(br#"{"cursor":"c"}"#)
3578            .expect_err("a cursor-only body was read as a page");
3579        assert!(format!("{err:#}").contains("no records field"), "{err:#}");
3580        let page = parse_list_records(br#"{"records":[]}"#).unwrap();
3581        assert!(page.records.is_empty());
3582    }
3583
3584    /// Serve one oversized-but-well-formed body and point `host` at it.
3585    async fn serve_oversized(host: &str, shape: &str) -> String {
3586        let filler = "x".repeat(crate::net::MAX_BODY_BYTES);
3587        let body = shape.replace("PAD", &filler);
3588        assert!(body.len() > crate::net::MAX_BODY_BYTES);
3589        let base = crate::net::tests::serve_body(body.into_bytes()).await;
3590        let port: u16 = base
3591            .trim_end_matches('/')
3592            .rsplit(':')
3593            .next()
3594            .unwrap()
3595            .parse()
3596            .unwrap();
3597        crate::net::test_host_override(host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
3598        format!("http://{host}:{port}")
3599    }
3600
3601    /// **The DID document is the most remote-controlled body of the lot.**
3602    ///
3603    /// For a `did:web:` the host comes straight out of the DID, so whoever
3604    /// supplies the DID chooses the server. The SSRF guard proves the address
3605    /// is public; it says nothing about the body being finite.
3606    #[tokio::test]
3607    async fn the_did_document_read_is_capped() {
3608        let base = serve_oversized(
3609            "did-doc-cap.test",
3610            r##"{"service":[{"id":"#atproto_pds","type":"AtprotoPersonalDataServer","serviceEndpoint":"https://pds.example"}],"pad":"PAD"}"##,
3611        )
3612        .await;
3613        let err = resolve_did_to_pds(
3614            &ssrf_test_client(),
3615            &base,
3616            "did:plc:ohutz6x5acjmpuulp3x7wxxc",
3617        )
3618        .await
3619        .expect_err("an oversized DID document was buffered whole");
3620        assert!(
3621            format!("{err:#}").contains("cap"),
3622            "failed for the wrong reason: {err:#}"
3623        );
3624    }
3625
3626    /// `resolver_base` is a user-influenced PDS host, as this function's own
3627    /// guard comment says.
3628    #[tokio::test]
3629    async fn the_resolve_handle_read_is_capped() {
3630        let base = serve_oversized(
3631            "resolve-handle-cap.test",
3632            r#"{"did":"did:plc:ohutz6x5acjmpuulp3x7wxxc","pad":"PAD"}"#,
3633        )
3634        .await;
3635        let err = resolve_handle(&ssrf_test_client(), &base, "alice.example.com")
3636            .await
3637            .expect_err("an oversized resolveHandle body was buffered whole");
3638        assert!(
3639            format!("{err:#}").contains("cap"),
3640            "failed for the wrong reason: {err:#}"
3641        );
3642    }
3643
3644    /// **The sidecar's body is capped like every other body we read.**
3645    ///
3646    /// `/internal/repo` proxies whatever the account's PDS returned, so its
3647    /// size is remote-controlled by a host the reader chose and we did not.
3648    /// Every other response in this codebase goes through
3649    /// [`crate::net::read_capped`]; this one buffered the whole thing with
3650    /// `resp.json()`, so the 8 MB ceiling that bounds the direct PDS client
3651    /// simply did not exist on the sidecar backend — which is the default.
3652    #[tokio::test]
3653    async fn the_sidecar_client_caps_the_body_it_will_buffer() {
3654        // Well-formed, and past the cap. The guard has to fire on size, not
3655        // on the shape being wrong.
3656        let filler = "x".repeat(crate::net::MAX_BODY_BYTES);
3657        let body = format!(r#"{{"ok":true,"data":{{"records":[],"pad":"{filler}"}}}}"#);
3658        assert!(body.len() > crate::net::MAX_BODY_BYTES);
3659        let base = crate::net::tests::serve_body(body.into_bytes()).await;
3660        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
3661        let err = client
3662            .list_records(
3663                "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
3664                "app.feather.subscription",
3665                None,
3666                None,
3667            )
3668            .await
3669            .expect_err("an oversized sidecar body was buffered whole");
3670        assert!(
3671            format!("{err:#}").contains("cap"),
3672            "failed for the wrong reason: {err:#}"
3673        );
3674    }
3675
3676    /// The sidecar path needs the records guard too, not only the envelope
3677    /// one: `{"ok":true,"data":{}}` is what a proxy makes of an empty or
3678    /// unexpected upstream body.
3679    #[tokio::test]
3680    async fn the_sidecar_client_refuses_a_data_object_without_records() {
3681        let base = crate::net::tests::serve_body(br#"{"ok":true,"data":{}}"#.to_vec()).await;
3682        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
3683        let err = client
3684            .list_records(
3685                "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
3686                "app.feather.subscription",
3687                None,
3688                None,
3689            )
3690            .await
3691            .expect_err("`data: {}` was read as an empty repo");
3692        assert!(format!("{err:#}").contains("no records field"), "{err:#}");
3693    }
3694
3695    /// An exactly-full final page dropped nothing, so it must not warn that it
3696    /// did: `>=` reported truncation whenever the last page landed flush.
3697    #[test]
3698    fn an_exactly_full_page_is_not_a_truncation() {
3699        let page = |n: usize| -> Vec<RecordEntry> {
3700            (0..n)
3701                .map(|i| RecordEntry {
3702                    uri: format!("at://did:plc:x/c/{i}"),
3703                    cid: None,
3704                    value: Value::Null,
3705                })
3706                .collect()
3707        };
3708        let mut out = Vec::new();
3709        assert!(
3710            !extend_truncating(&mut out, page(3), 3),
3711            "a page that exactly fills the cap dropped nothing"
3712        );
3713        assert_eq!(out.len(), 3);
3714        assert!(
3715            extend_truncating(&mut out, page(1), 3),
3716            "one more IS a drop"
3717        );
3718        assert_eq!(out.len(), 3);
3719    }
3720
3721    /// **The reading walk USES the truncating accumulator.** The helper being
3722    /// correct is not the point — the previous round's bug was a guard that
3723    /// existed and was not called. Driven through a real server: one page of
3724    /// five records under a cap of three.
3725    #[tokio::test]
3726    async fn the_reading_walk_returns_a_truncated_archive_rather_than_an_error() {
3727        let records: Vec<Value> = (0..5)
3728            .map(|i| serde_json::json!({"uri": format!("at://did:plc:x/c/{i}"), "value": {}}))
3729            .collect();
3730        let body = serde_json::json!({ "records": records }).to_string();
3731        let base = crate::net::tests::serve_body(body.into_bytes()).await;
3732        let port: u16 = base
3733            .trim_end_matches('/')
3734            .rsplit(':')
3735            .next()
3736            .unwrap()
3737            .parse()
3738            .unwrap();
3739        crate::net::test_host_override(
3740            "truncating-pds.test",
3741            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
3742        );
3743        let client = PdsClient::anonymous(
3744            ssrf_test_client(),
3745            format!("http://truncating-pds.test:{port}"),
3746            "did:plc:x",
3747        );
3748
3749        let walk = client
3750            .list_recent_matching("site.standard.document", 3, 100, |_| true)
3751            .await
3752            .expect("a big archive must be readable, not an error");
3753        assert_eq!(
3754            walk.records.len(),
3755            3,
3756            "the walk did not truncate to its cap"
3757        );
3758        assert!(
3759            !walk.complete,
3760            "a truncated walk must not report completeness"
3761        );
3762
3763        // The fail-closed walk still refuses the same overshoot.
3764        let err = client
3765            .list_all_records("community.lexicon.rss.subscription")
3766            .await;
3767        assert!(
3768            err.is_ok() || format!("{:#}", err.unwrap_err()).contains("cap"),
3769            "the sub_ref walk must keep its refusal"
3770        );
3771    }
3772
3773    /// **A write is not "succeeded" because the status was 200.** The sidecar's
3774    /// `delete_record` and `apply_writes` discard the body entirely, so a
3775    /// `200 {"error": …}` reported success: the UI showed a reader
3776    /// unsubscribed while the record was still in their repo, and a whole
3777    /// batch of writes vanished silently.
3778    #[tokio::test]
3779    async fn the_sidecar_client_refuses_a_200_error_envelope_on_writes() {
3780        let base = crate::net::tests::serve_body(
3781            br#"{"ok":true,"data":{"error":"InvalidRequest","message":"nope"}}"#.to_vec(),
3782        )
3783        .await;
3784        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
3785        let did = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
3786        let err = client
3787            .delete_subscription(did, "rk1")
3788            .await
3789            .expect_err("a failed delete was reported as success");
3790        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
3791
3792        let err = client
3793            .apply_writes(
3794                did,
3795                &[WriteOp::Delete {
3796                    collection: lexicon::nsid::SUBSCRIPTION.to_string(),
3797                    rkey: "rk1".to_string(),
3798                }],
3799            )
3800            .await
3801            .expect_err("a failed batch was reported as success");
3802        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
3803    }
3804
3805    /// The sidecar proxies the PDS's body, so the same 2xx envelope arrives
3806    /// through `RepoOk.data` — a defaulted `Value` that deserialised into an
3807    /// empty page just as happily. Driven through the real client.
3808    #[tokio::test]
3809    async fn the_sidecar_client_refuses_a_200_error_envelope() {
3810        let base = crate::net::tests::serve_body(
3811            br#"{"ok":true,"data":{"error":"InvalidRequest","message":"nope"}}"#.to_vec(),
3812        )
3813        .await;
3814        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
3815        let err = client
3816            .list_records(
3817                "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
3818                "app.feather.subscription",
3819                None,
3820                None,
3821            )
3822            .await
3823            .expect_err("an error envelope was read as an empty page");
3824        assert!(
3825            format!("{err:#}").contains("InvalidRequest"),
3826            "failed for the wrong reason: {err:#}"
3827        );
3828    }
3829
3830    /// **Every listRecords caller refuses a 2xx error envelope, not just the
3831    /// anonymous one.** `oauth::xrpc::Repo` reads `records` off the JSON with
3832    /// `unwrap_or(Array([]))` and the sidecar's `RepoOk.data` is a defaulted
3833    /// `Value`, so a PDS answering 200 with an envelope reached
3834    /// `resolve_subscriptions` as `Ok(empty)` — which is not the fail-closed
3835    /// branch, so `sync_sub_refs` DELETEd the DID's whole `sub_ref` projection:
3836    /// one bad response revokes a reader's access to every feed they have.
3837    #[test]
3838    fn an_error_envelope_is_refused_whatever_shape_it_arrives_in() {
3839        let envelope = serde_json::json!({"error": "InvalidRequest", "message": "bad cursor"});
3840        let err = reject_error_envelope(&envelope).expect_err("an envelope passed as data");
3841        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
3842        // A real page, and an empty real page, are both data.
3843        reject_error_envelope(&serde_json::json!({"records": []})).expect("an empty page is data");
3844        reject_error_envelope(&serde_json::json!({"records": [], "cursor": "c"})).unwrap();
3845    }
3846
3847    /// **A 200 carrying an error envelope is not an empty page.** `records` is
3848    /// `#[serde(default)]`, so `{"error": "...", "message": "..."}` on a 200
3849    /// deserialised as zero records — and a walk over a stranger's documents
3850    /// then returned a healthy, empty feed instead of an error. Some PDS
3851    /// implementations do answer 200 for application-level failures.
3852    #[test]
3853    fn a_200_with_an_error_envelope_is_not_an_empty_page() {
3854        let err = parse_list_records(br#"{"error":"InvalidRequest","message":"bad cursor"}"#)
3855            .expect_err("an error envelope parsed as a page");
3856        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
3857        let page = parse_list_records(br#"{"records":[]}"#).expect("an empty page is a page");
3858        assert!(page.records.is_empty() && page.cursor.is_none());
3859    }
3860
3861    /// **Both invariants, now read out of the bytes rather than out of a
3862    /// `Value`.** Parsing once is the point of the change; parsing once while
3863    /// quietly dropping a guard would be a much worse trade, and these are the
3864    /// shapes those guards exist for.
3865    #[test]
3866    fn parsing_a_page_from_bytes_keeps_both_invariants() {
3867        // Each shape names the reason it must fail for. Accepting either
3868        // message would let the envelope guard be deleted without a test
3869        // noticing, because an error envelope also has no `records` field — so
3870        // it keeps failing, for a reason that stops applying the day a PDS
3871        // returns an envelope alongside a records array.
3872        for (label, body, because) in [
3873            (
3874                "an error envelope on a 2xx",
3875                &br#"{"error":"InvalidRequest","message":"bad cursor"}"#[..],
3876                "error envelope",
3877            ),
3878            (
3879                "an envelope that also carries records",
3880                &br#"{"error":"InvalidRequest","records":[]}"#[..],
3881                "error envelope",
3882            ),
3883            (
3884                "a body with no records field",
3885                &br#"{"cursor":"c"}"#[..],
3886                "no records",
3887            ),
3888            ("a proxy's empty object", &br#"{}"#[..], "no records"),
3889            ("an empty body", &b""[..], "no records"),
3890        ] {
3891            let err = parse_list_records(body)
3892                .map(|p| panic!("{label} was read as a page of {} records", p.records.len()))
3893                .unwrap_err();
3894            let msg = format!("{err:#}");
3895            assert!(
3896                msg.contains(because),
3897                "{label} should have failed on {because:?}, got: {msg}"
3898            );
3899        }
3900        let page = parse_list_records(br#"{"records":[],"cursor":"c"}"#)
3901            .expect("a genuinely empty page is still a page");
3902        assert!(page.records.is_empty());
3903        assert_eq!(page.cursor.as_deref(), Some("c"));
3904    }
3905
3906    #[test]
3907    fn list_records_envelope_deserializes() {
3908        let resp: ListRecordsResponse =
3909            serde_json::from_value(subscription_list_json()).expect("envelope");
3910        assert_eq!(resp.records.len(), 2);
3911        assert_eq!(resp.cursor.as_deref(), Some("3ksub0002"));
3912        assert_eq!(resp.records[0].cid.as_deref(), Some("bafyreisubone"));
3913    }
3914
3915    #[test]
3916    fn record_entry_rkey_is_last_uri_segment() {
3917        let resp: ListRecordsResponse =
3918            serde_json::from_value(subscription_list_json()).expect("envelope");
3919        assert_eq!(resp.records[0].rkey(), Some("3ksub0001"));
3920        assert_eq!(resp.records[1].rkey(), Some("3ksub0002"));
3921    }
3922
3923    #[test]
3924    fn record_value_parses_into_lexicon_subscription() {
3925        let resp: ListRecordsResponse =
3926            serde_json::from_value(subscription_list_json()).expect("envelope");
3927
3928        let full: Subscription = resp.records[0].parse().expect("parse full sub");
3929        assert_eq!(full.r#type, lexicon::nsid::SUBSCRIPTION);
3930        assert_eq!(full.url, "https://example.com/feed.xml");
3931        assert_eq!(full.title.as_deref(), Some("Example Blog"));
3932        assert_eq!(full.site_url.as_deref(), Some("https://example.com/"));
3933        assert_eq!(full.fetch_hint, Some(lexicon::FetchHint::Hourly));
3934
3935        let minimal: Subscription = resp.records[1].parse().expect("parse minimal sub");
3936        assert_eq!(minimal.url, "https://blog.example.org/atom.xml");
3937        assert!(minimal.title.is_none());
3938    }
3939
3940    fn ssrf_test_client() -> Client {
3941        Client::builder()
3942            .user_agent(crate::USER_AGENT)
3943            .build()
3944            .unwrap()
3945    }
3946
3947    /// A hostile `did:web` whose host is the cloud-metadata address must be
3948    /// REFUSED before any request leaves the box — the DID-document fetch now
3949    /// routes through the SSRF guard (`guarded_get_no_privacy`), which rejects
3950    /// link-local / metadata targets.
3951    #[tokio::test]
3952    async fn resolve_did_web_blocks_metadata_host() {
3953        let client = ssrf_test_client();
3954        let err = resolve_did_to_pds(&client, "https://plc.directory", "did:web:169.254.169.254")
3955            .await
3956            .unwrap_err()
3957            .to_string();
3958        assert!(
3959            err.contains("forbidden") || err.contains("internal"),
3960            "expected an SSRF refusal, got: {err}"
3961        );
3962    }
3963
3964    /// A `did:web` pointing at loopback is likewise blocked (internal service
3965    /// reflection).
3966    #[tokio::test]
3967    async fn resolve_did_web_blocks_loopback_host() {
3968        let client = ssrf_test_client();
3969        let err = resolve_did_to_pds(&client, "https://plc.directory", "did:web:127.0.0.1")
3970            .await
3971            .unwrap_err()
3972            .to_string();
3973        assert!(
3974            err.contains("forbidden") || err.contains("internal"),
3975            "expected an SSRF refusal, got: {err}"
3976        );
3977    }
3978
3979    /// `resolve_handle` against a metadata/loopback resolver base is also guarded
3980    /// (the base can come from a prior hostile DID-doc resolution).
3981    #[tokio::test]
3982    async fn resolve_handle_blocks_metadata_resolver_base() {
3983        let client = ssrf_test_client();
3984        let err = resolve_handle(&client, "http://169.254.169.254", "alice.example.com")
3985            .await
3986            .unwrap_err()
3987            .to_string();
3988        assert!(
3989            err.contains("forbidden") || err.contains("internal"),
3990            "expected an SSRF refusal, got: {err}"
3991        );
3992    }
3993
3994    /// A resolved `serviceEndpoint` that targets an internal host is rejected at
3995    /// resolve time via [`crate::net::assert_public_target`], so it can never be
3996    /// handed to a raw XRPC client.
3997    #[tokio::test]
3998    async fn service_endpoint_internal_target_rejected() {
3999        assert!(crate::net::assert_public_target("http://169.254.169.254/")
4000            .await
4001            .is_err());
4002        assert!(crate::net::assert_public_target("http://127.0.0.1:3000/")
4003            .await
4004            .is_err());
4005        // A public endpoint literal passes.
4006        assert!(crate::net::assert_public_target("https://1.1.1.1/")
4007            .await
4008            .is_ok());
4009    }
4010
4011    /// A `PdsClient` pointed at an internal `pds_base`, as an attacker-controlled
4012    /// DID document could arrange between the `assert_public_target` at resolve
4013    /// time and the request.
4014    fn internal_target_client(pds_base: &str) -> PdsClient {
4015        PdsClient::new(
4016            ssrf_test_client(),
4017            pds_base,
4018            "did:plc:victim",
4019            Auth::Session(SessionAuth {
4020                did: "did:plc:victim".to_string(),
4021                handle: None,
4022                access_jwt: "session-bearer-must-not-leak".to_string(),
4023                refresh_jwt: None,
4024            }),
4025        )
4026    }
4027
4028    /// **Regression (v0.2.8):** every `com.atproto.repo.*` WRITE must go through
4029    /// the SSRF guard, not the shared client. Before the fix only `list_records`
4030    /// was guarded, so `createRecord` / `putRecord` / `deleteRecord` /
4031    /// `applyWrites` would happily deliver the session bearer to
4032    /// `169.254.169.254` or loopback on a rebound host.
4033    #[tokio::test]
4034    async fn every_repo_write_is_refused_against_an_internal_pds() {
4035        for base in [
4036            "http://169.254.169.254",
4037            "http://127.0.0.1:9",
4038            "http://[::1]",
4039        ] {
4040            let client = internal_target_client(base);
4041            let sub = Subscription::new("https://example.com/feed.xml", "2026-08-13T00:00:00Z");
4042
4043            let mut errors = vec![
4044                client
4045                    .create_record(lexicon::nsid::SUBSCRIPTION, &sub)
4046                    .await
4047                    .unwrap_err()
4048                    .to_string(),
4049                client
4050                    .put_record(lexicon::nsid::SUBSCRIPTION, "rkey", &sub)
4051                    .await
4052                    .unwrap_err()
4053                    .to_string(),
4054                client
4055                    .delete_record(lexicon::nsid::SUBSCRIPTION, "rkey")
4056                    .await
4057                    .unwrap_err()
4058                    .to_string(),
4059            ];
4060            errors.push(
4061                client
4062                    .apply_writes(&[WriteOp::Delete {
4063                        collection: lexicon::nsid::SUBSCRIPTION.to_string(),
4064                        rkey: "rkey".to_string(),
4065                    }])
4066                    .await
4067                    .unwrap_err()
4068                    .to_string(),
4069            );
4070
4071            for err in errors {
4072                assert!(
4073                    err.contains("forbidden") || err.contains("internal"),
4074                    "{base}: expected an SSRF refusal, got: {err}"
4075                );
4076            }
4077        }
4078    }
4079
4080    /// **Regression (v0.2.8):** the app password travels in the request BODY,
4081    /// where reqwest's cross-origin header sanitisation cannot protect it — so
4082    /// `createSession` is guarded too, and a rebound/internal `pds_base` never
4083    /// receives it.
4084    #[tokio::test]
4085    async fn app_password_login_is_refused_against_an_internal_pds() {
4086        let client = ssrf_test_client();
4087        for base in ["http://169.254.169.254", "http://127.0.0.1:9"] {
4088            let err = login_with_app_password(&client, base, "alice.example.com", "hunter2-app-pw")
4089                .await
4090                .unwrap_err()
4091                .to_string();
4092            assert!(
4093                err.contains("forbidden") || err.contains("internal"),
4094                "{base}: expected an SSRF refusal, got: {err}"
4095            );
4096        }
4097    }
4098
4099    /// An anonymous client is read-only: the write paths fail closed on
4100    /// [`Auth::bearer`] before any socket work, so `Auth::Anonymous` can never
4101    /// become a credential-less write primitive against a stranger's PDS.
4102    #[tokio::test]
4103    async fn anonymous_client_cannot_write() {
4104        let client = PdsClient::anonymous(
4105            ssrf_test_client(),
4106            "https://pds.example.com",
4107            "did:plc:stranger",
4108        );
4109        let err = client
4110            .delete_record(lexicon::nsid::SUBSCRIPTION, "rkey")
4111            .await
4112            .unwrap_err()
4113            .to_string();
4114        assert!(
4115            err.contains("no credentials") || err.contains("anonymous") || err.contains("bearer"),
4116            "expected a fail-closed auth error, got: {err}"
4117        );
4118    }
4119
4120    #[test]
4121    fn write_result_deserializes() {
4122        let wr: WriteResult = serde_json::from_value(json!({
4123            "uri": "at://did:plc:abc123/community.lexicon.rss.subscription/3ksubnew",
4124            "cid": "bafyreinew"
4125        }))
4126        .expect("write result");
4127        assert!(wr.uri.ends_with("3ksubnew"));
4128        assert_eq!(wr.cid.as_deref(), Some("bafyreinew"));
4129    }
4130
4131    #[test]
4132    fn did_document_finds_pds_endpoint() {
4133        let doc: DidDocument = serde_json::from_value(json!({
4134            "id": "did:plc:abc123",
4135            "service": [
4136                {
4137                    "id": "#atproto_pds",
4138                    "type": "AtprotoPersonalDataServer",
4139                    "serviceEndpoint": "https://pds.example.com/"
4140                }
4141            ]
4142        }))
4143        .expect("did doc");
4144        assert_eq!(
4145            doc.pds_endpoint().as_deref(),
4146            Some("https://pds.example.com")
4147        );
4148    }
4149
4150    #[test]
4151    fn did_document_without_pds_yields_none() {
4152        let doc: DidDocument = serde_json::from_value(json!({
4153            "id": "did:plc:abc123",
4154            "service": []
4155        }))
4156        .expect("did doc");
4157        assert!(doc.pds_endpoint().is_none());
4158    }
4159
4160    #[test]
4161    fn session_auth_deserializes_create_session_shape() {
4162        let session: SessionAuth = serde_json::from_value(json!({
4163            "did": "did:plc:abc123",
4164            "handle": "alice.example.com",
4165            "accessJwt": "eyJh...access",
4166            "refreshJwt": "eyJh...refresh"
4167        }))
4168        .expect("session");
4169        assert_eq!(session.did, "did:plc:abc123");
4170        assert_eq!(session.handle.as_deref(), Some("alice.example.com"));
4171        let auth = Auth::Session(session);
4172        assert_eq!(auth.bearer().expect("bearer"), "eyJh...access");
4173    }
4174
4175    #[test]
4176    fn oauth_variant_carries_no_direct_bearer() {
4177        let auth = Auth::Oauth(OauthPlaceholder::default());
4178        assert!(
4179            auth.bearer().is_err(),
4180            "Auth::Oauth carries no direct bearer — the sidecar owns the OAuth path"
4181        );
4182    }
4183
4184    #[test]
4185    fn anonymous_variant_carries_no_bearer() {
4186        let err = Auth::Anonymous.bearer().unwrap_err().to_string();
4187        assert!(
4188            err.contains("anonymous"),
4189            "the anonymous refusal must name itself, got: {err}"
4190        );
4191    }
4192
4193    #[test]
4194    fn anonymous_client_targets_the_requested_repo() {
4195        let client = PdsClient::anonymous(
4196            ssrf_test_client(),
4197            "https://pds.example.com/",
4198            "did:plc:abc123",
4199        );
4200        // The trailing slash is trimmed so `xrpc_url` joins cleanly.
4201        assert_eq!(client.pds_base(), "https://pds.example.com");
4202        assert_eq!(client.did(), "did:plc:abc123");
4203        // …and it holds no credential.
4204        assert!(client.auth.bearer().is_err());
4205    }
4206
4207    /// The regression test for the defect this milestone fixes: `list_records`
4208    /// used to send on the shared client, bypassing the SSRF guard entirely. It
4209    /// now routes through `net::guarded_get_no_privacy`, so an internal
4210    /// `pds_base` is refused before a packet leaves the box. Hermetic — the hosts
4211    /// are IP literals, rejected without any DNS lookup or connect.
4212    #[tokio::test]
4213    async fn list_records_blocks_internal_pds_base() {
4214        for base in ["http://169.254.169.254", "http://127.0.0.1:1"] {
4215            let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
4216            let err = client
4217                .list_records(lexicon::nsid::SUBSCRIPTION, Some(1), None)
4218                .await
4219                .unwrap_err()
4220                .to_string();
4221            assert!(
4222                err.contains("forbidden") || err.contains("internal"),
4223                "expected an SSRF refusal for {base}, got: {err}"
4224            );
4225        }
4226    }
4227
4228    /// The guard is not anonymous-only: an *authenticated* client reading a
4229    /// hostile PDS base is blocked identically. (That path was only ever safe by
4230    /// accident of usage.)
4231    #[tokio::test]
4232    async fn list_records_guard_applies_to_authed_clients_too() {
4233        let auth = Auth::Session(SessionAuth {
4234            did: "did:plc:x".to_string(),
4235            handle: None,
4236            access_jwt: "x".to_string(),
4237            refresh_jwt: None,
4238        });
4239        let client = PdsClient::new(
4240            ssrf_test_client(),
4241            "http://169.254.169.254",
4242            "did:plc:x",
4243            auth,
4244        );
4245        let err = client
4246            .list_records(lexicon::nsid::SUBSCRIPTION, Some(1), None)
4247            .await
4248            .unwrap_err()
4249            .to_string();
4250        assert!(
4251            err.contains("forbidden") || err.contains("internal"),
4252            "expected an SSRF refusal, got: {err}"
4253        );
4254    }
4255
4256    #[test]
4257    fn apply_writes_ops_render_tagged_union() {
4258        let create = WriteOp::Create {
4259            collection: lexicon::nsid::SUBSCRIPTION.to_string(),
4260            rkey: None,
4261            value: json!({"url": "https://example.com/feed.xml"}),
4262        };
4263        let update = WriteOp::Update {
4264            collection: lexicon::nsid::READ_STATE.to_string(),
4265            rkey: "feedhash01".to_string(),
4266            value: json!({"feedUrl": "https://example.com/feed.xml"}),
4267        };
4268        let delete = WriteOp::Delete {
4269            collection: lexicon::nsid::SAVED.to_string(),
4270            rkey: "3ksaved01".to_string(),
4271        };
4272
4273        assert_eq!(
4274            create.to_json()["$type"],
4275            json!("com.atproto.repo.applyWrites#create")
4276        );
4277        // A create with no explicit rkey omits the field (server assigns a tid).
4278        assert!(create.to_json().get("rkey").is_none());
4279
4280        assert_eq!(
4281            update.to_json()["$type"],
4282            json!("com.atproto.repo.applyWrites#update")
4283        );
4284        assert_eq!(update.to_json()["rkey"], json!("feedhash01"));
4285
4286        assert_eq!(
4287            delete.to_json()["$type"],
4288            json!("com.atproto.repo.applyWrites#delete")
4289        );
4290        assert_eq!(delete.to_json()["rkey"], json!("3ksaved01"));
4291    }
4292
4293    #[test]
4294    fn read_state_flush_creates_first_then_updates() {
4295        // A cursor whose PDS record does NOT yet exist (pds_created = false) must
4296        // become a CREATE op at its stable rkey — NOT a bare update, which would
4297        // error on the missing record and (applyWrites being atomic per-repo) drop
4298        // the whole batch on a feed's first flush.
4299        let fresh = (
4300            "rs-fresh".to_string(),
4301            ReadState::new("https://a.example/feed.xml", None, "2026-07-12T00:00:00Z"),
4302            false,
4303        );
4304        // An already-created cursor updates in place.
4305        let existing = (
4306            "rs-existing".to_string(),
4307            ReadState::new(
4308                "https://b.example/feed.xml",
4309                Some("2026-07-11T00:00:00Z".to_string()),
4310                "2026-07-12T00:00:00Z",
4311            ),
4312            true,
4313        );
4314
4315        let ops = read_state_write_ops(&[fresh, existing]).expect("build ops");
4316        assert_eq!(ops.len(), 2);
4317
4318        // First op: a create carrying the stable rkey (put/create, not update).
4319        let create = ops[0].to_json();
4320        assert_eq!(
4321            create["$type"],
4322            json!("com.atproto.repo.applyWrites#create"),
4323            "first flush of a new feed must CREATE its readState record"
4324        );
4325        assert_eq!(create["rkey"], json!("rs-fresh"));
4326        // The created record omits readThrough (F1): backlog not implicitly read.
4327        assert!(create["value"].get("readThrough").is_none());
4328
4329        // Second op: an update for the already-created record.
4330        let update = ops[1].to_json();
4331        assert_eq!(
4332            update["$type"],
4333            json!("com.atproto.repo.applyWrites#update")
4334        );
4335        assert_eq!(update["rkey"], json!("rs-existing"));
4336
4337        // Both ride the SAME batch — batching is preserved.
4338        assert_eq!(ops.len(), 2);
4339    }
4340
4341    #[test]
4342    fn urlencode_escapes_did_colons_and_keeps_unreserved() {
4343        assert_eq!(urlencode("did:plc:abc123"), "did%3Aplc%3Aabc123");
4344        assert_eq!(
4345            urlencode("community.lexicon.rss.subscription"),
4346            "community.lexicon.rss.subscription"
4347        );
4348        assert_eq!(urlencode("a b&c"), "a%20b%26c");
4349    }
4350
4351    #[test]
4352    fn xrpc_record_not_found_is_detected() {
4353        let err = AtProtoError::Xrpc {
4354            status: StatusCode::BAD_REQUEST,
4355            error: "RecordNotFound".to_string(),
4356            message: Some("Could not locate record".to_string()),
4357        };
4358        assert!(err.is_record_not_found());
4359    }
4360
4361    // -- reader-facing CRUD: rkey extraction --------------------------------
4362
4363    #[test]
4364    fn write_result_extracts_rkey_from_uri() {
4365        let wr: WriteResult = serde_json::from_value(json!({
4366            "uri": "at://did:plc:abc123/community.lexicon.rss.subscription/3ksubnew",
4367            "cid": "bafyreinew"
4368        }))
4369        .expect("write result");
4370        assert_eq!(wr.rkey(), Some("3ksubnew"));
4371        assert_eq!(wr.into_rkey(), "3ksubnew");
4372    }
4373
4374    // -- reader-facing CRUD: deterministic sort orders ----------------------
4375    //
4376    // The `list_*_sorted` wrappers only add an ordering on top of the network
4377    // `list_*` read, so we exercise the *comparator* here on representative
4378    // data (parsed from a listRecords-shaped envelope) with no network.
4379
4380    // -- reader-facing CRUD: bulk applyWrites shape (OPML import) ------------
4381
4382    /// **Bulk subscribe, through the real client, asserted on the bytes it
4383    /// sent.** The test this replaces built the `WriteOp::Create` ops itself
4384    /// ("mirror what `add_subscriptions_bulk` builds") and asserted on its own
4385    /// construction; the function was never called, and writing every feed
4386    /// into the wrong collection with server-assigned rkeys left the suite
4387    /// green. Three atproto sort tests that re-implemented the comparator
4388    /// inline are deleted alongside — `lexicon::sort_tests` fails their
4389    /// mutation, and they added nothing but a misleading name.
4390    #[tokio::test]
4391    async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
4392        let (base, log) =
4393            crate::net::tests::serve_json_capturing(br#"{"ok":true,"data":{}}"#.to_vec()).await;
4394        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
4395        let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
4396            .map(|i| {
4397                crate::vetted::VettedSubscription::new(&lexicon::Subscription::new(
4398                    format!("https://f{i}.example/feed.xml"),
4399                    "2026-07-12T00:00:00.000Z",
4400                ))
4401            })
4402            .collect();
4403
4404        let rkeys = client
4405            .add_subscriptions_bulk("did:plc:ewvi7nxzyoun6zhxrhs64oiz", &subs)
4406            .await
4407            .expect("bulk write failed");
4408
4409        let sent = log.lock().unwrap().clone();
4410        assert_eq!(
4411            sent.len(),
4412            1,
4413            "expected one applyWrites request, got {sent:?}"
4414        );
4415        let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
4416            .expect("request body is JSON");
4417        let writes = body["writes"].as_array().expect("writes array");
4418        assert_eq!(writes.len(), 3);
4419        for (i, w) in writes.iter().enumerate() {
4420            assert_eq!(
4421                w["collection"],
4422                lexicon::nsid::SUBSCRIPTION,
4423                "write {i} went to the wrong collection"
4424            );
4425            assert_eq!(
4426                w["rkey"].as_str(),
4427                Some(rkeys[i].as_str()),
4428                "write {i} does not carry the rkey the client returned"
4429            );
4430        }
4431        let mut sorted = rkeys.clone();
4432        sorted.sort();
4433        assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
4434        assert_eq!(
4435            rkeys.iter().collect::<std::collections::HashSet<_>>().len(),
4436            3,
4437            "rkeys must be distinct"
4438        );
4439    }
4440
4441    /// **The walk stops on a repeated cursor.** `MAX_LIST_PAGES`, the
4442    /// same-cursor guard and the `got > 0` guard had no test; only
4443    /// `extend_bounded` was covered directly. A PDS that echoes the same
4444    /// cursor forever would otherwise be walked for 200 pages.
4445    #[tokio::test]
4446    async fn list_all_records_stops_on_a_repeated_cursor() {
4447        let body = serde_json::json!({
4448            "records": [{"uri": "at://did:plc:x/c/1", "value": {}}],
4449            "cursor": "same-every-time"
4450        })
4451        .to_string();
4452        let base = crate::net::tests::serve_body(body.into_bytes()).await;
4453        let port: u16 = base
4454            .trim_end_matches('/')
4455            .rsplit(':')
4456            .next()
4457            .unwrap()
4458            .parse()
4459            .unwrap();
4460        crate::net::test_host_override(
4461            "repeated-cursor.test",
4462            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
4463        );
4464        let client = PdsClient::anonymous(
4465            ssrf_test_client(),
4466            format!("http://repeated-cursor.test:{port}"),
4467            "did:plc:x",
4468        );
4469        let records = client.list_all_records("c").await.expect("walk failed");
4470        // Page 1: cursor None → "same". Page 2: "same" again → stop, after
4471        // taking that page. Two pages, not two hundred.
4472        assert_eq!(records.len(), 2, "a repeated cursor was followed");
4473    }
4474
4475    // -- bounding the parse before it allocates -----------------------------
4476
4477    /// **Strings cannot invent structure.** The subtle half of the bound: an
4478    /// article containing a million commas is one node, and counting naively
4479    /// would refuse it.
4480    #[test]
4481    fn structure_inside_a_string_is_not_structure() {
4482        let prose = format!(
4483            r#"{{"records":[{{"uri":"at://d/c/r","value":{{"t":"{}"}}}}]}}"#,
4484            "a,b,[c],{d}:e,".repeat(50_000)
4485        );
4486        let bound = count_structural_chars(prose.as_bytes());
4487        assert!(
4488            bound < 100,
4489            "a page of prose full of punctuation was counted as {bound} nodes"
4490        );
4491        assert!(
4492            parse_list_records(prose.as_bytes()).is_ok(),
4493            "a legitimate page of prose was refused"
4494        );
4495    }
4496
4497    #[test]
4498    fn a_node_explosion_is_refused_before_it_is_parsed() {
4499        // ~8 MB of the cheapest node there is, which is the measured attack.
4500        let mut body = String::from(r#"{"records":[{"uri":"at://d/c/r","value":["#);
4501        for _ in 0..1_200_000 {
4502            body.push_str("{},");
4503        }
4504        body.push_str(r#"{}]}]}"#);
4505        assert!(
4506            body.len() > 3_000_000,
4507            "the probe body is {} bytes",
4508            body.len()
4509        );
4510
4511        let bound = count_structural_chars(body.as_bytes());
4512        assert!(
4513            bound > MAX_LIST_STRUCTURAL_CHARS,
4514            "the attack shape was counted as only {bound} nodes"
4515        );
4516        let err = parse_list_records(body.as_bytes())
4517            .expect_err("a node explosion was parsed rather than refused");
4518        assert!(
4519            format!("{err:#}").contains("structural characters"),
4520            "failed for the wrong reason: {err:#}"
4521        );
4522    }
4523
4524    /// **The densest page the lexicons permit must fit, with room.**
4525    ///
4526    /// This is the floor under [`MAX_LIST_STRUCTURAL_CHARS`], and it is the reason the cap
4527    /// is 640 000 rather than the ~150 000 that would otherwise hold the memory
4528    /// claim comfortably. A `readState` record carries up to
4529    /// [`crate::lexicon::ReadState::MAX_IDS`] read ids, and a page carries 100 of
4530    /// them — far denser in nodes than a page of articles, which is mostly
4531    /// prose. Tighten the cap below this and a reader with a lot of history
4532    /// stops being able to sync at all.
4533    #[test]
4534    fn a_full_read_state_page_fits_under_the_cap() {
4535        let ids: Vec<String> = (0..crate::lexicon::ReadState::MAX_IDS)
4536            .map(|i| format!("https://example.com/blog/post-{i}"))
4537            .collect();
4538        let records: Vec<serde_json::Value> = (0..100)
4539            .map(|i| {
4540                serde_json::json!({
4541                    "uri": format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/community.lexicon.rss.readState/3lab{i}"),
4542                    "cid": "bafyreiabc123def456ghi789jkl012mno345pqr678stu901",
4543                    "value": {
4544                        "$type": "community.lexicon.rss.readState",
4545                        "feedUrl": "https://example.com/feed.xml",
4546                        "readThrough": "2026-07-11T09:30:00Z",
4547                        // **BOTH arrays, because the lexicon permits both.**
4548                        // Filling only `readIds` counted 203 503 nodes, so a cap
4549                        // as low as 300 000 passed every test in the suite while
4550                        // refusing the very page this test exists to protect.
4551                        "readIds": ids,
4552                        "unreadIds": ids,
4553                    }
4554                })
4555            })
4556            .collect();
4557        let body = serde_json::json!({ "records": records }).to_string();
4558        let bound = count_structural_chars(body.as_bytes());
4559        assert!(
4560            bound < MAX_LIST_STRUCTURAL_CHARS,
4561            "the densest legitimate page counts {bound} nodes against a cap of {MAX_LIST_STRUCTURAL_CHARS}"
4562        );
4563        // And it is dense enough to be the floor the cap was chosen for: a page
4564        // counting only a fifth of the cap would pass the assertion above while
4565        // leaving the cap free to drop far below real traffic.
4566        assert!(
4567            bound > MAX_LIST_STRUCTURAL_CHARS / 2,
4568            "this page counts only {bound} nodes, so it is no longer the floor \
4569             `MAX_LIST_STRUCTURAL_CHARS` was measured against and a much tighter cap would \
4570             pass it"
4571        );
4572        assert!(
4573            parse_list_records(body.as_bytes()).is_ok(),
4574            "a full read-state page was refused"
4575        );
4576    }
4577
4578    /// The write path takes the same guard as the listing path.
4579    #[tokio::test]
4580    async fn the_write_path_refuses_a_node_explosion() {
4581        let mut data = String::from(r#"{"ok":true,"data":{"records":["#);
4582        for _ in 0..700_000 {
4583            data.push_str("{},");
4584        }
4585        data.push_str(r#"{}]}}"#);
4586        let base = crate::net::tests::serve_body(data.into_bytes()).await;
4587        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
4588        let err = client
4589            .delete_record("did:plc:ewvi7nxzyoun6zhxrhs64oiz", "c", "r")
4590            .await
4591            .expect_err("a node explosion reached the parser on the write path");
4592        assert!(
4593            format!("{err:#}").contains("structural characters"),
4594            "failed for the wrong reason: {err:#}"
4595        );
4596    }
4597
4598    #[test]
4599    fn an_ordinary_page_is_nowhere_near_the_structure_bound() {
4600        // **A page, not a record.** This served `paged_bodies(1, 17_000, false)` —
4601        // ONE page holding ONE record — which counts about eleven characters, so
4602        // the old `< 1_000` assertion held by three orders of magnitude and would
4603        // have passed with the cap at 1 000. It also made the doc comment's "a
4604        // page of 100 documents measures 40 000" untested, and that figure was
4605        // wrong by 10x.
4606        let records: Vec<serde_json::Value> = (0..100)
4607            .map(|i| {
4608                serde_json::json!({
4609                    "uri": format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.document/3lab{i}"),
4610                    "cid": "bafyreiabc123def456ghi789jkl012mno345pqr678stu901",
4611                    "value": {
4612                        "$type": "site.standard.document",
4613                        "title": "A reasonably typical post title",
4614                        "path": format!("/posts/{i}"),
4615                        "site": "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
4616                        "publishedAt": "2026-07-11T09:30:00Z",
4617                        "description": "x".repeat(120),
4618                        "textContent": "y".repeat(15_000),
4619                    }
4620                })
4621            })
4622            .collect();
4623        let body = serde_json::json!({ "records": records }).to_string();
4624        // A real page of documents is mostly prose: 1.5 MB on the wire for 4 003
4625        // counted characters.
4626        assert!(
4627            body.len() > 1_000_000,
4628            "the probe page is only {} bytes, so it is not a full page",
4629            body.len(),
4630        );
4631        let counted = count_structural_chars(body.as_bytes());
4632        assert!(
4633            (3_500..4_500).contains(&counted),
4634            "a page of 100 documents counted {counted}, not the ~4 003 the cap's \
4635             doc comment claims — the ordinary-traffic end of the bracket moved",
4636        );
4637        assert!(
4638            counted * 100 < MAX_LIST_STRUCTURAL_CHARS,
4639            "ordinary traffic is within 100x of the cap ({counted} against \
4640             {MAX_LIST_STRUCTURAL_CHARS}), which is not the headroom the cap claims",
4641        );
4642    }
4643
4644    /// **An escaped quote does not end the string**, asserted without a magic
4645    /// number: the same document with the escape replaced by a plain letter has
4646    /// the same structure, so it must count the same. Get the escape wrong and the
4647    /// scanner leaves the string early and counts the rest as structure.
4648    #[test]
4649    fn an_escaped_quote_does_not_end_the_string() {
4650        let escaped = br#"{"records":[{"uri":"a\"b","value":{}}],"cursor":"x"}"#;
4651        let plain = br#"{"records":[{"uri":"axb","value":{}}],"cursor":"x"}"#;
4652        assert_eq!(
4653            count_structural_chars(escaped),
4654            count_structural_chars(plain),
4655            "an escaped quote changed the structure count"
4656        );
4657        let backslash = br#"{"records":[],"cursor":"x\\"}"#;
4658        let letter = br#"{"records":[],"cursor":"xy"}"#;
4659        assert_eq!(
4660            count_structural_chars(backslash),
4661            count_structural_chars(letter),
4662            "an escaped backslash changed the structure count"
4663        );
4664    }
4665
4666    /// A string is itself a node, so a page of strings costs more than a page of
4667    /// numbers. Without that, an array of a million short strings reads as cheap.
4668    #[test]
4669    fn a_string_counts_as_a_node() {
4670        assert!(
4671            count_structural_chars(br#"["a","b","c"]"#) > count_structural_chars(br#"[1,1,1]"#),
4672            "strings were not counted, so an array of them looks free"
4673        );
4674    }
4675
4676    #[tokio::test]
4677    async fn the_sidecar_refuses_a_node_explosion_too() {
4678        let mut data =
4679            String::from(r#"{"ok":true,"data":{"records":[{"uri":"at://d/c/r","value":["#);
4680        for _ in 0..1_200_000 {
4681            data.push_str("{},");
4682        }
4683        data.push_str(r#"{}]}]}}"#);
4684        let base = crate::net::tests::serve_body(data.into_bytes()).await;
4685        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
4686        let err = client
4687            .list_records("did:plc:ewvi7nxzyoun6zhxrhs64oiz", "c", None, None)
4688            .await
4689            .expect_err("a node explosion reached the parser");
4690        assert!(
4691            format!("{err:#}").contains("structural characters"),
4692            "failed for the wrong reason: {err:#}"
4693        );
4694    }
4695
4696    // -- walk byte budget ---------------------------------------------------
4697
4698    /// What a parsed value really costs, counted independently of the code
4699    /// under test: every node occupies a `Value`, wherever it sits.
4700    fn node_count(v: &serde_json::Value) -> usize {
4701        1 + match v {
4702            serde_json::Value::Array(a) => a.iter().map(node_count).sum::<usize>(),
4703            serde_json::Value::Object(o) => o.values().map(node_count).sum::<usize>(),
4704            _ => 0,
4705        }
4706    }
4707
4708    fn record_of(value: serde_json::Value) -> RecordEntry {
4709        RecordEntry {
4710            uri: "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/c/3lab".to_string(),
4711            cid: Some("bafyreiabc123def456ghi789jkl012mno345pqr678stu901".to_string()),
4712            value,
4713        }
4714    }
4715
4716    /// **The estimate must never under-report, on any shape.**
4717    ///
4718    /// The version this replaces charged serialized length, which is accurate
4719    /// on prose-shaped records and 42x optimistic on the shapes an attacker
4720    /// picks. A bound that is only correct on benign input is not a bound.
4721    #[test]
4722    fn the_estimate_charges_every_node_at_least_what_a_parsed_value_costs() {
4723        let deep: serde_json::Value =
4724            serde_json::from_str(&format!("{}{}", "[".repeat(100), "]".repeat(100))).unwrap();
4725        let shapes: Vec<(&str, serde_json::Value)> = vec![
4726            ("100 nested empty arrays", deep),
4727            (
4728                "4096 empty arrays",
4729                serde_json::json!(vec![serde_json::json!([]); 4096]),
4730            ),
4731            ("4096 empty strings", serde_json::json!(vec![""; 4096])),
4732            (
4733                "4096 nulls",
4734                serde_json::json!(vec![serde_json::Value::Null; 4096]),
4735            ),
4736            ("4096 bools", serde_json::json!(vec![true; 4096])),
4737            ("4096 small numbers", serde_json::json!(vec![0; 4096])),
4738            (
4739                "object with short keys",
4740                serde_json::Value::Object(
4741                    (0..4096)
4742                        .map(|i| (format!("k{i}"), serde_json::json!([])))
4743                        .collect(),
4744                ),
4745            ),
4746            (
4747                "a realistic document",
4748                serde_json::json!({
4749                    "$type": "site.standard.document",
4750                    "title": "A post with a reasonably typical title",
4751                    "path": "/posts/one",
4752                    "publishedAt": "2026-07-11T09:30:00Z",
4753                    "textContent": "x".repeat(17_000),
4754                }),
4755            ),
4756        ];
4757        for (label, value) in shapes {
4758            let entry = record_of(value);
4759            let charged = approx_bytes(&entry);
4760            let floor = node_count(&entry.value) * std::mem::size_of::<serde_json::Value>();
4761            assert!(
4762                charged >= floor,
4763                "{label}: charged {charged} for {} nodes, which cannot cost less than {floor}",
4764                node_count(&entry.value)
4765            );
4766            let wire = serde_json::to_vec(&entry.value).unwrap().len();
4767            assert!(
4768                charged >= wire,
4769                "{label}: charged {charged}, under the {wire} bytes it takes on the wire alone"
4770            );
4771        }
4772    }
4773
4774    /// **Known answers, taken from a real allocator elsewhere.**
4775    ///
4776    /// The property above models `Value` nodes and nothing else, which is how an
4777    /// object-shaped under-charge of about half slipped past it: a
4778    /// `serde_json::Map` is a `BTreeMap` whose leaf is allocated whole, so the
4779    /// entries' own nodes are not the cost.
4780    ///
4781    /// **This test does not measure anything.** The two figures were obtained
4782    /// with a counting global allocator against the `serde_json` in this
4783    /// lockfile and are hardcoded here, because a global allocator is not
4784    /// something to install in the suite for one assertion. That makes this a
4785    /// tripwire for the *estimate* changing, not for the *real cost* changing: a
4786    /// dependency or toolchain bump that grows a map's true footprint leaves this
4787    /// green and the estimate quietly short again. Re-taking these numbers is the
4788    /// price of trusting them.
4789    #[test]
4790    fn the_estimate_covers_shapes_measured_against_a_real_allocator() {
4791        let many_small = serde_json::json!(vec![serde_json::json!({"a": 0}); 5000]);
4792        let mut deep = serde_json::json!({"a": 0});
4793        for _ in 0..99 {
4794            deep = serde_json::json!({ "a": deep });
4795        }
4796        for (label, value, measured) in [
4797            ("5000 one-key objects", many_small, 3_430_000usize),
4798            ("a 100-deep chain of one-key objects", deep, 63_350),
4799        ] {
4800            let charged = approx_bytes(&record_of(value));
4801            assert!(
4802                charged >= measured,
4803                "{label}: charged {charged} against {measured} bytes actually held"
4804            );
4805        }
4806    }
4807
4808    #[test]
4809    fn the_estimate_counts_the_uri_and_cid_too() {
4810        let bare = RecordEntry {
4811            uri: String::new(),
4812            cid: None,
4813            value: serde_json::json!(null),
4814        };
4815        let addressed = record_of(serde_json::json!(null));
4816        assert!(
4817            approx_bytes(&addressed) > approx_bytes(&bare),
4818            "a record's own identifiers are retained alongside its value"
4819        );
4820    }
4821
4822    #[test]
4823    fn the_budget_admits_a_page_that_exactly_fills_it() {
4824        let page = vec![record_of(serde_json::json!({"t": "x".repeat(1000)}))];
4825        let exact: usize = page.iter().map(approx_bytes).sum();
4826        assert!(
4827            ByteBudget::new(exact).admit(&page),
4828            "a page that exactly fits was refused; the fence-post is one byte out"
4829        );
4830        assert!(
4831            !ByteBudget::new(exact - 1).admit(&page),
4832            "a page one byte over the budget was admitted"
4833        );
4834    }
4835
4836    #[test]
4837    fn a_refused_page_leaves_the_running_total_alone() {
4838        let small = vec![record_of(serde_json::json!({"t": "x".repeat(100)}))];
4839        let huge = vec![record_of(serde_json::json!({"t": "x".repeat(100_000)}))];
4840        let cost: usize = small.iter().map(approx_bytes).sum();
4841        let mut budget = ByteBudget::new(cost * 3);
4842
4843        assert!(budget.admit(&small), "the first page fits");
4844        let after_one = budget.used();
4845        assert!(after_one > 0, "an admitted page must be charged");
4846
4847        assert!(!budget.admit(&huge), "the oversized page must be refused");
4848        assert_eq!(
4849            budget.used(),
4850            after_one,
4851            "a refused page moved the total — either charged, or reset"
4852        );
4853        assert!(
4854            budget.admit(&small),
4855            "the walk could not continue against the total it had before the refusal"
4856        );
4857    }
4858
4859    /// Build `pages` responses, each holding one record of about `bytes`, each
4860    /// pointing at the next. Returns the base URL and what one page costs.
4861    ///
4862    /// **Pages that differ is the whole point.** A walk served the same body
4863    /// twice stops on its repeated-cursor guard, so every test built on the
4864    /// fixed-body server refuses on page one and never exercises accumulation
4865    /// at all — which is how a per-page budget once passed a whole suite.
4866    pub(crate) fn paged_bodies(
4867        pages: usize,
4868        bytes: usize,
4869        envelope: bool,
4870    ) -> (Vec<Vec<u8>>, usize) {
4871        let record = |i: usize| {
4872            serde_json::json!({
4873                "uri": format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/c/3lab{i}"),
4874                "cid": "bafyreiabc123def456ghi789jkl012mno345pqr678stu901",
4875                "value": { "t": "x".repeat(bytes) }
4876            })
4877        };
4878        let bodies = (0..pages)
4879            .map(|i| {
4880                let mut page = serde_json::json!({ "records": [record(i)] });
4881                if i + 1 < pages {
4882                    page["cursor"] = serde_json::json!(format!("p{}", i + 1));
4883                }
4884                if envelope {
4885                    page = serde_json::json!({ "ok": true, "data": page });
4886                }
4887                page.to_string().into_bytes()
4888            })
4889            .collect();
4890        let entry: RecordEntry = serde_json::from_value(record(0)).unwrap();
4891        (bodies, approx_bytes(&entry))
4892    }
4893
4894    async fn host_for(bodies: Vec<Vec<u8>>, host: &str) -> (String, u16) {
4895        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
4896        let port: u16 = base
4897            .trim_end_matches('/')
4898            .rsplit(':')
4899            .next()
4900            .unwrap()
4901            .parse()
4902            .unwrap();
4903        crate::net::test_host_override(host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
4904        (format!("http://{host}:{port}"), port)
4905    }
4906
4907    // ---- #177: one malformed envelope ---------------------------------------
4908
4909    /// A record with no `uri`, which is what #177 probed on `main`.
4910    fn malformed_page(cursor: Option<&str>) -> Vec<u8> {
4911        let mut page = serde_json::json!({
4912            "records": [
4913                { "uri": "at://did:plc:x/c/3labGOOD", "value": {} },
4914                { "cid": "bafy", "value": {} },
4915            ]
4916        });
4917        if let Some(c) = cursor {
4918            page["cursor"] = serde_json::json!(c);
4919        }
4920        page.to_string().into_bytes()
4921    }
4922
4923    #[test]
4924    fn one_malformed_envelope_is_counted_not_fatal_to_the_page() {
4925        let page = parse_list_records(&malformed_page(None))
4926            .expect("one malformed envelope failed the whole page");
4927        assert_eq!(page.records.len(), 1, "the good record was not kept");
4928        assert_eq!(page.records[0].uri, "at://did:plc:x/c/3labGOOD");
4929        assert_eq!(page.malformed, 1, "the malformed record was not counted");
4930    }
4931
4932    /// The publications listing in `standard_site::fetch` reads a stranger's
4933    /// repo through this walk: it skips, counts, and keeps paging.
4934    #[tokio::test]
4935    async fn the_skipping_walk_skips_malformed_records_and_keeps_paging() {
4936        let only_bad = serde_json::json!({
4937            "records": [{ "cid": "bafy", "value": {} }], "cursor": "p1"
4938        })
4939        .to_string()
4940        .into_bytes();
4941        let bodies = vec![
4942            only_bad,
4943            malformed_page(Some("p2")),
4944            serde_json::json!({ "records": [] })
4945                .to_string()
4946                .into_bytes(),
4947        ];
4948        let (base, _) = host_for(bodies, "skipping-walk-malformed.test").await;
4949        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
4950        let (records, skipped) = client
4951            .list_all_records_skipping_within("c", &mut ByteBudget::new(MAX_LIST_BYTES))
4952            .await
4953            .expect("a malformed record failed a stranger's walk");
4954        assert_eq!(
4955            records.len(),
4956            1,
4957            "the good record behind the bad page was lost"
4958        );
4959        assert_eq!(skipped, 2, "skipped records were not counted");
4960    }
4961
4962    /// Pages of nothing but junk records, each ~`junk` bytes, with fresh
4963    /// cursors, so only the budget can stop the walk.
4964    fn junk_pages(n: usize, junk: usize) -> Vec<Vec<u8>> {
4965        (0..n)
4966            .map(|i| {
4967                serde_json::json!({
4968                    "records": [{ "cid": "bafy", "value": "x".repeat(junk) }],
4969                    "cursor": format!("p{}", i + 1),
4970                })
4971                .to_string()
4972                .into_bytes()
4973            })
4974            .collect()
4975    }
4976
4977    /// Review of #224: skipped records were never charged, so a stranger's
4978    /// repo serving junk pages walked all MAX_LIST_PAGES of them — gigabytes
4979    /// per poll — under a budget meant to stop at 128 MiB.
4980    #[tokio::test]
4981    async fn skipped_records_are_charged_against_the_budget() {
4982        let (base, _) = host_for(junk_pages(40, 256 * 1024), "junk-recent.test").await;
4983        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
4984        let mut budget = ByteBudget::new(1024 * 1024);
4985        let walk = client
4986            .list_recent_matching_within("c", 100, &mut budget, 25, |_| true)
4987            .await
4988            .unwrap();
4989        assert!(
4990            !walk.complete,
4991            "40 pages of junk were walked to the end under a 1 MiB budget"
4992        );
4993
4994        let (base, _) = host_for(junk_pages(40, 256 * 1024), "junk-skipping.test").await;
4995        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
4996        client
4997            .list_all_records_skipping_within("c", &mut ByteBudget::new(1024 * 1024))
4998            .await
4999            .expect_err("40 pages of junk were walked to the end under a 1 MiB budget");
5000    }
5001
5002    /// Review of #224: a page that skipped a record was charged its whole
5003    /// wire size AND its good records again, so a large publication near the
5004    /// budget failed only because one tiny malformed record sat beside it.
5005    #[tokio::test]
5006    async fn a_skipped_record_does_not_double_charge_its_page() {
5007        let page = |i: usize, with_bad: bool| {
5008            let mut records = vec![serde_json::json!({
5009                "uri": format!("at://did:plc:x/c/3lab{i}"), "value": "v".repeat(100_000)
5010            })];
5011            if with_bad {
5012                records.push(serde_json::json!({ "cid": "b", "value": {} }));
5013            }
5014            let mut body = serde_json::json!({ "records": records });
5015            if i < 4 {
5016                body["cursor"] = serde_json::json!(format!("p{}", i + 1));
5017            }
5018            body.to_string().into_bytes()
5019        };
5020        // A budget that holds the five clean pages with room to spare, and
5021        // less than twice that.
5022        let clean: Vec<_> = (0..5).map(|i| page(i, false)).collect();
5023        let (base, _) = host_for(clean, "double-charge-clean.test").await;
5024        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5025        let mut budget = ByteBudget::new(MAX_LIST_BYTES);
5026        let (clean_records, _) = client
5027            .list_all_records_skipping_within("c", &mut budget)
5028            .await
5029            .unwrap();
5030        let fits = budget.used() + budget.used() / 2;
5031
5032        let mixed: Vec<_> = (0..5).map(|i| page(i, true)).collect();
5033        let (base, _) = host_for(mixed, "double-charge-mixed.test").await;
5034        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5035        let (records, skipped) = client
5036            .list_all_records_skipping_within("c", &mut ByteBudget::new(fits))
5037            .await
5038            .expect("one tiny malformed record per page failed a walk that fits");
5039        assert_eq!(records.len(), clean_records.len());
5040        assert_eq!(skipped, 5);
5041
5042        // The documents walk, the same way: a page that skipped a record pays
5043        // its wire size once, not that and its kept records again.
5044        let mixed: Vec<_> = (0..5).map(|i| page(i, true)).collect();
5045        let (base, _) = host_for(mixed, "double-charge-recent.test").await;
5046        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5047        let walk = client
5048            .list_recent_matching_within("c", 100, &mut ByteBudget::new(fits), 25, |_| true)
5049            .await
5050            .unwrap();
5051        assert!(
5052            walk.complete,
5053            "one tiny malformed record per page cut short a walk that fits"
5054        );
5055        assert_eq!(walk.records.len(), clean_records.len());
5056    }
5057
5058    /// Third review of #224: charging a page that skipped a record only its
5059    /// wire size let a walk retain what it never paid for — a parsed record
5060    /// can hold up to 42x its wire size. Pages of dense `[[],[],…]` values, each
5061    /// with one malformed record beside them, retained 18x the budget.
5062    #[tokio::test]
5063    async fn a_page_that_skipped_a_record_still_pays_for_what_it_keeps() {
5064        let dense = format!("[{}]", vec!["[]"; 2000].join(","));
5065        let pages = |with_bad: bool| -> Vec<Vec<u8>> {
5066            (0..3)
5067                .map(|i| {
5068                    let mut records: Vec<String> = (0..99)
5069                        .map(|r| {
5070                            format!(r#"{{"uri":"at://did:plc:x/c/3l{i}x{r}","value":{dense}}}"#)
5071                        })
5072                        .collect();
5073                    if with_bad {
5074                        records.push(r#"{"cid":"b","value":{}}"#.to_string());
5075                    }
5076                    let cursor = if i < 2 {
5077                        format!(r#","cursor":"p{}""#, i + 1)
5078                    } else {
5079                        String::new()
5080                    };
5081                    format!(r#"{{"records":[{}]{cursor}}}"#, records.join(",")).into_bytes()
5082                })
5083                .collect()
5084        };
5085        const BUDGET: usize = 4 * 1024 * 1024;
5086        // Control: without the malformed records the walk is refused.
5087        let (base, _) = host_for(pages(false), "retained-control.test").await;
5088        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5089        client
5090            .list_all_records_skipping_within("c", &mut ByteBudget::new(BUDGET))
5091            .await
5092            .expect_err("control: the clean pages fit a budget they exceed");
5093
5094        let (base, _) = host_for(pages(true), "retained-skipping.test").await;
5095        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5096        client
5097            .list_all_records_skipping_within("c", &mut ByteBudget::new(BUDGET))
5098            .await
5099            .expect_err("one malformed record per page bought pages the budget refuses");
5100
5101        // The documents walk, with pages that each FIT the remaining budget
5102        // but together exceed it: only charging what the kept records retain
5103        // can stop it. (Pages each bigger than the budget are stopped by the
5104        // transient check alone and prove nothing about the charge.)
5105        let small_dense = format!("[{}]", vec!["[]"; 120].join(","));
5106        let fitting_pages: Vec<Vec<u8>> = (0..6)
5107            .map(|i| {
5108                let mut records: Vec<String> = (0..99)
5109                    .map(|r| {
5110                        format!(r#"{{"uri":"at://did:plc:x/c/3m{i}x{r}","value":{small_dense}}}"#)
5111                    })
5112                    .collect();
5113                records.push(r#"{"cid":"b","value":{}}"#.to_string());
5114                let cursor = if i < 5 {
5115                    format!(r#","cursor":"q{}""#, i + 1)
5116                } else {
5117                    String::new()
5118                };
5119                format!(r#"{{"records":[{}]{cursor}}}"#, records.join(",")).into_bytes()
5120            })
5121            .collect();
5122        let (base, _) = host_for(fitting_pages, "retained-recent.test").await;
5123        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5124        let walk = client
5125            .list_recent_matching_within("c", 10_000, &mut ByteBudget::new(BUDGET), 100, |_| true)
5126            .await
5127            .unwrap();
5128        let retained: usize = walk.records.iter().map(approx_bytes).sum();
5129        assert!(
5130            retained <= BUDGET,
5131            "retained {retained} under a {BUDGET}-byte budget"
5132        );
5133        assert!(
5134            !walk.complete,
5135            "pages totalling more than the budget were all kept"
5136        );
5137    }
5138
5139    /// The reader's own repo: this walk feeds `replace_sub_refs`, so skipping
5140    /// would drop a subscription silently. It refuses, by type.
5141    #[tokio::test]
5142    async fn the_own_repo_walk_refuses_a_page_with_a_malformed_record() {
5143        let (base, _) = host_for(vec![malformed_page(None)], "own-repo-malformed.test").await;
5144        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5145        let err = client
5146            .list_all_records("c")
5147            .await
5148            .expect_err("a page with a malformed record was accepted");
5149        let refused = err
5150            .downcast_ref::<MalformedRecords>()
5151            .unwrap_or_else(|| panic!("refused for the wrong reason: {err:#}"));
5152        assert_eq!(refused.count, 1);
5153        assert_eq!(refused.collection, "c");
5154    }
5155
5156    #[tokio::test]
5157    async fn the_sidecar_walk_refuses_a_page_with_a_malformed_record() {
5158        let body = serde_json::json!({
5159            "ok": true,
5160            "data": { "records": [
5161                { "uri": "at://did:plc:x/c/3labGOOD", "value": {} },
5162                { "cid": "bafy", "value": {} },
5163            ]}
5164        })
5165        .to_string()
5166        .into_bytes();
5167        let base = crate::net::tests::serve_bodies_in_sequence(vec![body]).await;
5168        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5169        let err = client
5170            .list_all_records("did:plc:x", "c")
5171            .await
5172            .expect_err("a page with a malformed record was accepted");
5173        assert!(
5174            err.downcast_ref::<MalformedRecords>().is_some(),
5175            "refused for the wrong reason: {err:#}"
5176        );
5177    }
5178
5179    /// A stranger's publication: skipping is right here, and the walk must keep
5180    /// paging past a page whose ONLY records were malformed — that page is not
5181    /// the end of the collection.
5182    #[tokio::test]
5183    async fn a_publication_walk_skips_malformed_records_and_keeps_paging() {
5184        let only_bad = serde_json::json!({
5185            "records": [{ "cid": "bafy", "value": {} }], "cursor": "p1"
5186        })
5187        .to_string()
5188        .into_bytes();
5189        let bodies = vec![
5190            only_bad,
5191            malformed_page(Some("p2")),
5192            serde_json::json!({ "records": [] })
5193                .to_string()
5194                .into_bytes(),
5195        ];
5196        let (base, _) = host_for(bodies, "publication-malformed.test").await;
5197        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5198        let mut budget = ByteBudget::new(MAX_LIST_BYTES);
5199        let walk = client
5200            .list_recent_matching_within("c", 100, &mut budget, 25, |_| true)
5201            .await
5202            .expect("a malformed record failed a stranger's publication walk");
5203        assert_eq!(
5204            walk.records.len(),
5205            1,
5206            "the good record behind the bad page was lost"
5207        );
5208        assert_eq!(walk.malformed, 2, "skipped records were not counted");
5209        assert!(
5210            walk.complete,
5211            "the walk stopped at a page of only malformed records"
5212        );
5213    }
5214
5215    /// **The budget is spent across pages, not reset by each one.**
5216    ///
5217    /// The single test this project most needed and did not have. Without it,
5218    /// moving the budget's construction inside the page loop — making the cap
5219    /// 200x weaker and effectively inert — passed every test in the suite.
5220    #[tokio::test]
5221    async fn a_refusing_walk_spends_its_budget_across_pages() {
5222        let (bodies, per_page) = paged_bodies(3, 4096, false);
5223        let (base, _) = host_for(bodies, "budget-accumulate.test").await;
5224        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5225
5226        let err = client
5227            .list_all_records_within("c", &mut ByteBudget::new(per_page * 2))
5228            .await
5229            .expect_err("three pages cannot fit in a two-page budget");
5230        let msg = format!("{err:#}");
5231        assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
5232        assert!(
5233            msg.contains("2 held"),
5234            "the walk did not keep exactly the two pages that fit: {msg}"
5235        );
5236    }
5237
5238    #[tokio::test]
5239    async fn a_truncating_walk_keeps_the_pages_that_fit() {
5240        let (bodies, per_page) = paged_bodies(3, 4096, false);
5241        let (base, _) = host_for(bodies, "budget-accumulate-trunc.test").await;
5242        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5243
5244        let walk = client
5245            .list_recent_matching_within("c", 100, &mut ByteBudget::new(per_page * 2), 100, |_| {
5246                true
5247            })
5248            .await
5249            .expect("an additive walk truncates rather than failing");
5250        assert_eq!(
5251            walk.records.len(),
5252            2,
5253            "the pages that fit were not kept, or the refused one was"
5254        );
5255        assert!(
5256            !walk.complete,
5257            "a walk stopped by the budget called itself complete"
5258        );
5259    }
5260
5261    /// **The budget must not bind before the record cap does, with room spare.**
5262    ///
5263    /// The walks that carry `MAX_LIST_RECORDS` REFUSE when a bound is hit, and a
5264    /// refusal drops the reader into `resolve_subscriptions`' fail-closed branch
5265    /// — so an account near the record cap would serve a stale projection on
5266    /// every poll, forever. The figure quoted in `MAX_LIST_BYTES`'s own comment
5267    /// is this calculation, and a review caught that figure being wrong by a
5268    /// factor of two because nothing computed it. This does.
5269    ///
5270    /// Double, not merely under: the margin is what stops a slightly longer
5271    /// title or one more optional field from turning a working account into a
5272    /// permanently failing one.
5273    #[test]
5274    fn a_full_subscription_repo_fits_the_budget_twice_over() {
5275        let record = record_of(serde_json::json!({
5276            "$type": "community.lexicon.rss.subscription",
5277            "url": "https://example.com/blog/feed.xml",
5278            "title": "Some Blog With A Longish Name",
5279            "siteUrl": "https://example.com/blog",
5280            "createdAt": "2026-07-11T09:30:00Z",
5281            "folder": "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/community.lexicon.rss.folder/3lab999",
5282            "fetchHint": "hourly",
5283        }));
5284        let per_record = approx_bytes(&record);
5285        let full_repo = per_record * MAX_LIST_RECORDS;
5286        assert!(
5287            full_repo * 2 <= MAX_LIST_BYTES,
5288            "a full repo charges {per_record} B x {MAX_LIST_RECORDS} = {} MB against a {} MB \
5289             budget — too close for a walk whose verdict is a refusal",
5290            full_repo / (1024 * 1024),
5291            MAX_LIST_BYTES / (1024 * 1024)
5292        );
5293    }
5294
5295    /// **Running out of pages is a refusal, not a short answer.**
5296    ///
5297    /// The three refusing walks fell out of `for _ in 0..MAX_LIST_PAGES` into a
5298    /// bare `Ok(out)`, so a repo bigger than the page budget returned a truncated
5299    /// list that looks exactly like a complete one. `resolve_subscriptions` needs
5300    /// an `Err` to take its fail-closed branch; given `Ok` it hands the short list
5301    /// to `replace_sub_refs`, which DELETEs the reader's whole `sub_ref`
5302    /// projection and reinserts only what it was given. Everything past the cap
5303    /// is gone from their account, on an ordinary poll, with no attacker.
5304    ///
5305    /// `extend_bounded`'s refusal cannot catch this: `MAX_LIST_PAGES` x the 100
5306    /// records we ask for is exactly `MAX_LIST_RECORDS`, so against any server
5307    /// that honours `limit` the page budget runs out first, every time.
5308    #[tokio::test]
5309    async fn a_walk_that_runs_out_of_pages_refuses_rather_than_truncating() {
5310        // One more page than the budget, every page still offering a cursor.
5311        let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
5312            .map(|i| {
5313                serde_json::json!({
5314                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
5315                    "cursor": format!("p{}", i + 1),
5316                })
5317                .to_string()
5318                .into_bytes()
5319            })
5320            .collect();
5321        let (base, _) = host_for(bodies, "pages-exhausted.test").await;
5322        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5323
5324        let err = client
5325            .list_all_records("c")
5326            .await
5327            .expect_err("a truncated list was returned as a complete one");
5328        let msg = format!("{err:#}");
5329        assert!(
5330            msg.contains("did not finish"),
5331            "failed for the wrong reason: {msg}"
5332        );
5333    }
5334
5335    /// **A walk that finishes cleanly across several pages still returns `Ok`.**
5336    ///
5337    /// The refusal's dangerous direction. Removing the flag's reset makes *every*
5338    /// multi-page walk refuse, which puts a reader with more than one page of
5339    /// records permanently into the fail-closed branch — and a review found that
5340    /// mutation surviving on the sidecar walk, which is the default backend,
5341    /// because nothing walked it to a clean finish and asserted success.
5342    #[tokio::test]
5343    async fn the_sidecar_walk_that_finishes_cleanly_returns_the_records() {
5344        let mut bodies: Vec<Vec<u8>> = (0..3)
5345            .map(|i| {
5346                serde_json::json!({
5347                    "ok": true,
5348                    "data": {
5349                        "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
5350                        "cursor": format!("p{}", i + 1),
5351                    }
5352                })
5353                .to_string()
5354                .into_bytes()
5355            })
5356            .collect();
5357        // **The terminator CARRIES a record.** Ending on an empty page left a
5358        // second mutation alive: drop the last page's records and the assertion
5359        // below still counts three, because the last page had none to drop. A
5360        // real PDS ends on a partial page, and that page's records are the ones
5361        // an off-by-one loses.
5362        bodies.push(
5363            serde_json::json!({
5364                "ok": true,
5365                "data": { "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }] }
5366            })
5367            .to_string()
5368            .into_bytes(),
5369        );
5370        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
5371        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5372        let records = client
5373            .list_all_records("did:plc:ewvi7nxzyoun6zhxrhs64oiz", "c")
5374            .await
5375            .expect("a walk that ran out of records is not a short list");
5376        assert_eq!(
5377            records.len(),
5378            4,
5379            "the pages that were served were not all kept"
5380        );
5381        assert!(
5382            records.iter().any(|r| r.uri.ends_with("3labLAST")),
5383            "the LAST page's records were dropped — the walk kept the right \
5384             count only because every page held one: {:?}",
5385            records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
5386        );
5387    }
5388
5389    /// **The page cap is pinned exactly, not to within one.**
5390    ///
5391    /// `the_sidecar_walk_that_runs_out_of_pages_refuses` serves
5392    /// `MAX_LIST_PAGES + 1` pages, so a budget one page SHORT refuses too and
5393    /// that mutation survives it. A walk whose last allowed request is the
5394    /// terminating one must come back `Ok` — which fails the moment the loop
5395    /// allows one page fewer, and is the direction that costs a reader their
5396    /// subscriptions.
5397    #[tokio::test]
5398    async fn a_sidecar_walk_that_terminates_on_its_last_allowed_page_succeeds() {
5399        let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
5400            .map(|i| {
5401                serde_json::json!({
5402                    "ok": true,
5403                    "data": {
5404                        "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
5405                        "cursor": format!("p{}", i + 1),
5406                    }
5407                })
5408                .to_string()
5409                .into_bytes()
5410            })
5411            .collect();
5412        // Request number `MAX_LIST_PAGES` — the last the loop allows — is the one
5413        // that terminates, and it carries a record of its own.
5414        bodies.push(
5415            serde_json::json!({
5416                "ok": true,
5417                "data": { "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }] }
5418            })
5419            .to_string()
5420            .into_bytes(),
5421        );
5422        assert_eq!(bodies.len(), MAX_LIST_PAGES);
5423        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
5424        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5425
5426        let records = client
5427            .list_all_records("did:plc:ewvi7nxzyoun6zhxrhs64oiz", "c")
5428            .await
5429            .expect("a walk that terminated inside its budget is not a short list");
5430        assert_eq!(
5431            records.len(),
5432            MAX_LIST_PAGES,
5433            "a walk that used its whole page budget and finished lost records",
5434        );
5435    }
5436
5437    /// **The direct walk's page cap, pinned exactly.**
5438    ///
5439    /// Twin of `a_sidecar_walk_that_terminates_on_its_last_allowed_page_succeeds`
5440    /// for the anonymous client. Verified needed: with only the `+ 1` refusal test
5441    /// above, `for _ in 0..MAX_LIST_PAGES - 1` left all 914 tests passing.
5442    #[tokio::test]
5443    async fn a_direct_walk_that_terminates_on_its_last_allowed_page_succeeds() {
5444        let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
5445            .map(|i| {
5446                serde_json::json!({
5447                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
5448                    "cursor": format!("p{}", i + 1),
5449                })
5450                .to_string()
5451                .into_bytes()
5452            })
5453            .collect();
5454        bodies.push(
5455            serde_json::json!({
5456                "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
5457            })
5458            .to_string()
5459            .into_bytes(),
5460        );
5461        assert_eq!(bodies.len(), MAX_LIST_PAGES);
5462        let (base, _) = host_for(bodies, "last-allowed-page.test").await;
5463        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5464
5465        let records = client
5466            .list_all_records("c")
5467            .await
5468            .expect("a walk that terminated inside its budget is not a short list");
5469        assert_eq!(
5470            records.len(),
5471            MAX_LIST_PAGES,
5472            "a walk that used its whole page budget and finished lost records",
5473        );
5474        assert!(
5475            records.iter().any(|r| r.uri.ends_with("3labLAST")),
5476            "the LAST page's records were dropped",
5477        );
5478    }
5479
5480    /// **The TRUNCATING walk's page cap, pinned exactly — it reports completeness
5481    /// rather than refusing, so an off-by-one here is a silent short read.**
5482    ///
5483    /// A publication whose archive needs exactly the page budget to exhaust is
5484    /// `complete`; one page fewer makes it `complete = false`, which
5485    /// `store_publication` treats as a partial read. Verified needed:
5486    /// `for _ in 0..MAX_LIST_PAGES - 1` on this walk left all 914 tests passing.
5487    #[tokio::test]
5488    async fn a_truncating_walk_that_exhausts_on_its_last_allowed_page_is_complete() {
5489        let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
5490            .map(|i| {
5491                serde_json::json!({
5492                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
5493                    "cursor": format!("p{}", i + 1),
5494                })
5495                .to_string()
5496                .into_bytes()
5497            })
5498            .collect();
5499        bodies.push(
5500            serde_json::json!({
5501                "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
5502            })
5503            .to_string()
5504            .into_bytes(),
5505        );
5506        assert_eq!(bodies.len(), MAX_LIST_PAGES);
5507        let (base, _) = host_for(bodies, "last-allowed-page-truncating.test").await;
5508        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5509
5510        // `max_records` well above what is served, so the cap under test is the
5511        // PAGE budget and not the record one.
5512        let walk = client
5513            .list_recent_matching("c", MAX_LIST_PAGES * 10, 1, |_| true)
5514            .await
5515            .expect("walk failed");
5516        assert_eq!(
5517            walk.records.len(),
5518            MAX_LIST_PAGES,
5519            "a walk that used its whole page budget and exhausted the collection \
5520             lost records",
5521        );
5522        assert!(
5523            walk.complete,
5524            "a collection that ran out on the last allowed page was reported as a \
5525             partial read, which is a starvation warning for a complete archive",
5526        );
5527    }
5528
5529    /// The sidecar walk refuses a short list too — and it is the default backend.
5530    #[tokio::test]
5531    async fn the_sidecar_walk_that_runs_out_of_pages_refuses() {
5532        let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
5533            .map(|i| {
5534                serde_json::json!({
5535                    "ok": true,
5536                    "data": {
5537                        "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
5538                        "cursor": format!("p{}", i + 1),
5539                    }
5540                })
5541                .to_string()
5542                .into_bytes()
5543            })
5544            .collect();
5545        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
5546        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5547        let err = client
5548            .list_all_records("did:plc:ewvi7nxzyoun6zhxrhs64oiz", "c")
5549            .await
5550            .expect_err("a truncated list was returned as a complete one");
5551        assert!(
5552            format!("{err:#}").contains("did not finish"),
5553            "failed for the wrong reason: {err:#}"
5554        );
5555    }
5556
5557    /// **A budget passed to two walks is spent by both of them.**
5558    ///
5559    /// The reason it is passed rather than constructed: a publication read runs
5560    /// a second walk while still holding the first's records, so two independent
5561    /// ceilings let one read hold twice the bound. Here the first walk spends
5562    /// the budget and the second finds it spent.
5563    #[tokio::test]
5564    async fn two_walks_sharing_a_budget_do_not_each_get_the_whole_of_it() {
5565        let (bodies, per_page) = paged_bodies(4, 4096, false);
5566        let (base, _) = host_for(bodies, "budget-shared.test").await;
5567        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5568        let mut budget = ByteBudget::new(per_page * 3);
5569
5570        let err = client
5571            .list_all_records_within("c", &mut budget)
5572            .await
5573            .expect_err("four pages cannot fit a three-page budget");
5574        assert!(format!("{err:#}").contains("3 held"), "{err:#}");
5575
5576        // Same budget, nothing left in it.
5577        let err = client
5578            .list_all_records_within("c", &mut budget)
5579            .await
5580            .expect_err("the second walk was handed a fresh ceiling");
5581        assert!(
5582            format!("{err:#}").contains("0 held"),
5583            "the second walk kept something out of an exhausted budget: {err:#}"
5584        );
5585    }
5586
5587    /// **A transient page has to fit what is LEFT of the budget.**
5588    ///
5589    /// Measuring it against the ceiling lets a walk that has already retained
5590    /// most of its budget hold a further ceiling's worth of page on top. The
5591    /// filter keeps nothing here, so the running total cannot stop the walk and
5592    /// only the remaining-budget comparison can.
5593    #[tokio::test]
5594    async fn a_transient_page_must_fit_what_is_left_not_the_ceiling() {
5595        let (bodies, per_page) = paged_bodies(3, 4096, false);
5596        let (base, _) = host_for(bodies, "budget-remaining.test").await;
5597        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5598
5599        // Ceiling of two and a half pages, two of them already spent.
5600        let mut budget = ByteBudget::new(per_page * 5 / 2);
5601        let spent = vec![
5602            record_of(serde_json::json!({ "t": "x".repeat(4096) })),
5603            record_of(serde_json::json!({ "t": "x".repeat(4096) })),
5604        ];
5605        assert!(budget.admit(&spent), "the pre-spend has to fit");
5606        assert!(
5607            budget.remaining() < per_page,
5608            "and has to leave less than a page"
5609        );
5610
5611        let walk = client
5612            .list_recent_matching_within("c", 100, &mut budget, 100, |_| false)
5613            .await
5614            .expect("an additive walk truncates rather than failing");
5615        assert!(
5616            !walk.complete,
5617            "a page larger than the remaining budget was walked past"
5618        );
5619    }
5620
5621    /// **A filter that keeps nothing must not let the walk run unbounded.**
5622    ///
5623    /// The running total charges what is kept, so a filter matching nothing
5624    /// charges zero and the total can never stop the walk. What it holds is
5625    /// another matter: each page is fully parsed before the filter sees it, and
5626    /// `read_capped`'s 8 MB bounds the wire, not the tree. Only the per-page
5627    /// bound stands between that and the box.
5628    #[tokio::test]
5629    async fn a_filter_that_keeps_nothing_still_cannot_outrun_the_budget() {
5630        let (bodies, per_page) = paged_bodies(3, 4096, false);
5631        let (base, _) = host_for(bodies, "budget-filtered.test").await;
5632        let client = PdsClient::anonymous(ssrf_test_client(), base, "did:plc:x");
5633
5634        let walk = client
5635            .list_recent_matching_within("c", 100, &mut ByteBudget::new(per_page / 2), 100, |_| {
5636                false
5637            })
5638            .await
5639            .expect("an additive walk truncates rather than failing");
5640        assert!(
5641            !walk.complete,
5642            "a page too large to hold was walked past because the filter dropped it"
5643        );
5644        assert!(walk.records.is_empty(), "the filter kept nothing");
5645    }
5646
5647    #[tokio::test]
5648    async fn the_sidecar_walk_spends_its_budget_across_pages() {
5649        let (bodies, per_page) = paged_bodies(3, 4096, true);
5650        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
5651        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5652
5653        let err = client
5654            .list_all_records_within(
5655                "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
5656                "app.feather.subscription",
5657                &mut ByteBudget::new(per_page * 2),
5658            )
5659            .await
5660            .expect_err("the sidecar walk was the one with no budget at all");
5661        let msg = format!("{err:#}");
5662        assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
5663        assert!(
5664            msg.contains("2 held"),
5665            "did not accumulate across pages: {msg}"
5666        );
5667    }
5668
5669    /// **Shapes that used to read as a healthy empty page.**
5670    ///
5671    /// Each of these was accepted by the `Value` route as `records: []`, and an
5672    /// empty page is not inert: `resolve_subscriptions` passes it to
5673    /// `replace_sub_refs`, which DELETEs the reader's projection and rewrites
5674    /// what it was handed. A page that is wrong in this direction costs them
5675    /// every feed.
5676    #[test]
5677    fn a_page_that_is_not_a_listing_is_never_read_as_an_empty_one() {
5678        for (label, body) in [
5679            (
5680                "a non-string error alongside records",
5681                &br#"{"error":404,"records":[]}"#[..],
5682            ),
5683            (
5684                "an object error alongside records",
5685                &br#"{"error":{"code":"x"},"records":[]}"#[..],
5686            ),
5687            (
5688                "a duplicated records key, the second one empty",
5689                &br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#[..],
5690            ),
5691            ("an explicit null records", &br#"{"records":null}"#[..]),
5692        ] {
5693            assert!(
5694                parse_list_records(body).is_err(),
5695                "{label} was read as a page"
5696            );
5697        }
5698    }
5699
5700    /// **A non-string `error` is an envelope, and is reported as one.**
5701    ///
5702    /// Typing the field as a `String` made these fail as "invalid type" — the
5703    /// wrong reason for the exact shape the guard exists for, which is the same
5704    /// looseness that once let the guard be deleted unnoticed. So the reason is
5705    /// asserted, not just the refusal.
5706    #[test]
5707    fn a_non_string_error_is_reported_as_an_envelope() {
5708        for body in [
5709            &br#"{"error":404,"records":[]}"#[..],
5710            &br#"{"error":{"code":"x"},"records":[]}"#[..],
5711            &br#"{"error":[],"records":[]}"#[..],
5712            &br#"{"error":true,"records":[]}"#[..],
5713        ] {
5714            let err = parse_list_records(body)
5715                .expect_err("a non-string error envelope was read as an empty page");
5716            assert!(
5717                format!("{err:#}").contains("error envelope"),
5718                "{} failed for the wrong reason: {err:#}",
5719                String::from_utf8_lossy(body)
5720            );
5721        }
5722        // An empty name IS an envelope, as it was before this work: the route
5723        // this replaced keyed on `as_str`, so `Some("")` bailed. Exempting it
5724        // was a loosening made on speculation about proxy conventions, and a
5725        // loosening in this direction is a page accepted that used to be
5726        // refused.
5727        for body in [
5728            &br#"{"error":"","records":[]}"#[..],
5729            // Not zero on the wire, but zero once read: an exemption keyed on
5730            // `as_f64` swallowed anything that underflows.
5731            &br#"{"error":1e-400,"records":[]}"#[..],
5732        ] {
5733            let err = parse_list_records(body).expect_err("this is an envelope");
5734            assert!(
5735                format!("{err:#}").contains("error envelope"),
5736                "{} failed for the wrong reason: {err:#}",
5737                String::from_utf8_lossy(body)
5738            );
5739        }
5740        // The four spellings of "no error". `null` is what an ordinary listing
5741        // carries; `false` and integer `0` are a proxy convention, and refusing
5742        // those would fail a good page outright.
5743        for body in [
5744            &br#"{"error":null,"records":[]}"#[..],
5745            &br#"{"error":false,"records":[]}"#[..],
5746            &br#"{"error":0,"records":[]}"#[..],
5747            &br#"{"records":[]}"#[..],
5748        ] {
5749            assert!(
5750                parse_list_records(body).is_ok(),
5751                "{} is not an error envelope",
5752                String::from_utf8_lossy(body)
5753            );
5754        }
5755    }
5756
5757    /// **A non-string `error` is named by its type, never by its contents.**
5758    ///
5759    /// Rendering the value would serialise the whole attacker-chosen subtree
5760    /// before truncating it, allocating a full extra copy of up to the body cap
5761    /// — in a change whose purpose is cutting peak allocation. The earlier
5762    /// version of this did exactly that and the comment claimed otherwise.
5763    #[test]
5764    fn a_structured_error_is_named_by_its_type_not_serialised() {
5765        let payload = "s".repeat(20_000);
5766        let body = format!(r#"{{"error":{{"deep":"{payload}"}},"records":[]}}"#);
5767        let err = parse_list_records(body.as_bytes()).expect_err("an envelope is a refusal");
5768        let msg = format!("{err:#}");
5769        assert!(
5770            !msg.contains("ssss"),
5771            "the error's contents reached the message: {} chars",
5772            msg.len()
5773        );
5774        assert!(
5775            msg.contains("non-string error: object"),
5776            "it should name the shape instead: {msg}"
5777        );
5778    }
5779
5780    /// `data` absent is not `data` empty, on the sidecar envelope too.
5781    ///
5782    /// `{"ok":true}` is what a proxy makes of an unexpected upstream body, and
5783    /// reading it as a page of zero records is the wipe this whole family of
5784    /// guards exists to prevent.
5785    #[tokio::test]
5786    async fn the_sidecar_refuses_an_envelope_with_no_data() {
5787        for body in [&br#"{"ok":true}"#[..], &br#"{"ok":true,"data":{}}"#[..]] {
5788            let base = crate::net::tests::serve_body(body.to_vec()).await;
5789            let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5790            let err = client
5791                .list_records(
5792                    "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
5793                    "app.feather.subscription",
5794                    None,
5795                    None,
5796                )
5797                .await
5798                .expect_err("an envelope without a listing was read as an empty page");
5799            assert!(
5800                format!("{err:#}").contains("no records"),
5801                "{} failed for the wrong reason: {err:#}",
5802                String::from_utf8_lossy(body)
5803            );
5804        }
5805    }
5806
5807    /// **The sidecar gets the duplicated-key refusal too.**
5808    ///
5809    /// It was the one client still reading a listing through a `Value`, where a
5810    /// repeated key resolves last-wins — so a body carrying a second, empty
5811    /// `records` array read as a successful empty page, and an empty page on this
5812    /// path is `replace_sub_refs` deleting every `sub_ref` the reader has. It is
5813    /// also the default backend, so it was the one that mattered most.
5814    #[tokio::test]
5815    async fn the_sidecar_refuses_a_duplicated_records_key() {
5816        let base = crate::net::tests::serve_body(
5817            br#"{"ok":true,"data":{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}}"#
5818                .to_vec(),
5819        )
5820        .await;
5821        let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5822        let err = client
5823            .list_records(
5824                "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
5825                "app.feather.subscription",
5826                None,
5827                None,
5828            )
5829            .await
5830            .expect_err("a duplicated records key was read as an empty page");
5831        assert!(
5832            format!("{err:#}").contains("duplicate"),
5833            "failed for the wrong reason: {err:#}"
5834        );
5835    }
5836
5837    /// A non-string `message` must not fail an otherwise good page.
5838    /// **A listing has to be an object.**
5839    ///
5840    /// serde's derived `Deserialize` takes a struct positionally too, so with
5841    /// every field defaulted `[null,null,[]]` bound `records` to an empty vector
5842    /// and read as a healthy page — and a body with no keys defeats the envelope
5843    /// guard and the duplicated-key refusal at the same time, because neither has
5844    /// anything to look at. Fourteen bytes, and `replace_sub_refs` deletes every
5845    /// feed the reader has.
5846    #[test]
5847    fn a_listing_that_is_not_an_object_is_not_a_page() {
5848        for body in [
5849            &b"[null,null,[]]"[..],
5850            &b"[null,null,[],null]"[..],
5851            &br#"[null,null,[{"uri":"at://d/c/r","value":{}}],"c"]"#[..],
5852            &b"[]"[..],
5853            &br#""a string""#[..],
5854            &b"0"[..],
5855            &b"true"[..],
5856        ] {
5857            assert!(
5858                parse_list_records(body).is_err(),
5859                "{} was read as a page",
5860                String::from_utf8_lossy(body)
5861            );
5862        }
5863    }
5864
5865    /// **The rendering is bounded in bytes, whatever the input is made of.**
5866    ///
5867    /// Counting characters bounds nothing a log cares about: 120 astral-plane
5868    /// code points are 480 bytes. The invariant is on the output's byte length.
5869    #[test]
5870    fn a_truncated_message_is_bounded_in_bytes() {
5871        for (label, input) in [
5872            ("ascii", "e".repeat(50_000)),
5873            ("astral", "\u{1f600}".repeat(20_000)),
5874            (
5875                "mixed",
5876                format!("{}{}", "e".repeat(200), "\u{1f600}".repeat(200)),
5877            ),
5878            (
5879                "just over in bytes, just under in chars",
5880                "\u{1f600}".repeat(40),
5881            ),
5882        ] {
5883            let out = truncate_for_message(&input);
5884            assert!(
5885                out.len() <= 200,
5886                "{label}: rendered {} bytes from {} bytes of input",
5887                out.len(),
5888                input.len()
5889            );
5890        }
5891        // Short inputs pass through untouched.
5892        assert_eq!(truncate_for_message("Boom"), "Boom");
5893    }
5894
5895    /// **The sidecar has two envelope layers, and both are guards.**
5896    ///
5897    /// `page_from_body` covers the PDS's, which arrives inside `data`. The
5898    /// sidecar's own can say `ok:false` or carry its own `error` on a 200 while
5899    /// `data` still holds something that reads as a perfectly good empty page —
5900    /// and an empty page here is `replace_sub_refs` deleting every feed.
5901    #[tokio::test]
5902    async fn the_sidecar_refuses_its_own_error_envelope_on_a_2xx() {
5903        for body in [
5904            &br#"{"ok":false,"error":"ExpiredToken","data":{"records":[]}}"#[..],
5905            &br#"{"ok":true,"error":"ExpiredToken","data":{"records":[]}}"#[..],
5906            &br#"{"ok":false,"data":{"records":[]}}"#[..],
5907        ] {
5908            let base = crate::net::tests::serve_body(body.to_vec()).await;
5909            let client = SidecarClient::new(Client::new(), base.clone(), base, "secret");
5910            let err = client
5911                .list_records(
5912                    "did:plc:ewvi7nxzyoun6zhxrhs64oiz",
5913                    "app.feather.subscription",
5914                    None,
5915                    None,
5916                )
5917                .await
5918                .expect_err("the sidecar's own envelope was read as a page");
5919            let msg = format!("{err:#}");
5920            assert!(
5921                msg.contains("sidecar answered 2xx"),
5922                "{} failed for the wrong reason: {msg}",
5923                String::from_utf8_lossy(body)
5924            );
5925        }
5926    }
5927
5928    /// **An unknown field's contents are still validated.**
5929    ///
5930    /// `IgnoredAny` skips without validating, so a body that is not valid JSON at
5931    /// all read as a healthy empty page where the route this replaced refused it.
5932    #[test]
5933    fn an_unknown_field_holding_invalid_json_is_not_a_page() {
5934        for body in [
5935            &b"{\"records\":[],\"x\":\"\xff\xfe\"}"[..],
5936            &br#"{"records":[],"x":"\ud800"}"#[..],
5937        ] {
5938            assert!(
5939                parse_list_records(body).is_err(),
5940                "{} was read as a page",
5941                String::from_utf8_lossy(body)
5942            );
5943        }
5944    }
5945
5946    #[test]
5947    fn a_non_string_message_does_not_cost_the_page() {
5948        let page = parse_list_records(br#"{"records":[],"message":5,"cursor":"c"}"#)
5949            .expect("message carries no guard; typing it strictly failed whole listings");
5950        assert_eq!(page.cursor.as_deref(), Some("c"));
5951    }
5952
5953    /// The name that reaches the log is bounded, because the PDS chooses it.
5954    #[test]
5955    fn an_enormous_error_name_is_truncated_before_it_reaches_a_log() {
5956        let huge = "e".repeat(50_000);
5957        let body = format!(r#"{{"error":"{huge}","records":[]}}"#);
5958        let err = parse_list_records(body.as_bytes()).expect_err("an envelope is a refusal");
5959        let msg = format!("{err:#}");
5960        assert!(
5961            msg.len() < 400,
5962            "the error message carried {} bytes of attacker-chosen text",
5963            msg.len()
5964        );
5965        assert!(
5966            msg.contains("50000 bytes"),
5967            "it should say what it dropped: {msg}"
5968        );
5969
5970        // Astral-plane code points: the bound must hold in BYTES, because a log
5971        // line is bytes. Counting characters made this four times the stated cap.
5972        let wide = "\u{1f600}".repeat(20_000);
5973        let body = format!(r#"{{"error":"{wide}","message":"{wide}","records":[]}}"#);
5974        let err = parse_list_records(body.as_bytes()).expect_err("an envelope is a refusal");
5975        let msg = format!("{err:#}");
5976        assert!(
5977            msg.len() < 400,
5978            "a wide-character error rendered {} bytes",
5979            msg.len()
5980        );
5981    }
5982
5983    // -- TID rkeys ----------------------------------------------------------
5984
5985    #[test]
5986    fn tid_rkeys_are_13_char_s32_and_monotonic() {
5987        let mut gen = TidGenerator::new();
5988        let mut prev: Option<String> = None;
5989        for _ in 0..1000 {
5990            let tid = gen.next();
5991            assert_eq!(tid.len(), 13, "a TID is 13 s32 chars");
5992            assert!(
5993                tid.bytes().all(|b| S32_ALPHABET.contains(&b)),
5994                "TID {tid} uses only the s32 alphabet"
5995            );
5996            if let Some(p) = &prev {
5997                assert!(*p < tid, "TIDs must be strictly increasing ({p} < {tid})");
5998            }
5999            prev = Some(tid);
6000        }
6001    }
6002
6003    #[test]
6004    fn tid_rkeys_are_valid_atproto_record_keys() {
6005        // atproto rkey charset: [A-Za-z0-9._~:-], length 1..=512, not "."/"..".
6006        let mut gen = TidGenerator::new();
6007        let tid = gen.next();
6008        assert!(is_valid_rkey(&tid), "{tid:?}");
6009        assert!(tid
6010            .bytes()
6011            .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'.' | b'_' | b'~' | b':' | b'-')));
6012    }
6013
6014    #[test]
6015    fn tid_values_round_trip_through_the_decoder() {
6016        // The decoder is the inverse of the encoder across the whole range a
6017        // TID can hold, boundaries included.
6018        let max_tid = (0x001f_ffff_ffff_ffffu64 << 10) | 0x3ff;
6019        for v in [0u64, 1, 31, 32, 1023, 1024, 1_000_000, max_tid] {
6020            let encoded = encode_s32_tid(v);
6021            assert_eq!(
6022                decode_s32_tid(&encoded),
6023                Some(v),
6024                "{v} encoded to {encoded}, which did not decode back"
6025            );
6026        }
6027
6028        // **A round trip alone proves too little.** Encoder and decoder share
6029        // the alphabet, so swapping two of its symbols round-trips perfectly
6030        // and still reads every real record key wrong. These two are the
6031        // known answer: a record key from a real atproto repo, and the value
6032        // it holds, computed independently of this code.
6033        assert_eq!(
6034            decode_s32_tid("3jzfcijpj2z2a"),
6035            Some(1_728_652_679_052_295_174)
6036        );
6037        assert_eq!(encode_s32_tid(1_728_652_679_052_295_174), "3jzfcijpj2z2a");
6038        assert_eq!(
6039            decode_s32_tid("3jzfcijpj2z2a").map(|raw| raw >> 10),
6040            Some(1_688_137_381_887_007),
6041            "that key was written at 2023-06-30T15:03:01.887007Z"
6042        );
6043    }
6044
6045    #[test]
6046    fn the_first_tid_of_a_generator_decodes_to_the_microsecond_it_was_minted() {
6047        let micros = || {
6048            std::time::SystemTime::now()
6049                .duration_since(std::time::UNIX_EPOCH)
6050                .map(|d| d.as_micros() as u64)
6051                .unwrap_or(0)
6052        };
6053        // The FIRST `next()` only. `TidGenerator` bumps a TID to `last + 1`
6054        // to stay strictly increasing, and on a generator whose clock id is
6055        // already at its maximum that carry lands in the timestamp bits — so a
6056        // later TID can decode a microsecond or two past when it was really
6057        // minted. A fresh generator has `last: 0`, where the bump cannot fire.
6058        let before = micros();
6059        let tid = TidGenerator::new().next();
6060        let after = micros();
6061        let raw = decode_s32_tid(&tid).expect("a generated TID must decode");
6062        let minted = raw >> 10;
6063        assert!(
6064            (before..=after).contains(&minted),
6065            "TID {tid} decoded to {minted}, outside the {before}..={after} window it was minted in"
6066        );
6067    }
6068
6069    #[test]
6070    fn the_decoder_rejects_strings_that_are_not_13_char_s32_values() {
6071        for rkey in [
6072            "",               // empty
6073            "self",           // the common non-TID rkey
6074            "3jzfcijpj2z2",   // 12 chars: one short
6075            "3jzfcijpj2z2aa", // 14 chars: one long
6076            "3jzfcijpj2z2A",  // uppercase is outside the s32 alphabet
6077            "3jzfcijpj2z-a",  // a legal rkey character, but not an s32 one
6078            "3jzfcijpj2z2!",  // not a legal rkey character at all
6079            "c222222222222",  // decodes with bit 63 set: the reserved top bit
6080            "k222222222222",  // decodes past 64 bits entirely
6081            "zzzzzzzzzzzzz",  // the largest 13-char s32 string
6082        ] {
6083            assert_eq!(
6084                decode_s32_tid(rkey),
6085                None,
6086                "{rkey:?} is not a 13-character s32 value"
6087            );
6088        }
6089    }
6090
6091    /// **The window is not a slug detector, and this is what that costs.**
6092    ///
6093    /// A 13-character slug beginning `3` decodes into the last few years just
6094    /// as a record key does, and nothing in the string tells them apart. These
6095    /// are read as dates, and pinning that here is the honest alternative to a
6096    /// doc comment claiming otherwise. The damage is bounded: a wrong date is
6097    /// an ordinary past instant that ages, sweeps and is outranked normally.
6098    #[test]
6099    fn a_slug_that_decodes_inside_the_window_is_read_as_a_date() {
6100        for (slug, reads_as) in [
6101            ("3hoursinparis", "2020-11-24T08:17:26Z"),
6102            ("3ideasforjune", "2021-08-12T00:19:38Z"),
6103            ("3jokesaweekly", "2023-02-12T15:50:26Z"),
6104        ] {
6105            assert_eq!(
6106                tid_timestamp(slug).map(crate::feed::fmt_time),
6107                Some(reads_as.to_string()),
6108                "{slug} is indistinguishable from a record key written then"
6109            );
6110        }
6111    }
6112
6113    #[test]
6114    fn a_tid_minted_slightly_ahead_of_our_clock_is_still_believed() {
6115        let now = chrono::Utc::now();
6116        let of = |at: chrono::DateTime<chrono::Utc>| {
6117            encode_s32_tid((at.timestamp_micros() as u64) << 10)
6118        };
6119        assert!(
6120            tid_timestamp(&of(now + chrono::Duration::seconds(2))).is_some(),
6121            "a PDS two seconds fast must not leave a fresh document undated"
6122        );
6123        assert_eq!(
6124            tid_timestamp(&of(now + chrono::Duration::hours(1))),
6125            None,
6126            "an hour ahead is a broken clock or a slug, not skew"
6127        );
6128    }
6129
6130    #[test]
6131    fn a_tid_timestamp_is_bounded_at_both_ends() {
6132        let now = chrono::Utc::now();
6133        let of = |micros: i64| encode_s32_tid((micros as u64) << 10);
6134
6135        // A TID minted now dates to now.
6136        let fresh = TidGenerator::new().next();
6137        let dated = tid_timestamp(&fresh).expect("a freshly minted TID has a timestamp");
6138        assert!(
6139            (now - chrono::Duration::minutes(1)..=now + chrono::Duration::minutes(1))
6140                .contains(&dated),
6141            "{fresh} dated to {dated}, not to now ({now})"
6142        );
6143
6144        // Before atproto existed: not a date.
6145        assert_eq!(
6146            tid_timestamp(&of(TID_FLOOR_MICROS - 1)),
6147            None,
6148            "a TID predating atproto must not date an entry"
6149        );
6150        assert!(
6151            tid_timestamp(&of(TID_FLOOR_MICROS)).is_some(),
6152            "the floor itself is a real instant"
6153        );
6154
6155        // In the future: not a date. A slug of 13 s32 characters lands here,
6156        // which is the case this bound exists for.
6157        let far_future = (now + chrono::Duration::days(365)).timestamp_micros();
6158        assert_eq!(
6159            tid_timestamp(&of(far_future)),
6160            None,
6161            "a TID from the future must not date an entry"
6162        );
6163        assert_eq!(
6164            tid_timestamp("abcdefghijklm"),
6165            None,
6166            "a 13-character slug decodes to the year 2192; it is not a date"
6167        );
6168    }
6169
6170    #[test]
6171    fn s32_encoding_is_ascending_for_ascending_values() {
6172        // The whole point of s32: numeric order == lexicographic string order.
6173        assert!(encode_s32_tid(1) < encode_s32_tid(2));
6174        assert!(encode_s32_tid(31) < encode_s32_tid(32));
6175        assert!(encode_s32_tid(1_000_000) < encode_s32_tid(1_000_001));
6176        // Ordering holds all the way to the largest real TID value (a 53-bit
6177        // microsecond timestamp shifted into bits 63..10, plus the clock id).
6178        let max_tid = (0x001f_ffff_ffff_ffffu64 << 10) | 0x3ff;
6179        assert!(encode_s32_tid(max_tid - 1) < encode_s32_tid(max_tid));
6180    }
6181    /// **Exceeding the record cap is an ERROR, not a silent truncation.**
6182    ///
6183    /// The page cap bounds how many requests a walk makes; it bounds the
6184    /// accumulated memory only if the server honours `limit=100`, and a host we
6185    /// did not choose has no obligation to. A review measured an 8 MB page
6186    /// holding ~95 000 minimal records and retaining 23 MB as
6187    /// `Vec&lt;RecordEntry&gt;` — 200 such pages is gigabytes on a 512 MB box.
6188    ///
6189    /// Truncating instead would be worse than the OOM it prevents. The caller
6190    /// of the live walk is `resolve_subscriptions`, whose result feeds
6191    /// `replace_sub_refs` — a `DELETE` plus reinsert of exactly what it was
6192    /// handed. A short list there is not a short list, it is **revoked access**
6193    /// to the feeds that fell off the end. That is the failure PR #167 was
6194    /// closed for reintroducing, so this returns `Err` and lets the existing
6195    /// fail-closed branch serve the last-known projection.
6196    #[test]
6197    fn exceeding_the_record_cap_is_an_error_not_a_truncation() {
6198        let page = |n: usize| -> Vec<RecordEntry> {
6199            (0..n)
6200                .map(|i| RecordEntry {
6201                    uri: format!("at://did:plc:x/c/{i}"),
6202                    cid: None,
6203                    value: serde_json::Value::Null,
6204                })
6205                .collect()
6206        };
6207
6208        let mut out = page(90);
6209        let err = extend_bounded(&mut out, page(20), 100, "c")
6210            .expect_err("a page past the cap was accepted");
6211        let msg = format!("{err:#}");
6212        assert!(msg.contains("100"), "the cap is not named: {msg}");
6213        assert_eq!(
6214            out.len(),
6215            90,
6216            "the partial page was kept — a truncated list must not survive the error"
6217        );
6218    }
6219
6220    #[test]
6221    fn accumulating_within_the_cap_succeeds() {
6222        let page = |n: usize| -> Vec<RecordEntry> {
6223            (0..n)
6224                .map(|i| RecordEntry {
6225                    uri: format!("at://did:plc:x/c/{i}"),
6226                    cid: None,
6227                    value: serde_json::Value::Null,
6228                })
6229                .collect()
6230        };
6231        let mut out = Vec::new();
6232        extend_bounded(&mut out, page(60), 100, "c").unwrap();
6233        extend_bounded(&mut out, page(40), 100, "c").unwrap();
6234        assert_eq!(out.len(), 100, "exactly the cap must be allowed");
6235    }
6236
6237    // -- applyWrites chunking (#240) ----------------------------------------
6238
6239    /// Every `applyWrites` call a fake saw, as the `writes` array it carried.
6240    pub(crate) type ApplyWritesLog = Arc<std::sync::Mutex<Vec<Vec<Value>>>>;
6241
6242    /// A fake that answers `applyWrites` the way a strict PDS does, on BOTH
6243    /// shapes this crate sends it in: the PDS's own
6244    /// `/xrpc/com.atproto.repo.applyWrites` (the direct and OAuth clients) and
6245    /// the sidecar's `/internal/repo` (`action: "applyWrites"`).
6246    ///
6247    /// It refuses what the reference PDS refuses — more than 200 writes
6248    /// (`InvalidRequest: Too many writes. Max: 200`, from
6249    /// `packages/pds/src/api/com/atproto/repo/applyWrites.ts`) and a body over
6250    /// the 150 KiB `jsonLimit` every reference PDS before atproto#4989 applied
6251    /// to it — so a test that sends an unchunked batch fails the way production
6252    /// would, rather than passing against a fake that accepts anything.
6253    ///
6254    /// Every call is logged, refused or not, so a test can assert that a later
6255    /// chunk was never SENT. `fail_call` (1-based) answers that call with a 500.
6256    pub(crate) async fn serve_apply_writes(fail_call: Option<usize>) -> (String, ApplyWritesLog) {
6257        use axum::body::Bytes;
6258        use axum::http::{StatusCode as Status, Uri};
6259        use axum::response::IntoResponse;
6260
6261        const PRE_4989_JSON_LIMIT: usize = 150 * 1024;
6262        let log: ApplyWritesLog = Arc::default();
6263        let sink = Arc::clone(&log);
6264        let app = axum::Router::new()
6265            .fallback(move |uri: Uri, body: Bytes| {
6266                let sink = Arc::clone(&sink);
6267                async move {
6268                    let sidecar = uri.path() == "/internal/repo";
6269                    let reply = |status: Status, error: &str, message: &str| {
6270                        let body = if sidecar {
6271                            json!({ "ok": false, "error": error, "message": message, "status": status.as_u16() })
6272                        } else {
6273                            json!({ "error": error, "message": message })
6274                        };
6275                        (status, axum::Json(body)).into_response()
6276                    };
6277                    let parsed: Value = serde_json::from_slice(&body).unwrap_or(Value::Null);
6278                    // Anything else a handler sends on the way (a folder
6279                    // listing, say) is answered empty and not logged, so the
6280                    // log and `fail_call` count `applyWrites` calls only.
6281                    let Some(writes) = parsed["writes"].as_array().cloned() else {
6282                        let empty = json!({ "records": [] });
6283                        return if sidecar {
6284                            axum::Json(json!({ "ok": true, "data": empty })).into_response()
6285                        } else {
6286                            axum::Json(empty).into_response()
6287                        };
6288                    };
6289                    let call = {
6290                        let mut calls = sink.lock().unwrap();
6291                        calls.push(writes.clone());
6292                        calls.len()
6293                    };
6294                    if body.len() > PRE_4989_JSON_LIMIT {
6295                        return reply(
6296                            Status::PAYLOAD_TOO_LARGE,
6297                            "PayloadTooLarge",
6298                            "request entity too large",
6299                        );
6300                    }
6301                    if writes.len() > 200 {
6302                        return reply(
6303                            Status::BAD_REQUEST,
6304                            "InvalidRequest",
6305                            "Too many writes. Max: 200",
6306                        );
6307                    }
6308                    if fail_call == Some(call) {
6309                        return reply(
6310                            Status::INTERNAL_SERVER_ERROR,
6311                            "InternalServerError",
6312                            "boom",
6313                        );
6314                    }
6315                    let data =
6316                        json!({ "commit": { "cid": "bafycommit", "rev": "3l" }, "results": [] });
6317                    if sidecar {
6318                        axum::Json(json!({ "ok": true, "data": data })).into_response()
6319                    } else {
6320                        axum::Json(data).into_response()
6321                    }
6322                }
6323            })
6324            // The fake must see an oversized body to refuse it, not have axum
6325            // refuse it first at its own 2 MB default.
6326            .layer(axum::extract::DefaultBodyLimit::disable());
6327        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
6328        let addr = listener.local_addr().unwrap();
6329        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
6330        (format!("http://{addr}"), log)
6331    }
6332
6333    /// `n` distinct vetted subscriptions, in a known order.
6334    fn vetted_subs(n: usize) -> Vec<crate::vetted::VettedSubscription> {
6335        (0..n)
6336            .map(|i| {
6337                crate::vetted::VettedSubscription::new(&lexicon::Subscription::new(
6338                    format!("https://f{i}.example/feed.xml"),
6339                    "2026-07-12T00:00:00.000Z",
6340                ))
6341            })
6342            .collect()
6343    }
6344
6345    /// The `rkey` of every write a run of calls carried, flattened in send order.
6346    pub(crate) fn sent_rkeys(log: &ApplyWritesLog) -> Vec<String> {
6347        log.lock()
6348            .unwrap()
6349            .iter()
6350            .flatten()
6351            .map(|w| w["rkey"].as_str().unwrap_or_default().to_string())
6352            .collect()
6353    }
6354
6355    /// How many writes each call carried, in send order.
6356    pub(crate) fn call_sizes(log: &ApplyWritesLog) -> Vec<usize> {
6357        log.lock().unwrap().iter().map(Vec::len).collect()
6358    }
6359
6360    const CHUNK_DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
6361
6362    fn chunk_sidecar(base: &str) -> SidecarClient {
6363        SidecarClient::new(Client::new(), base, base, "secret")
6364    }
6365
6366    /// **201 writes are two calls, 200 then 1, in order.** The OPML import used
6367    /// to send every feed in one `applyWrites`, which the reference PDS refuses
6368    /// past 200 — so any import over 200 feeds failed outright.
6369    #[tokio::test]
6370    async fn sidecar_bulk_add_of_201_is_two_calls_in_order() {
6371        let (base, log) = serve_apply_writes(None).await;
6372        let rkeys = chunk_sidecar(&base)
6373            .add_subscriptions_bulk(CHUNK_DID, &vetted_subs(201))
6374            .await
6375            .expect("a 201-feed import must succeed against a PDS that caps at 200");
6376
6377        assert_eq!(call_sizes(&log), vec![200, 1]);
6378        let urls: Vec<String> = log
6379            .lock()
6380            .unwrap()
6381            .iter()
6382            .flatten()
6383            .map(|w| w["value"]["url"].as_str().unwrap().to_string())
6384            .collect();
6385        let expected: Vec<String> = (0..201)
6386            .map(|i| format!("https://f{i}.example/feed.xml"))
6387            .collect();
6388        assert_eq!(urls, expected, "ops must keep input order across chunks");
6389        assert_eq!(
6390            sent_rkeys(&log),
6391            rkeys,
6392            "the returned rkeys are the ones written, in order"
6393        );
6394    }
6395
6396    /// 500 — the default per-DID cap, so the largest import a stock instance
6397    /// sends — is three calls.
6398    #[tokio::test]
6399    async fn sidecar_bulk_add_of_500_is_three_calls() {
6400        let (base, log) = serve_apply_writes(None).await;
6401        chunk_sidecar(&base)
6402            .add_subscriptions_bulk(CHUNK_DID, &vetted_subs(500))
6403            .await
6404            .expect("bulk write");
6405        assert_eq!(call_sizes(&log), vec![200, 200, 100]);
6406    }
6407
6408    /// Exactly the limit is ONE call: chunking must not split a batch that fits.
6409    #[tokio::test]
6410    async fn sidecar_bulk_add_of_exactly_200_is_one_call() {
6411        let (base, log) = serve_apply_writes(None).await;
6412        chunk_sidecar(&base)
6413            .add_subscriptions_bulk(CHUNK_DID, &vetted_subs(200))
6414            .await
6415            .expect("bulk write");
6416        assert_eq!(call_sizes(&log), vec![200]);
6417    }
6418
6419    /// Nothing to write is no call at all — the behaviour both live clients
6420    /// already had (the sidecar refuses an empty `writes[]` with a 400).
6421    #[tokio::test]
6422    async fn sidecar_bulk_add_of_nothing_sends_nothing() {
6423        let (base, log) = serve_apply_writes(None).await;
6424        let rkeys = chunk_sidecar(&base)
6425            .add_subscriptions_bulk(CHUNK_DID, &[])
6426            .await
6427            .expect("an empty import is not an error");
6428        assert!(rkeys.is_empty());
6429        assert!(
6430            call_sizes(&log).is_empty(),
6431            "an empty batch must not be sent"
6432        );
6433    }
6434
6435    /// **A failed chunk stops the run.** Chunk 2 of 3 fails: the call returns
6436    /// an error, and chunk 3 is never sent — sending it would commit writes
6437    /// after a gap, which no caller could describe as a prefix.
6438    #[tokio::test]
6439    async fn sidecar_bulk_add_stops_at_the_first_failed_chunk() {
6440        let (base, log) = serve_apply_writes(Some(2)).await;
6441        let err = chunk_sidecar(&base)
6442            .add_subscriptions_bulk(CHUNK_DID, &vetted_subs(500))
6443            .await
6444            .expect_err("a failed chunk must fail the call");
6445        assert_eq!(call_sizes(&log), vec![200, 200], "chunk 3 must NOT be sent");
6446        assert!(
6447            format!("{err:#}").contains("boom"),
6448            "the PDS's reason was lost: {err:#}"
6449        );
6450        let progress = ApplyWritesIncomplete::of(&err).expect("the error says how far it got");
6451        assert_eq!(
6452            (progress.landed, progress.in_doubt, progress.total),
6453            (200, 200, 500),
6454            "chunk 1 landed, chunk 2 is in doubt, chunk 3 was never sent"
6455        );
6456        // A plain `%err` log line still names the PDS's reason.
6457        assert!(err.to_string().contains("boom"), "{err}");
6458    }
6459
6460    /// A batch that fit in one call fails exactly as it did before chunking:
6461    /// the same message, and nothing landed.
6462    #[tokio::test]
6463    async fn a_single_chunk_failure_reads_as_it_always_did() {
6464        let (base, _log) = serve_apply_writes(Some(1)).await;
6465        let err = chunk_sidecar(&base)
6466            .add_subscriptions_bulk(CHUNK_DID, &vetted_subs(3))
6467            .await
6468            .expect_err("refused");
6469        let progress = ApplyWritesIncomplete::of(&err).expect("progress");
6470        assert_eq!((progress.landed, progress.in_doubt), (0, 3));
6471        assert!(
6472            !err.to_string().contains("applyWrites call"),
6473            "a one-call batch has no progress to report: {err}"
6474        );
6475        assert_eq!(
6476            format!("{err:#}").matches("boom").count(),
6477            1,
6478            "the cause must not print twice in the chain: {err:#}"
6479        );
6480    }
6481
6482    fn create_ops(n: usize, value_bytes: usize) -> Vec<WriteOp> {
6483        (0..n)
6484            .map(|i| WriteOp::Create {
6485                collection: lexicon::nsid::SUBSCRIPTION.to_string(),
6486                rkey: Some(format!("rk{i:05}")),
6487                value: json!({ "pad": "x".repeat(value_bytes) }),
6488            })
6489            .collect()
6490    }
6491
6492    /// The op-count boundary, on small ops the byte bound never touches.
6493    #[test]
6494    // A one-range Vec IS the expected value here: one chunk spanning the batch.
6495    #[allow(clippy::single_range_in_vec_init)]
6496    fn chunks_split_at_200_ops_and_not_before() {
6497        for (n, want) in [
6498            (0, vec![]),
6499            (1, vec![0..1]),
6500            (200, vec![0..200]),
6501            (201, vec![0..200, 200..201]),
6502            (500, vec![0..200, 200..400, 400..500]),
6503        ] {
6504            assert_eq!(chunk_writes(&create_ops(n, 8)), want, "{n} ops");
6505        }
6506    }
6507
6508    /// **The byte bound splits under 200 ops, and every chunk fits it.**
6509    #[test]
6510    fn chunks_split_on_bytes_and_each_fits() {
6511        let ops = create_ops(40, 10_000);
6512        let chunks = chunk_writes(&ops);
6513        assert!(chunks.len() > 1, "400 KB went out as {chunks:?}");
6514        let mut next = 0;
6515        for range in &chunks {
6516            assert_eq!(range.start, next, "chunks must be consecutive: {chunks:?}");
6517            next = range.end;
6518            let body = json!({
6519                "repo": CHUNK_DID,
6520                "writes": ops[range.clone()].iter().map(WriteOp::to_json).collect::<Vec<_>>(),
6521            })
6522            .to_string();
6523            assert!(
6524                body.len() <= APPLY_WRITES_MAX_BYTES + 200,
6525                "a {}-byte body for {range:?}",
6526                body.len()
6527            );
6528        }
6529        assert_eq!(next, ops.len(), "every op, once");
6530    }
6531
6532    /// **The byte bound does not split an ordinary import.** 200 subscriptions
6533    /// with a title and a site URL each fit one call, so the common OPML
6534    /// import pays one round trip per 200 feeds and no more — the figure the
6535    /// bound's doc comment rests on.
6536    #[test]
6537    fn a_realistic_200_feed_import_is_one_call() {
6538        let ops: Vec<WriteOp> = (0..200)
6539            .map(|i| {
6540                let mut sub = lexicon::Subscription::new(
6541                    format!("https://www.example-blog-{i:03}.com/feeds/posts/default.xml"),
6542                    "2026-07-12T00:00:00.000Z",
6543                );
6544                sub.title = Some(format!("An Example Blog With A Fairly Long Title {i}"));
6545                sub.site_url = Some(format!("https://www.example-blog-{i:03}.com/"));
6546                WriteOp::Create {
6547                    collection: lexicon::nsid::SUBSCRIPTION.to_string(),
6548                    rkey: Some(format!("3lab2c4d5e{i:03}")),
6549                    value: serde_json::to_value(&sub).unwrap(),
6550                }
6551            })
6552            .collect();
6553        // ~76 KB measured.
6554        let bytes: usize = ops.iter().map(|op| op.to_json().to_string().len()).sum();
6555        assert_eq!(chunk_writes(&ops), vec![0..200], "{bytes} bytes");
6556    }
6557
6558    /// An op bigger than the bound cannot be split: it goes alone, and the
6559    /// ops around it are not dragged into its call.
6560    #[test]
6561    fn an_oversized_op_goes_alone() {
6562        let mut ops = create_ops(3, 8);
6563        ops.insert(1, create_ops(1, APPLY_WRITES_MAX_BYTES + 1).remove(0));
6564        assert_eq!(chunk_writes(&ops), vec![0..1, 1..2, 2..4]);
6565    }
6566
6567    /// Read-state cursors as large as the lexicon allows: 1,000 ids each.
6568    fn big_cursors(n: usize) -> Vec<(String, ReadState, bool)> {
6569        (0..n)
6570            .map(|i| {
6571                let mut state = ReadState::new(
6572                    format!("https://f{i}.example/feed.xml"),
6573                    None,
6574                    "2026-07-12T00:00:00.000Z",
6575                );
6576                state.read_ids = (0..ReadState::MAX_IDS)
6577                    .map(|j| format!("https://f{i}.example/posts/{j:04}/an-entry-permalink"))
6578                    .collect();
6579                (format!("rk{i:04}"), state, i % 2 == 0)
6580            })
6581            .collect()
6582    }
6583
6584    /// **The byte bound splits a batch well under 200 ops.** Ten full cursors
6585    /// are ~500 KB: under the op cap, over every older reference PDS's 150 KiB
6586    /// body limit. Each call must fit, and together they must carry every
6587    /// cursor, once, in order.
6588    #[tokio::test]
6589    async fn sidecar_read_state_flush_splits_on_bytes_under_200_ops() {
6590        let (base, log) = serve_apply_writes(None).await;
6591        let cursors = big_cursors(10);
6592        chunk_sidecar(&base)
6593            .flush_read_states(CHUNK_DID, &cursors)
6594            .await
6595            .expect("a byte-heavy flush must succeed in chunks");
6596
6597        let sizes = call_sizes(&log);
6598        assert!(
6599            sizes.len() > 1,
6600            "a ~500 KB flush went out as one call: {sizes:?}"
6601        );
6602        let want: Vec<String> = cursors.iter().map(|(rkey, _, _)| rkey.clone()).collect();
6603        assert_eq!(sent_rkeys(&log), want, "every cursor, once, in order");
6604    }
6605
6606    /// The direct (app-password) client gets the same chunking: its
6607    /// `flush_read_states` is the same op list on a different wire.
6608    #[tokio::test]
6609    async fn direct_client_read_state_flush_of_201_is_two_calls() {
6610        let (base, log) = serve_apply_writes(None).await;
6611        let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
6612        let host = format!("chunk-direct-{port}.test");
6613        crate::net::test_host_override(&host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
6614        let client = PdsClient::new(
6615            ssrf_test_client(),
6616            format!("http://{host}:{port}"),
6617            CHUNK_DID,
6618            Auth::Session(SessionAuth {
6619                did: CHUNK_DID.to_string(),
6620                handle: None,
6621                access_jwt: "jwt".to_string(),
6622                refresh_jwt: None,
6623            }),
6624        );
6625        let cursors: Vec<(String, ReadState, bool)> = (0..201)
6626            .map(|i| {
6627                let feed = format!("https://f{i}.example/feed.xml");
6628                let state = ReadState::new(feed, None, "2026-07-12T00:00:00.000Z");
6629                (format!("rk{i:04}"), state, true)
6630            })
6631            .collect();
6632        client.flush_read_states(&cursors).await.expect("flush");
6633        assert_eq!(call_sizes(&log), vec![200, 1]);
6634        let want: Vec<String> = cursors.iter().map(|(rkey, _, _)| rkey.clone()).collect();
6635        assert_eq!(sent_rkeys(&log), want);
6636    }
6637}