Skip to main content

feather_reader/
repo.rs

1//! The one place the two repo backends are chosen between, and timed.
2//!
3//! Every `com.atproto.repo.*` call the reader makes goes through here, so the
4//! cutover is a single `match` rather than a swap at every call site — and,
5//! just as importantly, both backends are measured at the same boundary by the
6//! same wrapper. Timers placed separately on each path would be comparing the
7//! timers.
8//!
9//! The method list mirrors [`crate::atproto::SidecarClient`]'s reader surface
10//! exactly, because that is what `web.rs` already calls. Anything that had to be
11//! reshaped to fit would be a divergence between the two clients, and those
12//! belong in `xrpc.rs` where both can share the fix.
13
14use anyhow::{Context as _, Result};
15
16use crate::lexicon::{Folder, ReadState, Saved, Subscription};
17use crate::metrics::{timed, Backend};
18use crate::oauth;
19use crate::vetted::VettedSubscription;
20use crate::AppState;
21
22/// A dispatcher bound to one request's state.
23pub struct Repo<'a> {
24    state: &'a AppState,
25}
26
27impl AppState {
28    /// The repo client for this request, on whichever backend is configured.
29    pub fn repo(&self) -> Repo<'_> {
30        Repo { state: self }
31    }
32}
33
34impl Repo<'_> {
35    fn backend(&self) -> Backend {
36        self.state.config.repo_backend
37    }
38
39    /// The Rust client's runtime, or a clear error naming what is missing.
40    ///
41    /// Only reachable with the Rust backend selected; startup refuses that
42    /// combination when the runtime could not be built, so this is a
43    /// belt-and-braces path rather than the expected failure point.
44    fn rust(&self) -> Result<&oauth::runtime::OauthRuntime> {
45        self.state
46            .oauth
47            .as_deref()
48            .context("the rust repo backend is selected but its OAuth runtime is not configured")
49    }
50
51    /// Whether `did` has a session this backend could actually use.
52    ///
53    /// **A precondition check, not a repo operation — deliberately NOT wrapped
54    /// in `timed()`.** The background flusher calls this every round for every
55    /// DID holding dirty read-state; counting it would re-inflate the very
56    /// `flush_read_states` error count this exists to stop polluting (#117).
57    ///
58    /// Existence only: no decrypt, no staleness check, no refresh. A session
59    /// that is present but expired still counts as usable here, because the
60    /// refresh path is exactly what `session()` will do about it. The question
61    /// being answered is narrower — is there anything at all to work with, or
62    /// is this DID parked until the user signs in again?
63    ///
64    /// **The sidecar arm answers `true` unconditionally**, preserving today's
65    /// behaviour on that backend rather than guessing. The sidecar owns its own
66    /// session store and answering honestly would mean a loopback round trip per
67    /// DID per minute; prod runs `rust`, and that arm is deleted by #18.
68    pub async fn has_session(&self, did: &str) -> Result<bool> {
69        match self.backend() {
70            Backend::Sidecar => Ok(true),
71            Backend::Rust => {
72                let found: Option<(i64,)> =
73                    sqlx::query_as("SELECT 1 FROM oauth_session WHERE sub = ?1")
74                        .bind(did)
75                        .fetch_optional(&self.state.db)
76                        .await
77                        .with_context(|| format!("checking for an OAuth session for {did}"))?;
78                Ok(found.is_some())
79            }
80        }
81    }
82
83    /// Load a usable session for `did`, refreshing only when it is actually
84    /// stale.
85    ///
86    /// **Discovery is deferred to the refresh path.** Building the refresh
87    /// context eagerly would mean an authorization-server metadata fetch on
88    /// every repo call, which is both wasteful and would make the Rust backend
89    /// look slow in the comparison for a reason that is an artefact of the
90    /// wiring rather than the implementation. A live session needs no discovery
91    /// at all.
92    async fn session(&self, did: &str) -> Result<oauth::store::OAuthSession> {
93        let rust = self.rust()?;
94        let now = crate::store::now_unix();
95
96        let session = oauth::store::get_session(&self.state.db, &rust.codec, did)
97            .await?
98            .with_context(|| format!("no OAuth session for {did}"))?;
99        if !oauth::token::is_stale(session.expires_at, now) {
100            return Ok(session);
101        }
102
103        // **`oauth_refresh` is timed from HERE, and that placement is the whole
104        // point.** It sits after the `is_stale` early return, so it still counts
105        // refreshes rather than every repo call — but it now also covers
106        // DISCOVERY, which runs only on this branch and is therefore part of the
107        // refresh.
108        //
109        // A review found the earlier placement (inside `valid_session`) missed
110        // every refresh that failed in discovery: an unreachable PDS, and the
111        // issuer-mismatch check below. Those are the two likeliest refresh
112        // outages in production, and they recorded nothing at all — leaving
113        // exactly the situation this metric exists to end, where the only trace
114        // is an error on whatever repo call happened to trigger it.
115        //
116        // The cost is that a caller which waits behind another task's refresh
117        // and then finds the session already fresh still records one. That is
118        // the right trade: it did perform a discovery round trip, and
119        // over-counting successes is harmless where under-counting failures is
120        // not.
121        let metrics = &self.state.metrics;
122        crate::metrics::timed(metrics, Backend::Rust, "oauth_refresh", async {
123            // Stale: now the token endpoint is genuinely needed. `valid_session`
124            // re-reads under the subject lock, so a concurrent refresh that lands
125            // between the check above and the lock below is handled there rather
126            // than here.
127            let server = oauth::discovery::discover(
128                &self.state.http,
129                &session.aud,
130                rust.auth_method.as_str(),
131                // The grant's own issuer. Checked inside `discover` now, so no
132                // caller can omit it — this one and the callback remembered, and
133                // revocation did not.
134                Some(&session.issuer),
135            )
136            .await?;
137
138            // **The re-discovered issuer must be the one this session was issued
139            // by.** This is the worse of the two instances of the same hole: the
140            // refresh path sends the REFRESH TOKEN — long-lived, and the credential
141            // that mints every other one — to whatever endpoint discovery returns,
142            // and the session's stored issuer was being compared against nothing.
143            //
144            // Discovery's own checks are all internally consistent, so a hostile
145            // pair of documents satisfies every one of them. `store.rs` names this
146            // exact threat as the reason `issuer` is AAD-bound; the AAD protects the
147            // column from local tampering, and only this protects it from a network
148            // re-read.
149            //
150
151            let ctx = oauth::session::RefreshContext {
152                token_endpoint: &server.token_endpoint,
153                client_id: &rust.client_id,
154                auth_method: rust.auth_method,
155                client_key: rust.client_key.as_ref(),
156                revocation_endpoint: server.revocation_endpoint.as_deref(),
157            };
158            oauth::session::valid_session(
159                &self.state.db,
160                &rust.codec,
161                &self.state.http,
162                &rust.locks,
163                did,
164                &ctx,
165                now,
166            )
167            .await
168        })
169        .await
170    }
171
172    /// The owned pieces an [`oauth::xrpc::Repo`] borrows.
173    ///
174    /// Returned owned rather than assembled here because `Repo` borrows both,
175    /// and a borrow cannot outlive the call that created what it points at. The
176    /// dispatch macro builds the handle in the caller's scope instead.
177    async fn rust_parts(
178        &self,
179        did: &str,
180    ) -> Result<(oauth::store::OAuthSession, oauth::keys::SigningKey)> {
181        let session = self.session(did).await?;
182        let key = oauth::keys::SigningKey::from_jwk_json(&session.dpop_key_jwk, "session")
183            .context("unsealing the session's DPoP key")?;
184        Ok((session, key))
185    }
186}
187
188/// Generate a dispatch method whose two arms take the same arguments.
189///
190/// A macro rather than twenty hand-written matches: the point of this module is
191/// that the two backends cannot drift, and a hand-written arm is exactly where
192/// an argument gets dropped or reordered on one side only.
193macro_rules! dispatch {
194    // Explicit visibility and metric label. Used by the three subscription
195    // writers, which are private here so that reaching them through `Repo` means
196    // going through the vetting wrapper of the same name; the label is passed
197    // explicitly so the private method's `_unvetted` suffix does not leak into
198    // the metric.
199    //
200    // **`Repo` is the vetted path. The layer below is narrower than it was,
201    // and still not closed.**
202    //
203    // The named writers demand vetted records: `add_subscription`,
204    // `update_subscription`, `add_subscriptions_bulk` and `add_saved` on both
205    // backends take `&vetted::Vetted*`, whose inner field is private to
206    // `src/vetted.rs` and reachable only through a constructor that vets.
207    //
208    // **That sentence has now been wrong twice, so read what follows as the
209    // correction rather than the guarantee.** #150 closed three routes and its
210    // message implied the class; #154 said "there is nothing to hand them that
211    // skipped the check" and a review found a fourth,
212    // `SidecarClient::create_subscriptions_batch`, taking a raw
213    // `&[Subscription]`. Deleting it, this comment then said "the layer below
214    // now demands vetted records too" — and a second review found the *general*
215    // case the deleted function had been one instance of: `create_record`,
216    // `put_record` and `apply_writes` are generic over `T: Serialize`, and
217    // `lexicon::Subscription` derives `Serialize`, so any of them writes an
218    // unvetted record. Verified by compiling it.
219    //
220    // What changed: those three are now PRIVATE on `PdsClient` and
221    // `SidecarClient` — not `pub(crate)`, which was the first attempt and
222    // which a self-review caught doing nothing for the actual threat: a
223    // handler in `web.rs` is in this crate, so `pub(crate)` left it fully
224    // able to call them. Verified both ways by compiling a probe from `web.rs`:
225    // `pub(crate)` → builds clean; private → three "private method" errors.
226    // The vetted wrappers sit in the same `impl` block and need no visibility.
227    //
228    // The three on `oauth::xrpc::Repo` remain `pub` because
229    // `examples/oauth_spike.rs` is a separate crate target and drives them
230    // against a scratch collection.
231    //
232    // They are no longer generic over every `Serialize`, though: `create_record`
233    // and `put_record` take `vetted::WritableRecord`, a sealed trait whose
234    // implementors are the vetted wrappers plus the two lexicon records with
235    // nothing to vet. `create_record(nsid::SUBSCRIPTION, &raw_subscription)`
236    // does not compile, and a `compile_fail` doctest on the trait says so.
237    //
238    // This comment used to name removing `Serialize` from
239    // `lexicon::Subscription` as the way to end the class, and call it blocked:
240    // `VettedSubscription` is `#[serde(transparent)]` over it, so removing the
241    // derive forces a hand-written impl, and a record-shape slip there would
242    // silently migrate every reader's repo. The trait gets the same guarantee
243    // with the wire format still coming from one derive.
244    //
245    // What remains reachable is `apply_writes`, whose `WriteOp::Create.value`
246    // is a `serde_json::Value` — a hand-built record, not an accidental
247    // one-liner. Narrowed, and the remainder named rather than implied.
248    (
249        $(#[$meta:meta])*
250        $vis:vis $name:ident ( $( $arg:ident : $ty:ty ),* ) -> $ret:ty,
251        label: $label:expr,
252        sidecar: $sidecar:ident,
253        rust: $rust:ident
254    ) => {
255        $(#[$meta])*
256        $vis async fn $name(&self, did: &str $(, $arg: $ty)*) -> Result<$ret> {
257            timed(&self.state.metrics, self.backend(), $label, async {
258                match self.backend() {
259                    Backend::Sidecar => self.state.sidecar.$sidecar(did $(, $arg)*).await,
260                    Backend::Rust => {
261                        let (session, key) = self.rust_parts(did).await?;
262                        let repo = oauth::xrpc::Repo {
263                            http: &self.state.http,
264                            pool: &self.state.db,
265                            session: &session,
266                            key: &key,
267                        };
268                        repo.$rust($($arg),*).await
269                    }
270                }
271            })
272            .await
273        }
274    };
275
276    // The common case: public, and the metric label is the method name. Forwards
277    // to the arm above rather than repeating the body — a second copy of the
278    // match is exactly the backend drift this macro exists to prevent.
279    (
280        $(#[$meta:meta])*
281        $name:ident ( $( $arg:ident : $ty:ty ),* ) -> $ret:ty,
282        sidecar: $sidecar:ident,
283        rust: $rust:ident
284    ) => {
285        dispatch! {
286            $(#[$meta])*
287            pub $name ( $( $arg : $ty ),* ) -> $ret,
288            label: stringify!($name),
289            sidecar: $sidecar,
290            rust: $rust
291        }
292    };
293}
294
295impl Repo<'_> {
296    // ── subscriptions ────────────────────────────────────────────────────────
297
298    dispatch! {
299        /// The reader's feed list, in display order.
300        list_subscriptions_sorted() -> Vec<(String, Subscription)>,
301        sidecar: list_subscriptions_sorted,
302        rust: list_subscriptions_sorted
303    }
304
305    dispatch! {
306        /// The reader's feed list with the CID each record was listed at,
307        /// unsorted — for a read-modify-write that puts with `swapRecord`
308        /// (#149). `list_subscriptions_sorted` is unchanged for every caller
309        /// that only displays the list.
310        list_subscriptions_with_cids() -> Vec<(String, Option<String>, Subscription)>,
311        sidecar: list_subscriptions_with_cids,
312        rust: list_subscriptions_with_cids
313    }
314
315    dispatch! {
316        /// Raw subscribe. Private: reach it through [`Repo::add_subscription`].
317        add_subscription_unvetted(sub: &crate::vetted::VettedSubscription) -> String,
318        label: "add_subscription",
319        sidecar: add_subscription,
320        rust: add_subscription
321    }
322
323    /// Subscribe. Returns the new record's rkey.
324    ///
325    /// Vets the record first — see [`crate::vetted::VettedSubscription`].
326    pub async fn add_subscription(&self, did: &str, sub: &Subscription) -> Result<String> {
327        self.add_subscription_unvetted(did, &VettedSubscription::new(sub))
328            .await
329    }
330
331    dispatch! {
332        /// Unsubscribe by rkey.
333        remove_subscription(rkey: &str) -> (),
334        sidecar: remove_subscription,
335        rust: remove_subscription
336    }
337
338    dispatch! {
339        /// Raw update. Private: reach it through [`Repo::update_subscription`].
340        update_subscription_unvetted(
341            rkey: &str,
342            sub: &crate::vetted::VettedSubscription,
343            swap_record: Option<&str>
344        ) -> crate::atproto::WriteResult,
345        label: "update_subscription",
346        sidecar: update_subscription,
347        rust: update_subscription
348    }
349
350    /// Rename or re-folder a subscription.
351    ///
352    /// Vets the record first — see [`crate::vetted::VettedSubscription`].
353    pub async fn update_subscription(
354        &self,
355        did: &str,
356        rkey: &str,
357        sub: &Subscription,
358        swap_record: Option<&str>,
359    ) -> Result<crate::atproto::WriteResult> {
360        self.update_subscription_unvetted(did, rkey, &VettedSubscription::new(sub), swap_record)
361            .await
362    }
363
364    dispatch! {
365        /// Raw bulk add. Private: reach it through [`Repo::add_subscriptions_bulk`].
366        add_subscriptions_bulk_unvetted(subs: &[crate::vetted::VettedSubscription]) -> Vec<String>,
367        label: "add_subscriptions_bulk",
368        sidecar: add_subscriptions_bulk,
369        rust: add_subscriptions_bulk
370    }
371
372    /// OPML import — `applyWrites`, one call per 200 feeds. A failure part-way
373    /// leaves a prefix written; see [`crate::atproto::ApplyWritesIncomplete`].
374    ///
375    /// Vets every record first — see [`crate::vetted::VettedSubscription::all`]. An import is the path where the
376    /// URLs are least trustworthy: the file is arbitrary, and 200 of them arrive
377    /// at once with nobody reading each line.
378    pub async fn add_subscriptions_bulk(
379        &self,
380        did: &str,
381        subs: &[Subscription],
382    ) -> Result<Vec<String>> {
383        self.add_subscriptions_bulk_unvetted(did, &VettedSubscription::all(subs))
384            .await
385    }
386
387    // ── folders ──────────────────────────────────────────────────────────────
388
389    dispatch! {
390        /// Folders in display order.
391        list_folders_sorted() -> Vec<(String, Folder)>,
392        sidecar: list_folders_sorted,
393        rust: list_folders_sorted
394    }
395
396    dispatch! {
397        /// Folders with the CID each record was listed at, unsorted — the read
398        /// half of a rename that puts with `swapRecord` (#268).
399        list_folders_with_cids() -> Vec<(String, Option<String>, Folder)>,
400        sidecar: list_folders_with_cids,
401        rust: list_folders_with_cids
402    }
403
404    dispatch! {
405        /// Create a folder. Returns its rkey.
406        add_folder(folder: &Folder) -> String,
407        sidecar: add_folder,
408        rust: add_folder
409    }
410
411    dispatch! {
412        /// Rename a folder in place, on the condition that it is still at
413        /// `swap_record` when that is given (#268).
414        rename_folder(
415            rkey: &str,
416            folder: &Folder,
417            swap_record: Option<&str>
418        ) -> crate::atproto::WriteResult,
419        sidecar: rename_folder,
420        rust: rename_folder
421    }
422
423    dispatch! {
424        /// Delete a folder.
425        remove_folder(rkey: &str) -> (),
426        sidecar: remove_folder,
427        rust: remove_folder
428    }
429
430    // ── saved ────────────────────────────────────────────────────────────────
431
432    dispatch! {
433        /// Saved items, newest first.
434        list_saved_sorted() -> Vec<(String, Saved)>,
435        sidecar: list_saved_sorted,
436        rust: list_saved_sorted
437    }
438
439    dispatch! {
440        /// Raw save. Private: reach it through [`Repo::add_saved`].
441        add_saved_unvetted(saved: &crate::vetted::VettedSaved) -> String,
442        label: "add_saved",
443        sidecar: add_saved,
444        rust: add_saved
445    }
446
447    /// Save an entry. Returns its rkey.
448    ///
449    /// **Refuses rather than publishing a URL we would not render.** `url` is
450    /// required on this record, so unlike a subscription's `siteUrl` there is no
451    /// honest resting place for a rejected value — see [`crate::vetted::VettedSaved`].
452    /// The only caller already treats a failed PDS write as recoverable: it logs
453    /// and keeps the entry starred locally, so the reader loses nothing but the
454    /// hostile record.
455    pub async fn add_saved(&self, did: &str, saved: &Saved) -> Result<String> {
456        self.add_saved_unvetted(did, &crate::vetted::VettedSaved::new(saved)?)
457            .await
458    }
459
460    dispatch! {
461        /// Saved items in PDS order — the un-star path, which matches by URL and
462        /// does not care about display order.
463        list_saved() -> Vec<(String, Saved)>,
464        sidecar: list_saved,
465        rust: list_saved
466    }
467
468    dispatch! {
469        /// Unsave by rkey.
470        remove_saved(rkey: &str) -> (),
471        sidecar: remove_saved,
472        rust: remove_saved
473    }
474
475    // ── read state ───────────────────────────────────────────────────────────
476
477    dispatch! {
478        /// Every read cursor.
479        list_read_states() -> Vec<(String, ReadState)>,
480        sidecar: list_read_states,
481        rust: list_read_states
482    }
483
484    dispatch! {
485        /// Every record in `collection`, unparsed — for asking which rkeys
486        /// EXIST, which the typed listers cannot answer: they skip a record
487        /// this build cannot parse, and a skipped record reads as a missing one.
488        list_all_records(collection: &str) -> Vec<crate::atproto::RecordEntry>,
489        sidecar: list_all_records,
490        rust: list_all_records
491    }
492
493    dispatch! {
494        /// Flush dirty cursors via `applyWrites`, chunked — the hottest write path.
495        /// A failure part-way leaves a prefix written; see
496        /// [`crate::atproto::ApplyWritesIncomplete`].
497        flush_read_states(cursors: &[(String, ReadState, bool)]) -> (),
498        sidecar: flush_read_states,
499        rust: flush_read_states
500    }
501}
502
503#[cfg(test)]
504mod tests {
505    use super::*;
506    use crate::config::Config;
507
508    const DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
509
510    async fn state_with(backend: Backend, public_url: &str) -> anyhow::Result<AppState> {
511        let db = crate::store::init_url("sqlite::memory:").await?;
512        AppState::new(
513            Config {
514                repo_backend: backend,
515                public_url: public_url.to_string(),
516                oauth: crate::config::OauthConfig {
517                    // **A unique path per test, not the default.**
518                    //
519                    // `key_path` defaults to the RELATIVE `oauth-signing-key.json`,
520                    // so a rust-backend test run writes real (encrypted) key
521                    // material into whatever the working directory happens to be
522                    // — the repo root — and every later run then tries to decrypt
523                    // a file written under a different key and fails to boot.
524                    // Ambient filesystem state is not a thing a test should depend
525                    // on, and this is the same relative-path foot-gun the README
526                    // documents for containers.
527                    key_path: std::env::temp_dir().join(format!(
528                        "fr-test-oauth-key-{}-{:p}.json",
529                        std::process::id(),
530                        &db as *const _
531                    )),
532                    encryption_key: Some("a".repeat(43)),
533                    ..crate::config::OauthConfig::default()
534                },
535                ..Config::default()
536            },
537            db,
538        )
539    }
540
541    /// A state pointed at a mock sidecar, so a repo write actually goes out.
542    async fn sidecar_state(internal_url: &str) -> anyhow::Result<AppState> {
543        let db = crate::store::init_url("sqlite::memory:").await?;
544        AppState::new(
545            Config {
546                repo_backend: Backend::Sidecar,
547                public_url: "http://localhost:8080".to_string(),
548                sidecar: crate::config::SidecarConfig {
549                    public_url: internal_url.to_string(),
550                    internal_url: internal_url.to_string(),
551                    internal_secret: "test-secret".to_string(),
552                },
553                ..Config::default()
554            },
555            db,
556        )
557    }
558
559    /// A sidecar mock that keeps the body of the first repo write it is sent.
560    ///
561    /// The assertion has to be made on the BYTES ON THE WIRE. Checking the
562    /// `Subscription` we passed in would pass just as happily with the vet
563    /// deleted — the record only becomes safe on the way out.
564    async fn spawn_capturing_sidecar() -> (String, std::sync::Arc<std::sync::Mutex<Vec<String>>>) {
565        use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
566        let seen = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
567        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
568        let addr = listener.local_addr().unwrap();
569        let sink = seen.clone();
570        tokio::spawn(async move {
571            loop {
572                let Ok((mut sock, _)) = listener.accept().await else {
573                    break;
574                };
575                // **Read until the body is complete, not until the first
576                // syscall returns.** A single `read` gets whatever one TCP
577                // segment carried; if headers and body arrive separately, the
578                // "body" is empty and every `!contains(...)` assertion below
579                // passes for the wrong reason — a false green in the one test
580                // the mutation argument rests on.
581                let mut raw: Vec<u8> = Vec::new();
582                let mut chunk = [0u8; 4096];
583                while let Ok(n) = sock.read(&mut chunk).await {
584                    if n == 0 {
585                        break;
586                    }
587                    raw.extend_from_slice(&chunk[..n]);
588                    let Some(split) = raw.windows(4).position(|w| w == b"\r\n\r\n") else {
589                        continue;
590                    };
591                    let (head, body) = raw.split_at(split + 4);
592                    let want = String::from_utf8_lossy(head).lines().find_map(|l| {
593                        let (k, v) = l.split_once(':')?;
594                        k.eq_ignore_ascii_case("content-length")
595                            .then(|| v.trim().parse::<usize>().ok())?
596                    });
597                    if want.is_none_or(|want| body.len() >= want) {
598                        sink.lock()
599                            .unwrap()
600                            .push(String::from_utf8_lossy(body).to_string());
601                        break;
602                    }
603                }
604                let body = serde_json::json!({
605                    "ok": true,
606                    "data": {
607                        "uri": "at://did:plc:x/community.lexicon.rss.subscription/rk1",
608                        "cid": "bafyreiabc"
609                    }
610                })
611                .to_string();
612                let resp = format!(
613                    "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}",
614                    body.len(),
615                    body
616                );
617                let _ = sock.write_all(resp.as_bytes()).await;
618                let _ = sock.flush().await;
619            }
620        });
621        (format!("http://{addr}"), seen)
622    }
623
624    fn sub_with_site(site: &str) -> Subscription {
625        let mut sub = Subscription::new("https://example.com/feed.xml", "2026-01-01T00:00:00.000Z");
626        sub.site_url = Some(site.to_string());
627        sub
628    }
629
630    /// Every write path, against the one scheme that motivated the guard.
631    ///
632    /// Parameterised over the three writers rather than testing one, because the
633    /// vet is applied per-wrapper: a fourth writer, or a wrapper that forgets the
634    /// call, is precisely the regression this is here to catch.
635    #[tokio::test]
636    async fn no_writer_publishes_a_hostile_site_url() {
637        for hostile in [
638            "javascript:alert(1)",
639            "data:text/html;base64,PHNjcmlwdD4=",
640            "  javascript:alert(1)  ",
641            "vbscript:msgbox(1)",
642        ] {
643            let (url, seen) = spawn_capturing_sidecar().await;
644            let state = sidecar_state(&url).await.expect("sidecar state");
645            let sub = sub_with_site(hostile);
646
647            let _ = state.repo().add_subscription(DID, &sub).await;
648            let _ = state
649                .repo()
650                .update_subscription(DID, "rk1", &sub, None)
651                .await;
652            let _ = state
653                .repo()
654                .add_subscriptions_bulk(DID, std::slice::from_ref(&sub))
655                .await;
656
657            let bodies = seen.lock().unwrap().clone();
658            assert_eq!(
659                bodies.len(),
660                3,
661                "expected one body per writer, got {bodies:?}"
662            );
663            for (writer, body) in ["add", "update", "bulk"].iter().zip(&bodies) {
664                // Anchor the negative assertions below: a truncated or empty
665                // capture would satisfy every `!contains(...)` vacuously.
666                assert!(
667                    body.contains("https://example.com/feed.xml"),
668                    "{writer} captured no usable body, so the assertions that \
669                     follow would pass for the wrong reason: {body:?}"
670                );
671                assert!(
672                    !body.contains("javascript:")
673                        && !body.contains("data:")
674                        && !body.contains("vbscript:"),
675                    "{writer} published {hostile:?} to the PDS: {body}"
676                );
677                assert!(
678                    !body.contains("siteUrl"),
679                    "{writer} sent a rejected siteUrl as an empty string; it must be \
680                     omitted, so other clients render no link rather than a broken one: \
681                     {body}"
682                );
683            }
684        }
685    }
686
687    /// **A hostile saved URL must not reach the PDS either.**
688    ///
689    /// `community.lexicon.rss.saved` publishes `url` into the reader's own repo,
690    /// into a field `safe_link.rs` already treats as attacker-controlled on the
691    /// RENDER side — that is what `SafeLink::external` exists for. The write side
692    /// had no equivalent: `add_saved` went straight through `dispatch!` with no
693    /// vet, so `entries.url` rows that arrived before the ingest guard existed,
694    /// or by any future path that does not go through `feed.rs`, were published
695    /// verbatim when the reader starred them.
696    ///
697    /// Asserted on the bytes on the wire, for the same reason as the
698    /// subscription writers: the record only becomes unsafe on the way out.
699    #[tokio::test]
700    async fn starring_does_not_publish_a_hostile_url() {
701        for hostile in [
702            "javascript:alert(1)",
703            "data:text/html;base64,PHNjcmlwdD4=",
704            "vbscript:msgbox(1)",
705        ] {
706            let (url, seen) = spawn_capturing_sidecar().await;
707            let state = sidecar_state(&url).await.expect("sidecar state");
708            let mut saved = Saved::new(hostile, "2026-01-01T00:00:00.000Z");
709            saved.title = Some("Hostile".to_string());
710
711            let _ = state.repo().add_saved(DID, &saved).await;
712
713            let bodies = seen.lock().unwrap().clone();
714            assert!(
715                !bodies.iter().any(|b| b.contains("javascript:")
716                    || b.contains("data:")
717                    || b.contains("vbscript:")),
718                "starring published {hostile:?} to the PDS: {bodies:?}"
719            );
720        }
721    }
722
723    /// The guard must not eat the ordinary case.
724    #[tokio::test]
725    async fn a_legitimate_site_url_is_published_unchanged() {
726        let (url, seen) = spawn_capturing_sidecar().await;
727        let state = sidecar_state(&url).await.expect("sidecar state");
728        let sub = sub_with_site("https://example.com/blog");
729
730        let _ = state.repo().add_subscription(DID, &sub).await;
731
732        let bodies = seen.lock().unwrap().clone();
733        assert_eq!(bodies.len(), 1, "expected one write, got {bodies:?}");
734        assert!(
735            bodies[0].contains("https://example.com/blog"),
736            "a perfectly good site link was dropped: {}",
737            bodies[0]
738        );
739    }
740
741    /// **Both arms must record under the SAME operation name.**
742    ///
743    /// The comparison is two rows in one table keyed by (backend, operation). If
744    /// the arms tagged their calls differently — a rename on one side, a typo on
745    /// the other — the table would show two half-populated sets of rows and no
746    /// pair would ever line up. The macro derives the name from the method via
747    /// `stringify!` precisely so this cannot drift, and this pins it.
748    ///
749    /// Neither call can succeed here (there is no sidecar and no session), which
750    /// is the point: the name is recorded either way, and a failed call is what
751    /// the error columns exist to show.
752    #[tokio::test]
753    async fn both_backends_record_under_the_same_operation_name() {
754        let sidecar = state_with(Backend::Sidecar, "http://localhost:8080")
755            .await
756            .expect("sidecar state");
757        let _ = sidecar.repo().list_subscriptions_sorted(DID).await;
758
759        let rust = state_with(Backend::Rust, "http://localhost:8080")
760            .await
761            .expect("rust state");
762        let _ = rust.repo().list_subscriptions_sorted(DID).await;
763
764        let sidecar_rows = sidecar.metrics.snapshot();
765        let rust_rows = rust.metrics.snapshot();
766        let names: Vec<&str> = sidecar_rows
767            .iter()
768            .chain(rust_rows.iter())
769            .map(|row| row.op.as_str())
770            .collect();
771
772        assert_eq!(
773            names.len(),
774            2,
775            "each backend should record exactly one call"
776        );
777        assert_eq!(
778            names[0], names[1],
779            "the two backends tagged the same operation differently, so their rows \
780             can never be compared"
781        );
782        assert_eq!(names[0], "list_subscriptions_sorted");
783    }
784
785    /// **The three hand-typed metric labels must stay what they say.**
786    ///
787    /// Every other method derives its label from its own name via `stringify!`,
788    /// which is what makes drift impossible for them. The subscription writers
789    /// cannot: the macro-generated method behind each wrapper is named
790    /// `*_unvetted`, and that suffix must not reach the metrics table, so the
791    /// label is passed as a literal instead. A literal is exactly the thing that
792    /// can drift, so it is pinned here.
793    ///
794    /// Honest about the limit: this pins the labels, not the correspondence
795    /// between a label and its wrapper's name. Renaming a public wrapper without
796    /// touching its literal would still slip through — the residual cost of
797    /// hand-typing three of the fifteen.
798    ///
799    /// None of the calls can succeed (no session), which is the point: the label
800    /// is recorded either way.
801    #[tokio::test]
802    async fn the_three_hand_typed_labels_are_what_they_claim() {
803        let state = state_with(Backend::Rust, "http://localhost:8080")
804            .await
805            .expect("rust state");
806        let sub = Subscription::new("https://example.com/feed.xml", "2026-01-01T00:00:00.000Z");
807
808        let _ = state.repo().add_subscription(DID, &sub).await;
809        let _ = state
810            .repo()
811            .update_subscription(DID, "rk1", &sub, None)
812            .await;
813        let _ = state
814            .repo()
815            .add_subscriptions_bulk(DID, std::slice::from_ref(&sub))
816            .await;
817
818        let mut ops: Vec<String> = state
819            .metrics
820            .snapshot()
821            .into_iter()
822            .map(|row| row.op)
823            .collect();
824        ops.sort();
825        assert_eq!(
826            ops,
827            vec![
828                "add_subscription".to_string(),
829                "add_subscriptions_bulk".to_string(),
830                "update_subscription".to_string(),
831            ],
832            "a writer's metric label drifted from the operation it names, so its \
833             rows can never be compared against the other backend's"
834        );
835    }
836
837    /// A failed call is still recorded — as a FAILURE, not as a fast success.
838    #[tokio::test]
839    async fn a_failed_call_is_recorded_in_the_error_column() {
840        let state = state_with(Backend::Rust, "http://localhost:8080")
841            .await
842            .expect("rust state");
843        let result = state.repo().list_subscriptions_sorted(DID).await;
844        assert!(result.is_err(), "there is no session, so this must fail");
845
846        let snapshot = state.metrics.snapshot();
847        let stats = &snapshot[0].stats;
848        assert_eq!(stats.err_count, 1);
849        assert_eq!(
850            stats.ok_count, 0,
851            "a failure was counted as a success, which is exactly the reading \
852             that makes a broken backend look fast"
853        );
854        assert_eq!(stats.percentile(50.0), None, "no successes, so no p50");
855    }
856
857    /// **Selecting the Rust backend with an unusable OAuth config must not boot.**
858    ///
859    /// The alternative is a server that starts cleanly and then fails every
860    /// single repo call at request time — which looks like a PDS outage rather
861    /// than a configuration error, on a path the operator has just switched to.
862    #[tokio::test]
863    async fn the_rust_backend_refuses_to_start_on_an_unusable_oauth_config() {
864        // A public URL with a path is rejected by `ClientConfig::new` (it would
865        // publish a doubled client_id).
866        let err = match state_with(Backend::Rust, "https://feather-reader.com/oauth").await {
867            Err(err) => err,
868            Ok(_) => panic!("the rust backend booted with an unusable OAuth config"),
869        };
870        let rendered = format!("{err:#}");
871        assert!(
872            rendered.contains("FEATHERREADER_REPO_BACKEND=rust"),
873            "the error must name the switch that caused it: {rendered}"
874        );
875    }
876
877    /// The same bad config with the SIDECAR selected still boots: the Rust
878    /// runtime is unused, and refusing to start would block a rollback.
879    #[tokio::test]
880    async fn the_sidecar_backend_still_boots_with_an_unusable_oauth_config() {
881        let state = state_with(Backend::Sidecar, "https://feather-reader.com/oauth")
882            .await
883            .expect("the sidecar path must not be blocked by Rust-only config");
884        assert!(state.oauth.is_none());
885    }
886    /// **A refresh that fails in DISCOVERY must be counted.**
887    ///
888    /// Regression test for the defect a review found in the first version of
889    /// this metric: the span sat inside `valid_session`, but `Repo::session`
890    /// runs discovery BEFORE that — only on the stale branch, so discovery is
891    /// part of the refresh — and a failure there propagated via `?` recording
892    /// nothing at all.
893    ///
894    /// That silently excluded the two likeliest refresh outages: an unreachable
895    /// PDS, and the issuer-mismatch check. Those are precisely the events
896    /// `err_count` exists to move on, and they left the metric flat while the
897    /// error surfaced only on whatever repo call happened to trigger it — the
898    /// exact situation this work set out to end.
899    #[tokio::test]
900    async fn a_refresh_that_fails_in_discovery_is_counted() {
901        let state = state_with(Backend::Rust, "https://feather-reader.com")
902            .await
903            .expect("state");
904        let runtime = state.oauth.as_deref().expect("oauth runtime");
905        let did = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
906
907        // EXPIRED, so `Repo::session` takes the stale branch and reaches
908        // discovery. `pds.invalid` cannot resolve, so discovery fails.
909        oauth::store::put_session(
910            &state.db,
911            &runtime.codec,
912            &oauth::store::OAuthSession {
913                sub: did.into(),
914                issuer: "https://auth.invalid".into(),
915                aud: "https://pds.invalid".into(),
916                dpop_key_jwk: oauth::keys::SigningKey::generate("session-dpop")
917                    .to_jwk_json()
918                    .unwrap(),
919                access_token: "at".into(),
920                refresh_token: "rt".into(),
921                token_type: "DPoP".into(),
922                granted_scope: "atproto".into(),
923                expires_at: Some(crate::store::now_unix() - 1),
924            },
925        )
926        .await
927        .unwrap();
928
929        let err = state.repo().session(did).await;
930        assert!(err.is_err(), "an unreachable PDS must fail the refresh");
931
932        let row = state
933            .metrics
934            .snapshot()
935            .into_iter()
936            .find(|r| r.op == "oauth_refresh" && r.backend == Backend::Rust)
937            .expect("a refresh that failed in discovery recorded nothing at all");
938        assert_eq!(row.stats.err_count, 1);
939        assert_eq!(row.stats.ok_count, 0);
940    }
941
942    /// A session that is still fresh must record NO refresh — the property that
943    /// keeps the metric meaningful, since `session()` runs on every repo call.
944    #[tokio::test]
945    async fn a_fresh_session_records_no_refresh() {
946        let state = state_with(Backend::Rust, "https://feather-reader.com")
947            .await
948            .expect("state");
949        let runtime = state.oauth.as_deref().expect("oauth runtime");
950        let did = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
951        oauth::store::put_session(
952            &state.db,
953            &runtime.codec,
954            &oauth::store::OAuthSession {
955                sub: did.into(),
956                issuer: "https://auth.invalid".into(),
957                aud: "https://pds.invalid".into(),
958                dpop_key_jwk: oauth::keys::SigningKey::generate("session-dpop")
959                    .to_jwk_json()
960                    .unwrap(),
961                access_token: "at".into(),
962                refresh_token: "rt".into(),
963                token_type: "DPoP".into(),
964                granted_scope: "atproto".into(),
965                expires_at: Some(crate::store::now_unix() + 3600),
966            },
967        )
968        .await
969        .unwrap();
970
971        state.repo().session(did).await.expect("fresh session");
972        assert!(
973            state
974                .metrics
975                .snapshot()
976                .iter()
977                .all(|r| r.op != "oauth_refresh"),
978            "a fresh session recorded a refresh it never performed",
979        );
980    }
981}