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