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