Skip to main content

feather_reader/oauth/
xrpc.rs

1//! DPoP-bound `com.atproto.repo.*` calls against the user's PDS.
2//!
3//! The layer the 31-method reader surface is built on. It returns the same
4//! [`RecordEntry`] / [`WriteResult`] / [`WriteOp`] types the existing
5//! [`crate::atproto`] client does, so the typed wrappers above it (subscriptions,
6//! folders, saved, read-state) transfer at cutover rather than being rewritten.
7//!
8//! Two properties that are not obvious from the call shapes:
9//!
10//! * **The PDS comes from the session's `aud`**, never re-derived per call. It
11//!   is a property of the token set — a token is valid at one PDS — and it is
12//!   AAD-bound in the session row precisely so it cannot be repointed.
13//! * **Every request carries both the proof and the token.** A resource request
14//!   is `Authorization: DPoP <token>` plus a `DPoP` proof whose `ath` binds to
15//!   that token; either alone is useless.
16
17use anyhow::{bail, Context as _, Result};
18use reqwest::Client;
19use serde_json::{json, Value};
20use sqlx::SqlitePool;
21
22use crate::atproto::{RecordEntry, WriteOp, WriteResult};
23
24use super::dpop::Endpoint;
25use super::keys::SigningKey;
26use super::request::{self, DpopBody, DpopRequest, Retry};
27use super::store::OAuthSession;
28
29/// Cap on how many pages `list_all_records` will walk.
30///
31/// The same bound the existing client uses: a repo is user-controlled, and an
32/// unbounded walk is a denial-of-service against ourselves.
33const MAX_LIST_PAGES: usize = 50;
34
35/// Hard cap on records accumulated by one walk — 50 pages x the 100 we request.
36/// See [`crate::atproto::extend_bounded`] for why exceeding it is an error
37/// rather than a truncation.
38const MAX_LIST_RECORDS: usize = 5_000;
39
40/// An authenticated handle on one account's repo.
41pub struct Repo<'a> {
42    pub http: &'a Client,
43    pub pool: &'a SqlitePool,
44    pub session: &'a OAuthSession,
45    /// The session's DPoP key, already unsealed.
46    pub key: &'a SigningKey,
47}
48
49impl Repo<'_> {
50    /// The XRPC endpoint for a method, on THIS session's PDS.
51    fn url(&self, nsid: &str) -> String {
52        format!("{}/xrpc/{nsid}", self.session.aud.trim_end_matches('/'))
53    }
54
55    /// Send, and fail loudly on a non-2xx with the XRPC error if there is one.
56    async fn send_raw(
57        &self,
58        url: &str,
59        body: DpopBody<'_>,
60        nsid: &str,
61    ) -> Result<request::PostOutcome> {
62        let outcome = request::send_with_dpop(
63            self.http,
64            self.pool,
65            &DpopRequest {
66                endpoint: Endpoint::ResourceServer,
67                url,
68                key: self.key,
69                access_token: Some(&self.session.access_token),
70                body,
71                // A repo call is repeated on a nonce challenge, on the
72                // ASSUMPTION that a PDS answering with a challenge has not acted
73                // on the request. That is how conformant servers behave and what
74                // the reference relies on, but it is not a guarantee we can
75                // verify: a PDS or intermediary that emitted `use_dpop_nonce`
76                // after a write committed would yield a duplicate record from
77                // `createRecord` or `applyWrites`, both of which are
78                // non-idempotent here.
79                retry: Retry::Allowed,
80            },
81        )
82        .await?;
83
84        if !outcome.is_success() {
85            bail!(
86                "{nsid} failed: {}",
87                xrpc_error(&outcome.body, outcome.status)
88            );
89        }
90        Ok(outcome)
91    }
92
93    /// [`send_raw`](Self::send_raw), then the body as JSON.
94    ///
95    /// Every call that wants a `Value` goes through here; `listRecords` does
96    /// not, because turning its body into a `Value` and then into a page
97    /// materialises the records twice.
98    async fn send(&self, url: &str, body: DpopBody<'_>, nsid: &str) -> Result<Value> {
99        let outcome = self.send_raw(url, body, nsid).await?;
100        // A 200 with an empty body is legitimate for deleteRecord/applyWrites.
101        if outcome.body.is_empty() {
102            return Ok(Value::Null);
103        }
104        outcome.json()
105    }
106
107    /// Render an XRPC error body for a message.
108    ///
109    /// Only ever called on a NON-success, where the body is an error document
110    /// rather than a token or a record — and even then only the `error` and
111    /// `message` fields, never the raw bytes.
112    ///
113    /// **Bounded by length before it is parsed.** This is the error twin of
114    /// `send`, whose success branch goes through `PostOutcome::json` and its node
115    /// guard — so without this a hostile PDS answering **500** instead of 200 got
116    /// the whole amplification the guard exists to stop, on every write and on
117    /// every listing failure. An error document is a few dozen bytes; see
118    /// [`super::MAX_ERROR_BODY`].
119    fn error_fields(body: &[u8]) -> Option<String> {
120        if !super::error_body_worth_parsing(body) {
121            return None;
122        }
123        let value: Value = serde_json::from_slice(body).ok()?;
124        let kind = value.get("error").and_then(Value::as_str)?;
125        match value.get("message").and_then(Value::as_str) {
126            Some(message) => Some(format!("{kind}: {message}")),
127            None => Some(kind.to_string()),
128        }
129    }
130
131    /// One page of a collection.
132    pub async fn list_records(
133        &self,
134        collection: &str,
135        limit: Option<u32>,
136        cursor: Option<&str>,
137    ) -> Result<(Vec<RecordEntry>, Option<String>)> {
138        let mut url = url::Url::parse(&self.url("com.atproto.repo.listRecords"))
139            .context("building the listRecords URL")?;
140        {
141            let mut query = url.query_pairs_mut();
142            query.append_pair("repo", &self.session.sub);
143            query.append_pair("collection", collection);
144            if let Some(limit) = limit {
145                query.append_pair("limit", &limit.to_string());
146            }
147            if let Some(cursor) = cursor {
148                query.append_pair("cursor", cursor);
149            }
150        }
151
152        // **Raw bytes, parsed once.** This is the live backend's list walk, the
153        // largest body this client reads, and routing it through `Value` first
154        // held two copies of every page at the same time.
155        let outcome = self
156            .send_raw(
157                url.as_str(),
158                DpopBody::Query,
159                "com.atproto.repo.listRecords",
160            )
161            .await?;
162        // A 2xx carrying an error envelope is NOT an empty page: `records`
163        // defaulting to `[]` turned a PDS failure into `Ok(vec![])`, which
164        // `resolve_subscriptions` reads as "this DID follows nothing" and
165        // `sync_sub_refs` then writes through, revoking every `sub_ref`.
166        // Both invariants, through the one function every listRecords caller
167        // shares — this check was added here first and had to be fitted to the
168        // other two clients a round later. It reads bytes now rather than a
169        // `Value`; the invariants are the same ones, expressed as fields.
170        let page = crate::atproto::parse_list_records(&outcome.body)?;
171        let (records, cursor) = (page.records, page.cursor);
172        Ok((records, cursor))
173    }
174
175    /// Every record in a collection, following the cursor.
176    ///
177    /// Bounded at `MAX_LIST_PAGES`: the repo is user-controlled, so an
178    /// unbounded walk is a denial-of-service against ourselves. A cursor that
179    /// does not advance also terminates the walk rather than spinning.
180    pub async fn list_all_records(&self, collection: &str) -> Result<Vec<RecordEntry>> {
181        self.list_all_records_within(
182            collection,
183            &mut crate::atproto::ByteBudget::new(crate::atproto::MAX_LIST_BYTES),
184        )
185        .await
186    }
187
188    /// [`list_all_records`](Self::list_all_records) against a caller's budget.
189    pub(crate) async fn list_all_records_within(
190        &self,
191        collection: &str,
192        budget: &mut crate::atproto::ByteBudget,
193    ) -> Result<Vec<RecordEntry>> {
194        let mut out = Vec::new();
195        let max_bytes = budget.max();
196        let mut cursor: Option<String> = None;
197        let mut more_offered = false;
198
199        for _ in 0..MAX_LIST_PAGES {
200            let (page, next) = self
201                .list_records(collection, Some(100), cursor.as_deref())
202                .await?;
203            let got = page.len();
204            // Same cap, same reason as `atproto::extend_bounded`: MAX_LIST_PAGES
205            // bounds requests, not memory, unless the server honours our limit.
206            // This is the LIVE walk on `backend=rust` — its result reaches
207            // `replace_sub_refs`, so a truncation here is revoked access.
208            if !budget.admit(&page) {
209                anyhow::bail!(
210                    "listRecords for {collection} exceeded the {max_bytes}-byte cap \
211                     ({} held, {} bytes charged) — refusing to accumulate further",
212                    out.len(),
213                    budget.used(),
214                );
215            }
216            crate::atproto::extend_bounded(&mut out, page, MAX_LIST_RECORDS, collection)?;
217            match next {
218                // `got > 0` is not defensive tidiness -- it is a whole round
219                // trip. This project's own PDS returns a cursor ALONGSIDE a
220                // short page, so without it every list fetches a second, empty
221                // page before stopping. Both of the other clients have this
222                // guard; measured, its absence was most of the remaining gap
223                // against the sidecar.
224                //
225                // The cursor-repeat check is the separate concern: a server that
226                // hands back the same cursor forever would otherwise loop.
227                Some(next) if got > 0 && Some(&next) != cursor.as_ref() => {
228                    cursor = Some(next);
229                    more_offered = true;
230                }
231                _ => {
232                    // `break` with the flag cleared, rather than an early
233                    // `return`: an early return makes the post-loop check
234                    // unreachable, so the flag reads as dead state and a
235                    // reviewer hunts for the case that clears it. It also left
236                    // the "every walk refuses" mutation alive here.
237                    more_offered = false;
238                    break;
239                }
240            }
241        }
242        // **Running out of pages is a refusal, not a short answer.** Falling out
243        // of the loop used to return `Ok(out)`, so a repo bigger than the page
244        // budget produced a truncated list indistinguishable from a complete
245        // one — and `resolve_subscriptions` needs an `Err` for its fail-closed
246        // branch. Given `Ok`, it hands the short list to `replace_sub_refs`,
247        // which DELETEs the reader's whole `sub_ref` projection and reinserts
248        // only what it was given. `extend_bounded` cannot catch this either:
249        // `MAX_LIST_PAGES` x the 100 we request is `MAX_LIST_RECORDS`, so against
250        // a server that honours our limit the page budget runs out first.
251        //
252        // **The cap is on REQUESTS, though, so the record count it bites at is
253        // the server's page size x the budget — not a number our constants fix.**
254        // A PDS answering 50 a page reaches half as far; one answering more than
255        // asked trips `extend_bounded` instead. And the last allowed page is a
256        // FALSE refusal: terminating costs one extra request when a short page
257        // still carries a cursor, so a walk holding every record it will ever
258        // hold still refuses on a cursor it never followed.
259        //
260        // This client's caps are a QUARTER of the direct client's (50 pages,
261        // 5 000 records), so with `limit=100` honoured the same reader refuses
262        // here at ~4 900 records and works to ~19 900 on the sidecar. `Saved`
263        // walks this too, one record per starred article, where 4 900 is a
264        // plausible number for a real reader.
265        if more_offered {
266            anyhow::bail!(
267                "listRecords for {collection} did not finish within {MAX_LIST_PAGES} pages \
268                 ({} held, and the PDS still offered more) — refusing a short list",
269                out.len(),
270            );
271        }
272        Ok(out)
273    }
274
275    /// Create a record, letting the PDS assign the key.
276    pub async fn create_record<T: crate::vetted::WritableRecord>(
277        &self,
278        collection: &str,
279        record: &T,
280    ) -> Result<WriteResult> {
281        self.write(
282            "com.atproto.repo.createRecord",
283            json!({ "repo": self.session.sub, "collection": collection, "record": record }),
284        )
285        .await
286    }
287
288    /// Create or replace a record at a known key.
289    pub async fn put_record<T: crate::vetted::WritableRecord>(
290        &self,
291        collection: &str,
292        rkey: &str,
293        record: &T,
294    ) -> Result<WriteResult> {
295        self.write(
296            "com.atproto.repo.putRecord",
297            json!({
298                "repo": self.session.sub,
299                "collection": collection,
300                "rkey": rkey,
301                "record": record,
302            }),
303        )
304        .await
305    }
306
307    pub async fn delete_record(&self, collection: &str, rkey: &str) -> Result<()> {
308        let body = json!({ "repo": self.session.sub, "collection": collection, "rkey": rkey });
309        self.send(
310            &self.url("com.atproto.repo.deleteRecord"),
311            DpopBody::Json(serde_json::to_vec(&body)?),
312            "com.atproto.repo.deleteRecord",
313        )
314        .await
315        .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
316        Ok(())
317    }
318
319    /// A batch of writes in one round trip.
320    pub async fn apply_writes(&self, writes: &[WriteOp]) -> Result<()> {
321        if writes.is_empty() {
322            return Ok(());
323        }
324        let ops: Vec<Value> = writes.iter().map(WriteOp::to_json).collect();
325        let body = json!({ "repo": self.session.sub, "writes": ops });
326        self.send(
327            &self.url("com.atproto.repo.applyWrites"),
328            DpopBody::Json(serde_json::to_vec(&body)?),
329            "com.atproto.repo.applyWrites",
330        )
331        .await
332        .and_then(|v| crate::atproto::reject_error_envelope(&v))?;
333        Ok(())
334    }
335
336    async fn write(&self, nsid: &str, body: Value) -> Result<WriteResult> {
337        let value = self
338            .send(
339                &self.url(nsid),
340                DpopBody::Json(serde_json::to_vec(&body)?),
341                nsid,
342            )
343            .await?;
344        serde_json::from_value(value).with_context(|| format!("{nsid} returned no usable result"))
345    }
346}
347
348/// Format an XRPC failure for an error message.
349fn xrpc_error(body: &[u8], status: u16) -> String {
350    match Repo::error_fields(body) {
351        Some(detail) => format!("status {status} ({detail})"),
352        None => format!("status {status}"),
353    }
354}
355
356/// The reader's typed surface, over [`Repo`].
357///
358/// Deliberately thin: each method is one repo call plus a parse, and the
359/// orderings come from [`crate::lexicon::sort`], SHARED with the sidecar client
360/// so the two cannot disagree across the cutover. A divergence there would not
361/// be subtle — it would reorder the user's feed list the moment the
362/// implementation swapped.
363impl Repo<'_> {
364    /// List a collection and parse each record into `T`, paired with its rkey.
365    ///
366    /// An unparseable record is SKIPPED with a warning rather than failing the
367    /// list. Records are written by other clients and by future versions of this
368    /// one; one record this build cannot read must not black out the whole feed
369    /// list.
370    async fn list_typed<T: serde::de::DeserializeOwned>(
371        &self,
372        collection: &str,
373    ) -> Result<Vec<(String, T)>> {
374        let records = self.list_all_records(collection).await?;
375        let mut out = Vec::with_capacity(records.len());
376        for record in records {
377            let rkey = record.rkey().unwrap_or_default().to_string();
378            match record.parse::<T>() {
379                Ok(value) => out.push((rkey, value)),
380                Err(err) => tracing::warn!(
381                    collection,
382                    uri = %record.uri,
383                    error = %err,
384                    "skipping unparseable record in collection"
385                ),
386            }
387        }
388        Ok(out)
389    }
390
391    // ── subscriptions ────────────────────────────────────────────────────────
392
393    pub async fn list_subscriptions(&self) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
394        self.list_typed(crate::lexicon::nsid::SUBSCRIPTION).await
395    }
396
397    /// Every subscription, in the reader's deterministic order.
398    pub async fn list_subscriptions_sorted(
399        &self,
400    ) -> Result<Vec<(String, crate::lexicon::Subscription)>> {
401        let mut subs = self.list_subscriptions().await?;
402        subs.sort_by(crate::lexicon::sort::subscriptions);
403        Ok(subs)
404    }
405
406    /// Subscribe to a feed. Returns the new record's rkey so the caller can
407    /// address it (rename, delete) without re-listing.
408    pub async fn add_subscription(
409        &self,
410        sub: &crate::vetted::VettedSubscription,
411    ) -> Result<String> {
412        Ok(self
413            .create_record(crate::lexicon::nsid::SUBSCRIPTION, sub)
414            .await?
415            .into_rkey())
416    }
417
418    pub async fn remove_subscription(&self, rkey: &str) -> Result<()> {
419        self.delete_record(crate::lexicon::nsid::SUBSCRIPTION, rkey)
420            .await
421    }
422
423    /// Replace a subscription in place — retitle, refile, change cadence.
424    /// Returns the [`WriteResult`] rather than discarding it: the caller logs
425    /// the resulting record URI, and a `()` here would have silently dropped
426    /// that field from the log line after the cutover.
427    pub async fn update_subscription(
428        &self,
429        rkey: &str,
430        sub: &crate::vetted::VettedSubscription,
431    ) -> Result<WriteResult> {
432        self.put_record(crate::lexicon::nsid::SUBSCRIPTION, rkey, sub)
433            .await
434    }
435
436    /// Batch-add many subscriptions in one `applyWrites` — the OPML-import path.
437    ///
438    /// One round trip rather than N: an import of several hundred feeds is the
439    /// case this exists for.
440    /// Returns the new rkeys, which are assigned HERE rather than by the server:
441    /// client-side TIDs keep the imported feeds in input order and make the
442    /// batch reproducible. The sidecar client does the same, with the same
443    /// generator.
444    pub async fn add_subscriptions_bulk(
445        &self,
446        subs: &[crate::vetted::VettedSubscription],
447    ) -> Result<Vec<String>> {
448        let mut gen = crate::atproto::TidGenerator::new();
449        let mut rkeys = Vec::with_capacity(subs.len());
450        let mut writes = Vec::with_capacity(subs.len());
451        for sub in subs {
452            let rkey = gen.next();
453            writes.push(WriteOp::Create {
454                collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
455                rkey: Some(rkey.clone()),
456                // Propagated, NOT defaulted: `unwrap_or(Value::Null)` here would
457                // write a null record into the user's repo on a serialization
458                // failure rather than failing the import.
459                value: serde_json::to_value(sub)?,
460            });
461            rkeys.push(rkey);
462        }
463        self.apply_writes(&writes).await?;
464        Ok(rkeys)
465    }
466
467    // ── folders ──────────────────────────────────────────────────────────────
468
469    pub async fn list_folders(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
470        self.list_typed(crate::lexicon::nsid::FOLDER).await
471    }
472
473    pub async fn list_folders_sorted(&self) -> Result<Vec<(String, crate::lexicon::Folder)>> {
474        let mut folders = self.list_folders().await?;
475        folders.sort_by(crate::lexicon::sort::folders);
476        Ok(folders)
477    }
478
479    pub async fn add_folder(&self, folder: &crate::lexicon::Folder) -> Result<String> {
480        Ok(self
481            .create_record(crate::lexicon::nsid::FOLDER, folder)
482            .await?
483            .into_rkey())
484    }
485
486    /// Delete a folder. Subscriptions referencing it are left alone; a dangling
487    /// reference reads as "unfiled", which is the same behaviour the sidecar
488    /// client has.
489    pub async fn remove_folder(&self, rkey: &str) -> Result<()> {
490        self.delete_record(crate::lexicon::nsid::FOLDER, rkey).await
491    }
492
493    /// Returns the [`WriteResult`], matching the sidecar client — the caller
494    /// logs the record URI from it.
495    pub async fn rename_folder(
496        &self,
497        rkey: &str,
498        folder: &crate::lexicon::Folder,
499    ) -> Result<WriteResult> {
500        self.put_record(crate::lexicon::nsid::FOLDER, rkey, folder)
501            .await
502    }
503
504    // ── saved ────────────────────────────────────────────────────────────────
505
506    pub async fn list_saved(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
507        self.list_typed(crate::lexicon::nsid::SAVED).await
508    }
509
510    /// Saved entries, newest first.
511    pub async fn list_saved_sorted(&self) -> Result<Vec<(String, crate::lexicon::Saved)>> {
512        let mut saved = self.list_saved().await?;
513        saved.sort_by(crate::lexicon::sort::saved);
514        Ok(saved)
515    }
516
517    pub async fn add_saved(&self, saved: &crate::vetted::VettedSaved) -> Result<String> {
518        Ok(self
519            .create_record(crate::lexicon::nsid::SAVED, saved)
520            .await?
521            .into_rkey())
522    }
523
524    pub async fn remove_saved(&self, rkey: &str) -> Result<()> {
525        self.delete_record(crate::lexicon::nsid::SAVED, rkey).await
526    }
527
528    // ── read state ───────────────────────────────────────────────────────────
529
530    pub async fn list_read_states(&self) -> Result<Vec<(String, crate::lexicon::ReadState)>> {
531        self.list_typed(crate::lexicon::nsid::READ_STATE).await
532    }
533
534    /// Upsert one read cursor at its feed-derived rkey.
535    pub async fn put_read_state(
536        &self,
537        rkey: &str,
538        state: &crate::lexicon::ReadState,
539    ) -> Result<()> {
540        self.put_record(crate::lexicon::nsid::READ_STATE, rkey, state)
541            .await?;
542        Ok(())
543    }
544
545    /// Flush many dirty read cursors in one `applyWrites`.
546    ///
547    /// Read state changes on nearly every page view, so this is the hottest
548    /// write path in the app; one round trip per flush rather than per feed is
549    /// the whole point.
550    /// The `bool` is whether the record already exists in the PDS. It is not
551    /// optional bookkeeping: an `#update` on a missing record ERRORS, and
552    /// `applyWrites` is atomic per repo, so one not-yet-created cursor in the
553    /// batch would drop the whole DID's flush. The op builder is shared with the
554    /// sidecar client so the two cannot decide create-vs-update differently.
555    pub async fn flush_read_states(
556        &self,
557        cursors: &[(String, crate::lexicon::ReadState, bool)],
558    ) -> Result<()> {
559        if cursors.is_empty() {
560            return Ok(());
561        }
562        let writes = crate::atproto::read_state_write_ops(cursors)?;
563        self.apply_writes(&writes).await
564    }
565}
566
567#[cfg(test)]
568mod tests {
569    use super::*;
570
571    /// **Both clients must build the SAME write ops.**
572    ///
573    /// `flush_read_states` is the hottest write path in the app, and the
574    /// create-vs-update choice is the one part of it that cannot be got wrong
575    /// quietly: an `#update` on a record that does not exist errors, and
576    /// `applyWrites` is atomic per repo, so a single first-flush cursor in the
577    /// batch takes the whole DID's flush down with it.
578    ///
579    /// The first version of this method here ignored the flag and emitted
580    /// `#create` unconditionally, which would have broken every feed's first
581    /// flush after the cutover. Sharing the builder is what makes that
582    /// impossible rather than merely fixed.
583    #[test]
584    fn read_state_writes_choose_create_or_update_per_cursor() {
585        let state = crate::lexicon::ReadState::new(
586            "https://example.com/feed",
587            Some("2026-01-01T00:00:00Z".to_string()),
588            "2026-01-01T00:00:00Z",
589        );
590        let cursors = vec![
591            ("existing".to_string(), state.clone(), true),
592            ("brand-new".to_string(), state.clone(), false),
593        ];
594
595        let ops = crate::atproto::read_state_write_ops(&cursors).expect("ops build");
596        assert_eq!(ops.len(), 2);
597
598        let rendered: Vec<Value> = ops.iter().map(|op| op.to_json()).collect();
599        assert_eq!(
600            rendered[0]["$type"], "com.atproto.repo.applyWrites#update",
601            "an existing record must be UPDATED, not re-created"
602        );
603        assert_eq!(
604            rendered[1]["$type"], "com.atproto.repo.applyWrites#create",
605            "a first flush must CREATE, or the whole atomic batch fails"
606        );
607    }
608
609    /// **Bulk import, through the real `Repo`, asserted on the bytes.** The
610    /// test this replaces exercised `TidGenerator` directly and never called
611    /// `add_subscriptions_bulk`; the function could stop assigning rkeys and
612    /// return garbage with the suite green.
613    #[tokio::test]
614    async fn bulk_subscribe_writes_client_assigned_ordered_rkeys_to_the_right_collection() {
615        let (base, log) = crate::net::tests::serve_json_capturing(b"{}".to_vec()).await;
616        let port: u16 = base.rsplit(':').next().unwrap().parse().unwrap();
617        crate::net::test_host_override(
618            "bulk-pds.test",
619            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
620        );
621        let http = Client::new();
622        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
623        crate::store::init_schema(&pool).await.unwrap();
624        let key = SigningKey::generate("k");
625        let mut s = session();
626        s.aud = format!("http://bulk-pds.test:{port}");
627        let repo = repo(&http, &pool, &s, &key);
628        let subs: Vec<crate::vetted::VettedSubscription> = (0..3)
629            .map(|i| {
630                crate::vetted::VettedSubscription::new(&crate::lexicon::Subscription::new(
631                    format!("https://f{i}.example/feed.xml"),
632                    "2026-07-12T00:00:00.000Z",
633                ))
634            })
635            .collect();
636
637        let rkeys = repo
638            .add_subscriptions_bulk(&subs)
639            .await
640            .expect("bulk write failed");
641
642        let sent = log.lock().unwrap().clone();
643        assert_eq!(
644            sent.len(),
645            1,
646            "expected one applyWrites request, got {sent:?}"
647        );
648        let body: Value = serde_json::from_str(sent[0].split("\r\n\r\n").nth(1).unwrap())
649            .expect("request body is JSON");
650        let writes = body["writes"].as_array().expect("writes array");
651        assert_eq!(writes.len(), 3);
652        for (i, w) in writes.iter().enumerate() {
653            assert_eq!(w["collection"], crate::lexicon::nsid::SUBSCRIPTION);
654            assert_eq!(w["rkey"].as_str(), Some(rkeys[i].as_str()));
655        }
656        let mut sorted = rkeys.clone();
657        sorted.sort();
658        assert_eq!(rkeys, sorted, "client-assigned rkeys must ascend");
659    }
660
661    fn session() -> OAuthSession {
662        OAuthSession {
663            sub: "did:plc:ewvi7nxzyoun6zhxrhs64oiz".into(),
664            issuer: "https://pds.example.com".into(),
665            aud: "https://pds.example.com".into(),
666            dpop_key_jwk: "{}".into(),
667            access_token: "tok".into(),
668            refresh_token: "ref".into(),
669            token_type: "DPoP".into(),
670            granted_scope: "atproto".into(),
671            expires_at: None,
672        }
673    }
674
675    fn repo<'a>(
676        http: &'a Client,
677        pool: &'a SqlitePool,
678        session: &'a OAuthSession,
679        key: &'a SigningKey,
680    ) -> Repo<'a> {
681        Repo {
682            http,
683            pool,
684            session,
685            key,
686        }
687    }
688
689    /// **The PDS comes from the session**, so a call cannot be pointed at
690    /// another host by a caller who passes the wrong base.
691    #[tokio::test]
692    async fn endpoints_are_built_from_the_sessions_audience() {
693        let http = Client::new();
694        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
695        let key = SigningKey::generate("k");
696        let mut s = session();
697        s.aud = "https://pds.example.com/".into();
698        let repo = repo(&http, &pool, &s, &key);
699        assert_eq!(
700            repo.url("com.atproto.repo.listRecords"),
701            "https://pds.example.com/xrpc/com.atproto.repo.listRecords",
702            "a trailing slash on the audience must not double the separator"
703        );
704    }
705
706    /// An XRPC error document is summarised by its `error`/`message` fields
707    /// only — never by echoing the raw body, which on other paths holds tokens.
708    #[test]
709    fn an_xrpc_error_is_summarised_not_echoed() {
710        let body = br#"{"error":"InvalidRequest","message":"unknown collection"}"#;
711        let rendered = xrpc_error(body, 400);
712        assert!(rendered.contains("InvalidRequest"));
713        assert!(rendered.contains("unknown collection"));
714
715        // A body that is not an XRPC error contributes nothing but the status.
716        let opaque = xrpc_error(br#"{"access_token":"SECRET"}"#, 500);
717        assert_eq!(opaque, "status 500");
718        assert!(!opaque.contains("SECRET"));
719        assert_eq!(xrpc_error(b"<html>oops</html>", 502), "status 502");
720    }
721
722    /// Every repo call must fail closed on an internal target, like every other
723    /// outbound path — asserted on the guard's own error.
724    #[tokio::test]
725    async fn repo_calls_fail_closed_on_an_internal_pds() {
726        let http = Client::new();
727        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
728        super::super::store::init_schema(&pool).await.unwrap();
729        let key = SigningKey::generate("k");
730        let mut s = session();
731        s.aud = "http://127.0.0.1:2583".into();
732        let repo = repo(&http, &pool, &s, &key);
733
734        let err = repo
735            .list_records("app.feather.subscription", None, None)
736            .await
737            .expect_err("must refuse a loopback PDS");
738        assert!(
739            format!("{err:#}").contains("forbidden (internal) address"),
740            "failed for the wrong reason: {err:#}"
741        );
742    }
743
744    /// **A 200 carrying an error envelope must not read as an empty repo.**
745    ///
746    /// `records` was taken off the JSON with `unwrap_or(Array([]))`, so a PDS
747    /// answering `200 {"error": …}` produced `Ok(vec![])`. That is not the
748    /// fail-closed branch in `web::resolve_subscriptions`: `sync_sub_refs`
749    /// writes the empty set through and `replace_sub_refs` DELETEs the DID's
750    /// entire `sub_ref` projection — one bad response revokes the reader's
751    /// access to every feed they have. Driven through the real client against
752    /// a real server, because the bug was the missing CALL, not the check.
753    #[tokio::test]
754    async fn a_200_error_envelope_is_not_an_empty_repo() {
755        let base = crate::net::tests::serve_body(
756            br#"{"error":"InvalidRequest","message":"bad cursor"}"#.to_vec(),
757        )
758        .await;
759        let port: u16 = base
760            .trim_end_matches('/')
761            .rsplit(':')
762            .next()
763            .unwrap()
764            .parse()
765            .unwrap();
766        crate::net::test_host_override(
767            "envelope-pds.test",
768            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
769        );
770
771        let http = Client::new();
772        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
773        crate::store::init_schema(&pool).await.unwrap();
774        let key = SigningKey::generate("k");
775        let mut s = session();
776        s.aud = format!("http://envelope-pds.test:{port}");
777        let repo = repo(&http, &pool, &s, &key);
778
779        let err = repo
780            .list_records("app.feather.subscription", None, None)
781            .await
782            .expect_err("an error envelope was read as an empty page");
783        assert!(
784            format!("{err:#}").contains("InvalidRequest"),
785            "failed for the wrong reason: {err:#}"
786        );
787    }
788
789    /// **The live walk spends its budget across pages too.**
790    ///
791    /// This is the `backend=rust` walk whose result reaches `replace_sub_refs`,
792    /// so a bound that silently failed to accumulate here would revoke a
793    /// reader's access to every feed past the cut. Three pages, a two-page
794    /// budget: the walk must refuse, and must say it kept two.
795    #[tokio::test]
796    async fn the_live_walk_spends_its_budget_across_pages() {
797        let (bodies, per_page) = crate::atproto::tests::paged_bodies(3, 4096, false);
798        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
799        let port: u16 = base
800            .trim_end_matches('/')
801            .rsplit(':')
802            .next()
803            .unwrap()
804            .parse()
805            .unwrap();
806        crate::net::test_host_override(
807            "live-budget-pages.test",
808            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
809        );
810
811        let http = Client::new();
812        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
813        crate::store::init_schema(&pool).await.unwrap();
814        let key = SigningKey::generate("k");
815        let mut s = session();
816        s.aud = format!("http://live-budget-pages.test:{port}");
817        let repo = repo(&http, &pool, &s, &key);
818
819        let err = repo
820            .list_all_records_within(
821                "app.feather.subscription",
822                &mut crate::atproto::ByteBudget::new(per_page * 2),
823            )
824            .await
825            .expect_err("three pages cannot fit in a two-page budget");
826        let msg = format!("{err:#}");
827        assert!(msg.contains("byte cap"), "wrong bound reported: {msg}");
828        assert!(
829            msg.contains("2 held"),
830            "the live walk did not accumulate across pages: {msg}"
831        );
832    }
833
834    /// **The live walk refuses a short list.** Its result reaches
835    /// `replace_sub_refs`, so returning a truncated list as a complete one
836    /// deletes every subscription past the page budget.
837    #[tokio::test]
838    async fn the_live_walk_that_runs_out_of_pages_refuses() {
839        let bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES + 1)
840            .map(|i| {
841                serde_json::json!({
842                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
843                    "cursor": format!("p{}", i + 1),
844                })
845                .to_string()
846                .into_bytes()
847            })
848            .collect();
849        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
850        let port: u16 = base
851            .trim_end_matches('/')
852            .rsplit(':')
853            .next()
854            .unwrap()
855            .parse()
856            .unwrap();
857        crate::net::test_host_override(
858            "pages-exhausted-live.test",
859            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
860        );
861        let http = Client::new();
862        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
863        crate::store::init_schema(&pool).await.unwrap();
864        let key = SigningKey::generate("k");
865        let mut s = session();
866        s.aud = format!("http://pages-exhausted-live.test:{port}");
867        let repo = repo(&http, &pool, &s, &key);
868
869        let err = repo
870            .list_all_records("app.feather.subscription")
871            .await
872            .expect_err("a truncated list was returned as a complete one");
873        assert!(
874            format!("{err:#}").contains("did not finish"),
875            "failed for the wrong reason: {err:#}"
876        );
877    }
878
879    /// The live walk's clean finish must still be a success — see the sidecar's
880    /// twin for why this direction is the dangerous one.
881    #[tokio::test]
882    async fn the_live_walk_that_finishes_cleanly_returns_the_records() {
883        let mut bodies: Vec<Vec<u8>> = (0..3)
884            .map(|i| {
885                serde_json::json!({
886                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
887                    "cursor": format!("p{}", i + 1),
888                })
889                .to_string()
890                .into_bytes()
891            })
892            .collect();
893        // **The terminator CARRIES a record**, because a real PDS ends on a
894        // partial page and those are the records an off-by-one loses. Ending on
895        // an empty page kept "drop the last page's records" alive: the count
896        // below was right either way.
897        bodies.push(
898            serde_json::json!({
899                "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
900            })
901            .to_string()
902            .into_bytes(),
903        );
904        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
905        let port: u16 = base
906            .trim_end_matches('/')
907            .rsplit(':')
908            .next()
909            .unwrap()
910            .parse()
911            .unwrap();
912        crate::net::test_host_override(
913            "clean-finish-live.test",
914            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
915        );
916        let http = Client::new();
917        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
918        crate::store::init_schema(&pool).await.unwrap();
919        let key = SigningKey::generate("k");
920        let mut s = session();
921        s.aud = format!("http://clean-finish-live.test:{port}");
922        let repo = repo(&http, &pool, &s, &key);
923
924        let records = repo
925            .list_all_records("app.feather.subscription")
926            .await
927            .expect("a walk that ran out of records is not a short list");
928        assert_eq!(records.len(), 4);
929        assert!(
930            records.iter().any(|r| r.uri.ends_with("3labLAST")),
931            "the LAST page's records were dropped: {:?}",
932            records.iter().map(|r| r.uri.as_str()).collect::<Vec<_>>(),
933        );
934    }
935
936    /// **The page cap is pinned exactly, not to within one.**
937    ///
938    /// `the_live_walk_that_runs_out_of_pages_refuses` serves `MAX_LIST_PAGES + 1`
939    /// pages, so a budget one page SHORT refuses too and that mutation survives
940    /// it. A walk whose last allowed request is the terminating one must come
941    /// back `Ok` — and on this backend an `Err` is `replace_sub_refs` never
942    /// running, which is the direction that costs a reader their subscriptions.
943    #[tokio::test]
944    async fn a_live_walk_that_terminates_on_its_last_allowed_page_succeeds() {
945        let mut bodies: Vec<Vec<u8>> = (0..MAX_LIST_PAGES - 1)
946            .map(|i| {
947                serde_json::json!({
948                    "records": [{ "uri": format!("at://did:plc:x/c/3lab{i}"), "value": {} }],
949                    "cursor": format!("p{}", i + 1),
950                })
951                .to_string()
952                .into_bytes()
953            })
954            .collect();
955        bodies.push(
956            serde_json::json!({
957                "records": [{ "uri": "at://did:plc:x/c/3labLAST", "value": {} }]
958            })
959            .to_string()
960            .into_bytes(),
961        );
962        assert_eq!(bodies.len(), MAX_LIST_PAGES);
963        let base = crate::net::tests::serve_bodies_in_sequence(bodies).await;
964        let port: u16 = base
965            .trim_end_matches('/')
966            .rsplit(':')
967            .next()
968            .unwrap()
969            .parse()
970            .unwrap();
971        crate::net::test_host_override(
972            "last-allowed-page-live.test",
973            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
974        );
975        let http = Client::new();
976        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
977        crate::store::init_schema(&pool).await.unwrap();
978        let key = SigningKey::generate("k");
979        let mut s = session();
980        s.aud = format!("http://last-allowed-page-live.test:{port}");
981        let repo = repo(&http, &pool, &s, &key);
982
983        let records = repo
984            .list_all_records("app.feather.subscription")
985            .await
986            .expect("a walk that terminated inside its budget is not a short list");
987        assert_eq!(
988            records.len(),
989            MAX_LIST_PAGES,
990            "a walk that used its whole page budget and finished lost records",
991        );
992    }
993
994    /// **The ERROR twin of `send`, which a review found unguarded.**
995    ///
996    /// `send`'s success branch goes through `PostOutcome::json` and its cap; its
997    /// failure branch goes to `xrpc_error` → `error_fields`, which deserialised the
998    /// whole body. So a hostile PDS answering **500** instead of 200 with the same
999    /// explosion got the full amplification on every write and on every listing
1000    /// failure — the guard bypassed by a status code.
1001    ///
1002    /// Both directions: a small error body still yields its reason, because that
1003    /// string is what a reader's log needs to tell "your PDS said no" from "we
1004    /// broke".
1005    #[test]
1006    fn an_oversized_error_body_is_not_parsed_for_its_reason() {
1007        let small = br#"{"error":"InvalidSwap","message":"record changed"}"#;
1008        assert_eq!(
1009            Repo::error_fields(small).as_deref(),
1010            Some("InvalidSwap: record changed"),
1011            "a real error body must still render its reason",
1012        );
1013
1014        let mut huge = String::from(r#"{"error":"InvalidSwap","pad":["#);
1015        while huge.len() < crate::oauth::MAX_ERROR_BODY + 1_024 {
1016            huge.push_str("{},");
1017        }
1018        huge.push_str("{}]}");
1019        assert!(huge.len() > crate::oauth::MAX_ERROR_BODY);
1020        assert_eq!(
1021            Repo::error_fields(huge.as_bytes()),
1022            None,
1023            "an oversized error body was deserialised to fish out one string",
1024        );
1025    }
1026
1027    /// **The write path parses a PDS body too, and it had no bound before the
1028    /// parse.**
1029    ///
1030    /// `listRecords` is guarded by `parse_list_records`; every OTHER body this
1031    /// client turns into a `Value` goes through `PostOutcome::json` — the repo
1032    /// writers' responses, the PAR response, the token response, the session
1033    /// refresh. `read_capped` bounds the WIRE at 8 MB, which is the *input* to the
1034    /// amplification rather than a limit on it: 8 MB of the cheapest node shape
1035    /// measured 824 MB retained on a 512 MB box.
1036    ///
1037    /// Drives `delete_record`, which is the shortest route from a handler to
1038    /// `send` → `json`.
1039    #[tokio::test]
1040    async fn the_live_write_path_refuses_a_node_explosion() {
1041        let mut body = String::from(r#"{"uri":"at://d/c/r","value":["#);
1042        for _ in 0..1_200_000 {
1043            body.push_str("{},");
1044        }
1045        body.push_str("{}]}");
1046        assert!(
1047            crate::atproto::count_structural_chars(body.as_bytes())
1048                > crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
1049            "the probe body is not over the cap, so this test proves nothing",
1050        );
1051        let base = crate::net::tests::serve_body(body.into_bytes()).await;
1052        let port: u16 = base
1053            .trim_end_matches('/')
1054            .rsplit(':')
1055            .next()
1056            .unwrap()
1057            .parse()
1058            .unwrap();
1059        crate::net::test_host_override(
1060            "write-explosion.test",
1061            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1062        );
1063        let http = Client::new();
1064        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1065        crate::store::init_schema(&pool).await.unwrap();
1066        let key = SigningKey::generate("k");
1067        let mut s = session();
1068        s.aud = format!("http://write-explosion.test:{port}");
1069        let repo = repo(&http, &pool, &s, &key);
1070
1071        let err = repo
1072            .delete_record("c", "r")
1073            .await
1074            .expect_err("a node explosion on the write path was parsed rather than refused");
1075        assert!(
1076            format!("{err:#}").contains("structural characters"),
1077            "failed for the wrong reason: {err:#}"
1078        );
1079    }
1080
1081    /// The other direction: an ordinary write response still parses. Without this
1082    /// a guard that refused every body would pass the test above.
1083    #[tokio::test]
1084    async fn the_live_write_path_accepts_an_ordinary_response() {
1085        let base =
1086            crate::net::tests::serve_body(br#"{"commit":{"cid":"bafy","rev":"3lab"}}"#.to_vec())
1087                .await;
1088        let port: u16 = base
1089            .trim_end_matches('/')
1090            .rsplit(':')
1091            .next()
1092            .unwrap()
1093            .parse()
1094            .unwrap();
1095        crate::net::test_host_override(
1096            "write-ordinary.test",
1097            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1098        );
1099        let http = Client::new();
1100        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1101        crate::store::init_schema(&pool).await.unwrap();
1102        let key = SigningKey::generate("k");
1103        let mut s = session();
1104        s.aud = format!("http://write-ordinary.test:{port}");
1105        let repo = repo(&http, &pool, &s, &key);
1106
1107        repo.delete_record("c", "r")
1108            .await
1109            .expect("an ordinary write response was refused");
1110    }
1111
1112    /// **A duplicated `records` key must not be able to empty a page.**
1113    ///
1114    /// This is the live `backend=rust` walk, so an empty page here reaches
1115    /// `replace_sub_refs` and deletes the reader's subscriptions. serde refuses
1116    /// a repeated field outright; a `serde_json::Value` takes the last one
1117    /// silently, so routing this body through a `Value` first turns a smuggled
1118    /// second key into a successful, empty listing. The test exists to pin
1119    /// which of the two this client uses.
1120    #[tokio::test]
1121    async fn the_live_walk_refuses_a_duplicated_records_key() {
1122        let base = crate::net::tests::serve_body(
1123            br#"{"records":[{"uri":"at://d/c/r","value":{}}],"records":[]}"#.to_vec(),
1124        )
1125        .await;
1126        let port: u16 = base
1127            .trim_end_matches('/')
1128            .rsplit(':')
1129            .next()
1130            .unwrap()
1131            .parse()
1132            .unwrap();
1133        crate::net::test_host_override(
1134            "dup-records.test",
1135            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1136        );
1137
1138        let http = Client::new();
1139        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1140        crate::store::init_schema(&pool).await.unwrap();
1141        let key = SigningKey::generate("k");
1142        let mut s = session();
1143        s.aud = format!("http://dup-records.test:{port}");
1144        let repo = repo(&http, &pool, &s, &key);
1145
1146        let err = repo
1147            .list_records("app.feather.subscription", None, None)
1148            .await
1149            .expect_err("a duplicated records key was read as an empty page");
1150        assert!(
1151            format!("{err:#}").contains("duplicate"),
1152            "failed for the wrong reason: {err:#}"
1153        );
1154    }
1155
1156    /// **An empty 2xx body is not an empty repo either.** `send` maps a
1157    /// zero-length 2xx to `Value::Null` — deliberately, for `deleteRecord` and
1158    /// `applyWrites` — and while `listRecords` still went through it, that
1159    /// slipped past the envelope guard
1160    /// and `records.unwrap_or(Array([]))` produced `Ok(vec![])`: the same
1161    /// `sub_ref` wipe the guard was added to prevent, through the sibling door.
1162    #[tokio::test]
1163    async fn an_empty_200_body_is_not_an_empty_repo() {
1164        let base = crate::net::tests::serve_body(Vec::new()).await;
1165        let port: u16 = base
1166            .trim_end_matches('/')
1167            .rsplit(':')
1168            .next()
1169            .unwrap()
1170            .parse()
1171            .unwrap();
1172        crate::net::test_host_override(
1173            "empty-body.test",
1174            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1175        );
1176        let http = Client::new();
1177        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1178        crate::store::init_schema(&pool).await.unwrap();
1179        let key = SigningKey::generate("k");
1180        let mut s = session();
1181        s.aud = format!("http://empty-body.test:{port}");
1182        let repo = repo(&http, &pool, &s, &key);
1183
1184        let err = repo
1185            .list_records("app.feather.subscription", None, None)
1186            .await
1187            .expect_err("an empty body was read as an empty repo");
1188        assert!(
1189            format!("{err:#}").contains("no records"),
1190            "failed for the wrong reason: {err:#}"
1191        );
1192    }
1193
1194    /// The write paths discarded the body too: `delete_record` and
1195    /// `apply_writes` are `send(..).await?; Ok(())`, so a 200 carrying an
1196    /// error envelope reported success for a delete that did not happen.
1197    #[tokio::test]
1198    async fn a_200_error_envelope_is_not_a_successful_write() {
1199        let base = crate::net::tests::serve_body(
1200            br#"{"error":"InvalidRequest","message":"nope"}"#.to_vec(),
1201        )
1202        .await;
1203        let port: u16 = base
1204            .trim_end_matches('/')
1205            .rsplit(':')
1206            .next()
1207            .unwrap()
1208            .parse()
1209            .unwrap();
1210        crate::net::test_host_override(
1211            "envelope-write.test",
1212            std::net::SocketAddr::from(([127, 0, 0, 1], port)),
1213        );
1214        let http = Client::new();
1215        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1216        crate::store::init_schema(&pool).await.unwrap();
1217        let key = SigningKey::generate("k");
1218        let mut s = session();
1219        s.aud = format!("http://envelope-write.test:{port}");
1220        let repo = repo(&http, &pool, &s, &key);
1221
1222        let err = repo
1223            .delete_record("app.feather.subscription", "rk1")
1224            .await
1225            .expect_err("a failed delete was reported as success");
1226        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1227
1228        let err = repo
1229            .apply_writes(&[crate::atproto::WriteOp::Delete {
1230                collection: "app.feather.subscription".to_string(),
1231                rkey: "rk1".to_string(),
1232            }])
1233            .await
1234            .expect_err("a failed batch was reported as success");
1235        assert!(format!("{err:#}").contains("InvalidRequest"), "{err:#}");
1236    }
1237
1238    /// An empty batch must not produce a request at all — an `applyWrites` with
1239    /// no writes is a round trip that can only fail.
1240    #[tokio::test]
1241    async fn an_empty_batch_is_not_sent() {
1242        let http = Client::new();
1243        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1244        let key = SigningKey::generate("k");
1245        let mut s = session();
1246        // A target that would fail loudly if it were ever contacted.
1247        s.aud = "http://127.0.0.1:2583".into();
1248        let repo = repo(&http, &pool, &s, &key);
1249        assert!(repo.apply_writes(&[]).await.is_ok());
1250    }
1251}