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