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