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 — `applyWrites`, one call per 200 feeds. A failure part-way
358    /// leaves a prefix written; see [`crate::atproto::ApplyWritesIncomplete`].
359    ///
360    /// Vets every record first — see [`crate::vetted::VettedSubscription::all`]. An import is the path where the
361    /// URLs are least trustworthy: the file is arbitrary, and 200 of them arrive
362    /// at once with nobody reading each line.
363    pub async fn add_subscriptions_bulk(
364        &self,
365        did: &str,
366        subs: &[Subscription],
367    ) -> Result<Vec<String>> {
368        self.add_subscriptions_bulk_unvetted(did, &VettedSubscription::all(subs))
369            .await
370    }
371
372    // ── folders ──────────────────────────────────────────────────────────────
373
374    dispatch! {
375        /// Folders in display order.
376        list_folders_sorted() -> Vec<(String, Folder)>,
377        sidecar: list_folders_sorted,
378        rust: list_folders_sorted
379    }
380
381    dispatch! {
382        /// Create a folder. Returns its rkey.
383        add_folder(folder: &Folder) -> String,
384        sidecar: add_folder,
385        rust: add_folder
386    }
387
388    dispatch! {
389        /// Rename a folder in place.
390        rename_folder(rkey: &str, folder: &Folder) -> crate::atproto::WriteResult,
391        sidecar: rename_folder,
392        rust: rename_folder
393    }
394
395    dispatch! {
396        /// Delete a folder.
397        remove_folder(rkey: &str) -> (),
398        sidecar: remove_folder,
399        rust: remove_folder
400    }
401
402    // ── saved ────────────────────────────────────────────────────────────────
403
404    dispatch! {
405        /// Saved items, newest first.
406        list_saved_sorted() -> Vec<(String, Saved)>,
407        sidecar: list_saved_sorted,
408        rust: list_saved_sorted
409    }
410
411    dispatch! {
412        /// Raw save. Private: reach it through [`Repo::add_saved`].
413        add_saved_unvetted(saved: &crate::vetted::VettedSaved) -> String,
414        label: "add_saved",
415        sidecar: add_saved,
416        rust: add_saved
417    }
418
419    /// Save an entry. Returns its rkey.
420    ///
421    /// **Refuses rather than publishing a URL we would not render.** `url` is
422    /// required on this record, so unlike a subscription's `siteUrl` there is no
423    /// honest resting place for a rejected value — see [`crate::vetted::VettedSaved`].
424    /// The only caller already treats a failed PDS write as recoverable: it logs
425    /// and keeps the entry starred locally, so the reader loses nothing but the
426    /// hostile record.
427    pub async fn add_saved(&self, did: &str, saved: &Saved) -> Result<String> {
428        self.add_saved_unvetted(did, &crate::vetted::VettedSaved::new(saved)?)
429            .await
430    }
431
432    dispatch! {
433        /// Saved items in PDS order — the un-star path, which matches by URL and
434        /// does not care about display order.
435        list_saved() -> Vec<(String, Saved)>,
436        sidecar: list_saved,
437        rust: list_saved
438    }
439
440    dispatch! {
441        /// Unsave by rkey.
442        remove_saved(rkey: &str) -> (),
443        sidecar: remove_saved,
444        rust: remove_saved
445    }
446
447    // ── read state ───────────────────────────────────────────────────────────
448
449    dispatch! {
450        /// Every read cursor.
451        list_read_states() -> Vec<(String, ReadState)>,
452        sidecar: list_read_states,
453        rust: list_read_states
454    }
455
456    dispatch! {
457        /// Every record in `collection`, unparsed — for asking which rkeys
458        /// EXIST, which the typed listers cannot answer: they skip a record
459        /// this build cannot parse, and a skipped record reads as a missing one.
460        list_all_records(collection: &str) -> Vec<crate::atproto::RecordEntry>,
461        sidecar: list_all_records,
462        rust: list_all_records
463    }
464
465    dispatch! {
466        /// Flush dirty cursors via `applyWrites`, chunked — the hottest write path.
467        /// A failure part-way leaves a prefix written; see
468        /// [`crate::atproto::ApplyWritesIncomplete`].
469        flush_read_states(cursors: &[(String, ReadState, bool)]) -> (),
470        sidecar: flush_read_states,
471        rust: flush_read_states
472    }
473}
474
475#[cfg(test)]
476mod tests {
477    use super::*;
478    use crate::config::Config;
479
480    const DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
481
482    async fn state_with(backend: Backend, public_url: &str) -> anyhow::Result<AppState> {
483        let db = crate::store::init_url("sqlite::memory:").await?;
484        AppState::new(
485            Config {
486                repo_backend: backend,
487                public_url: public_url.to_string(),
488                oauth: crate::config::OauthConfig {
489                    // **A unique path per test, not the default.**
490                    //
491                    // `key_path` defaults to the RELATIVE `oauth-signing-key.json`,
492                    // so a rust-backend test run writes real (encrypted) key
493                    // material into whatever the working directory happens to be
494                    // — the repo root — and every later run then tries to decrypt
495                    // a file written under a different key and fails to boot.
496                    // Ambient filesystem state is not a thing a test should depend
497                    // on, and this is the same relative-path foot-gun the README
498                    // documents for containers.
499                    key_path: std::env::temp_dir().join(format!(
500                        "fr-test-oauth-key-{}-{:p}.json",
501                        std::process::id(),
502                        &db as *const _
503                    )),
504                    encryption_key: Some("a".repeat(43)),
505                    ..crate::config::OauthConfig::default()
506                },
507                ..Config::default()
508            },
509            db,
510        )
511    }
512
513    /// A state pointed at a mock sidecar, so a repo write actually goes out.
514    async fn sidecar_state(internal_url: &str) -> anyhow::Result<AppState> {
515        let db = crate::store::init_url("sqlite::memory:").await?;
516        AppState::new(
517            Config {
518                repo_backend: Backend::Sidecar,
519                public_url: "http://localhost:8080".to_string(),
520                sidecar: crate::config::SidecarConfig {
521                    public_url: internal_url.to_string(),
522                    internal_url: internal_url.to_string(),
523                    internal_secret: "test-secret".to_string(),
524                },
525                ..Config::default()
526            },
527            db,
528        )
529    }
530
531    /// A sidecar mock that keeps the body of the first repo write it is sent.
532    ///
533    /// The assertion has to be made on the BYTES ON THE WIRE. Checking the
534    /// `Subscription` we passed in would pass just as happily with the vet
535    /// deleted — the record only becomes safe on the way out.
536    async fn spawn_capturing_sidecar() -> (String, std::sync::Arc<std::sync::Mutex<Vec<String>>>) {
537        use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
538        let seen = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
539        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
540        let addr = listener.local_addr().unwrap();
541        let sink = seen.clone();
542        tokio::spawn(async move {
543            loop {
544                let Ok((mut sock, _)) = listener.accept().await else {
545                    break;
546                };
547                // **Read until the body is complete, not until the first
548                // syscall returns.** A single `read` gets whatever one TCP
549                // segment carried; if headers and body arrive separately, the
550                // "body" is empty and every `!contains(...)` assertion below
551                // passes for the wrong reason — a false green in the one test
552                // the mutation argument rests on.
553                let mut raw: Vec<u8> = Vec::new();
554                let mut chunk = [0u8; 4096];
555                while let Ok(n) = sock.read(&mut chunk).await {
556                    if n == 0 {
557                        break;
558                    }
559                    raw.extend_from_slice(&chunk[..n]);
560                    let Some(split) = raw.windows(4).position(|w| w == b"\r\n\r\n") else {
561                        continue;
562                    };
563                    let (head, body) = raw.split_at(split + 4);
564                    let want = String::from_utf8_lossy(head).lines().find_map(|l| {
565                        let (k, v) = l.split_once(':')?;
566                        k.eq_ignore_ascii_case("content-length")
567                            .then(|| v.trim().parse::<usize>().ok())?
568                    });
569                    if want.is_none_or(|want| body.len() >= want) {
570                        sink.lock()
571                            .unwrap()
572                            .push(String::from_utf8_lossy(body).to_string());
573                        break;
574                    }
575                }
576                let body = serde_json::json!({
577                    "ok": true,
578                    "data": {
579                        "uri": "at://did:plc:x/community.lexicon.rss.subscription/rk1",
580                        "cid": "bafyreiabc"
581                    }
582                })
583                .to_string();
584                let resp = format!(
585                    "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}",
586                    body.len(),
587                    body
588                );
589                let _ = sock.write_all(resp.as_bytes()).await;
590                let _ = sock.flush().await;
591            }
592        });
593        (format!("http://{addr}"), seen)
594    }
595
596    fn sub_with_site(site: &str) -> Subscription {
597        let mut sub = Subscription::new("https://example.com/feed.xml", "2026-01-01T00:00:00.000Z");
598        sub.site_url = Some(site.to_string());
599        sub
600    }
601
602    /// Every write path, against the one scheme that motivated the guard.
603    ///
604    /// Parameterised over the three writers rather than testing one, because the
605    /// vet is applied per-wrapper: a fourth writer, or a wrapper that forgets the
606    /// call, is precisely the regression this is here to catch.
607    #[tokio::test]
608    async fn no_writer_publishes_a_hostile_site_url() {
609        for hostile in [
610            "javascript:alert(1)",
611            "data:text/html;base64,PHNjcmlwdD4=",
612            "  javascript:alert(1)  ",
613            "vbscript:msgbox(1)",
614        ] {
615            let (url, seen) = spawn_capturing_sidecar().await;
616            let state = sidecar_state(&url).await.expect("sidecar state");
617            let sub = sub_with_site(hostile);
618
619            let _ = state.repo().add_subscription(DID, &sub).await;
620            let _ = state.repo().update_subscription(DID, "rk1", &sub).await;
621            let _ = state
622                .repo()
623                .add_subscriptions_bulk(DID, std::slice::from_ref(&sub))
624                .await;
625
626            let bodies = seen.lock().unwrap().clone();
627            assert_eq!(
628                bodies.len(),
629                3,
630                "expected one body per writer, got {bodies:?}"
631            );
632            for (writer, body) in ["add", "update", "bulk"].iter().zip(&bodies) {
633                // Anchor the negative assertions below: a truncated or empty
634                // capture would satisfy every `!contains(...)` vacuously.
635                assert!(
636                    body.contains("https://example.com/feed.xml"),
637                    "{writer} captured no usable body, so the assertions that \
638                     follow would pass for the wrong reason: {body:?}"
639                );
640                assert!(
641                    !body.contains("javascript:")
642                        && !body.contains("data:")
643                        && !body.contains("vbscript:"),
644                    "{writer} published {hostile:?} to the PDS: {body}"
645                );
646                assert!(
647                    !body.contains("siteUrl"),
648                    "{writer} sent a rejected siteUrl as an empty string; it must be \
649                     omitted, so other clients render no link rather than a broken one: \
650                     {body}"
651                );
652            }
653        }
654    }
655
656    /// **A hostile saved URL must not reach the PDS either.**
657    ///
658    /// `community.lexicon.rss.saved` publishes `url` into the reader's own repo,
659    /// into a field `safe_link.rs` already treats as attacker-controlled on the
660    /// RENDER side — that is what `SafeLink::external` exists for. The write side
661    /// had no equivalent: `add_saved` went straight through `dispatch!` with no
662    /// vet, so `entries.url` rows that arrived before the ingest guard existed,
663    /// or by any future path that does not go through `feed.rs`, were published
664    /// verbatim when the reader starred them.
665    ///
666    /// Asserted on the bytes on the wire, for the same reason as the
667    /// subscription writers: the record only becomes unsafe on the way out.
668    #[tokio::test]
669    async fn starring_does_not_publish_a_hostile_url() {
670        for hostile in [
671            "javascript:alert(1)",
672            "data:text/html;base64,PHNjcmlwdD4=",
673            "vbscript:msgbox(1)",
674        ] {
675            let (url, seen) = spawn_capturing_sidecar().await;
676            let state = sidecar_state(&url).await.expect("sidecar state");
677            let mut saved = Saved::new(hostile, "2026-01-01T00:00:00.000Z");
678            saved.title = Some("Hostile".to_string());
679
680            let _ = state.repo().add_saved(DID, &saved).await;
681
682            let bodies = seen.lock().unwrap().clone();
683            assert!(
684                !bodies.iter().any(|b| b.contains("javascript:")
685                    || b.contains("data:")
686                    || b.contains("vbscript:")),
687                "starring published {hostile:?} to the PDS: {bodies:?}"
688            );
689        }
690    }
691
692    /// The guard must not eat the ordinary case.
693    #[tokio::test]
694    async fn a_legitimate_site_url_is_published_unchanged() {
695        let (url, seen) = spawn_capturing_sidecar().await;
696        let state = sidecar_state(&url).await.expect("sidecar state");
697        let sub = sub_with_site("https://example.com/blog");
698
699        let _ = state.repo().add_subscription(DID, &sub).await;
700
701        let bodies = seen.lock().unwrap().clone();
702        assert_eq!(bodies.len(), 1, "expected one write, got {bodies:?}");
703        assert!(
704            bodies[0].contains("https://example.com/blog"),
705            "a perfectly good site link was dropped: {}",
706            bodies[0]
707        );
708    }
709
710    /// **Both arms must record under the SAME operation name.**
711    ///
712    /// The comparison is two rows in one table keyed by (backend, operation). If
713    /// the arms tagged their calls differently — a rename on one side, a typo on
714    /// the other — the table would show two half-populated sets of rows and no
715    /// pair would ever line up. The macro derives the name from the method via
716    /// `stringify!` precisely so this cannot drift, and this pins it.
717    ///
718    /// Neither call can succeed here (there is no sidecar and no session), which
719    /// is the point: the name is recorded either way, and a failed call is what
720    /// the error columns exist to show.
721    #[tokio::test]
722    async fn both_backends_record_under_the_same_operation_name() {
723        let sidecar = state_with(Backend::Sidecar, "http://localhost:8080")
724            .await
725            .expect("sidecar state");
726        let _ = sidecar.repo().list_subscriptions_sorted(DID).await;
727
728        let rust = state_with(Backend::Rust, "http://localhost:8080")
729            .await
730            .expect("rust state");
731        let _ = rust.repo().list_subscriptions_sorted(DID).await;
732
733        let sidecar_rows = sidecar.metrics.snapshot();
734        let rust_rows = rust.metrics.snapshot();
735        let names: Vec<&str> = sidecar_rows
736            .iter()
737            .chain(rust_rows.iter())
738            .map(|row| row.op.as_str())
739            .collect();
740
741        assert_eq!(
742            names.len(),
743            2,
744            "each backend should record exactly one call"
745        );
746        assert_eq!(
747            names[0], names[1],
748            "the two backends tagged the same operation differently, so their rows \
749             can never be compared"
750        );
751        assert_eq!(names[0], "list_subscriptions_sorted");
752    }
753
754    /// **The three hand-typed metric labels must stay what they say.**
755    ///
756    /// Every other method derives its label from its own name via `stringify!`,
757    /// which is what makes drift impossible for them. The subscription writers
758    /// cannot: the macro-generated method behind each wrapper is named
759    /// `*_unvetted`, and that suffix must not reach the metrics table, so the
760    /// label is passed as a literal instead. A literal is exactly the thing that
761    /// can drift, so it is pinned here.
762    ///
763    /// Honest about the limit: this pins the labels, not the correspondence
764    /// between a label and its wrapper's name. Renaming a public wrapper without
765    /// touching its literal would still slip through — the residual cost of
766    /// hand-typing three of the fifteen.
767    ///
768    /// None of the calls can succeed (no session), which is the point: the label
769    /// is recorded either way.
770    #[tokio::test]
771    async fn the_three_hand_typed_labels_are_what_they_claim() {
772        let state = state_with(Backend::Rust, "http://localhost:8080")
773            .await
774            .expect("rust state");
775        let sub = Subscription::new("https://example.com/feed.xml", "2026-01-01T00:00:00.000Z");
776
777        let _ = state.repo().add_subscription(DID, &sub).await;
778        let _ = state.repo().update_subscription(DID, "rk1", &sub).await;
779        let _ = state
780            .repo()
781            .add_subscriptions_bulk(DID, std::slice::from_ref(&sub))
782            .await;
783
784        let mut ops: Vec<String> = state
785            .metrics
786            .snapshot()
787            .into_iter()
788            .map(|row| row.op)
789            .collect();
790        ops.sort();
791        assert_eq!(
792            ops,
793            vec![
794                "add_subscription".to_string(),
795                "add_subscriptions_bulk".to_string(),
796                "update_subscription".to_string(),
797            ],
798            "a writer's metric label drifted from the operation it names, so its \
799             rows can never be compared against the other backend's"
800        );
801    }
802
803    /// A failed call is still recorded — as a FAILURE, not as a fast success.
804    #[tokio::test]
805    async fn a_failed_call_is_recorded_in_the_error_column() {
806        let state = state_with(Backend::Rust, "http://localhost:8080")
807            .await
808            .expect("rust state");
809        let result = state.repo().list_subscriptions_sorted(DID).await;
810        assert!(result.is_err(), "there is no session, so this must fail");
811
812        let snapshot = state.metrics.snapshot();
813        let stats = &snapshot[0].stats;
814        assert_eq!(stats.err_count, 1);
815        assert_eq!(
816            stats.ok_count, 0,
817            "a failure was counted as a success, which is exactly the reading \
818             that makes a broken backend look fast"
819        );
820        assert_eq!(stats.percentile(50.0), None, "no successes, so no p50");
821    }
822
823    /// **Selecting the Rust backend with an unusable OAuth config must not boot.**
824    ///
825    /// The alternative is a server that starts cleanly and then fails every
826    /// single repo call at request time — which looks like a PDS outage rather
827    /// than a configuration error, on a path the operator has just switched to.
828    #[tokio::test]
829    async fn the_rust_backend_refuses_to_start_on_an_unusable_oauth_config() {
830        // A public URL with a path is rejected by `ClientConfig::new` (it would
831        // publish a doubled client_id).
832        let err = match state_with(Backend::Rust, "https://feather-reader.com/oauth").await {
833            Err(err) => err,
834            Ok(_) => panic!("the rust backend booted with an unusable OAuth config"),
835        };
836        let rendered = format!("{err:#}");
837        assert!(
838            rendered.contains("FEATHERREADER_REPO_BACKEND=rust"),
839            "the error must name the switch that caused it: {rendered}"
840        );
841    }
842
843    /// The same bad config with the SIDECAR selected still boots: the Rust
844    /// runtime is unused, and refusing to start would block a rollback.
845    #[tokio::test]
846    async fn the_sidecar_backend_still_boots_with_an_unusable_oauth_config() {
847        let state = state_with(Backend::Sidecar, "https://feather-reader.com/oauth")
848            .await
849            .expect("the sidecar path must not be blocked by Rust-only config");
850        assert!(state.oauth.is_none());
851    }
852    /// **A refresh that fails in DISCOVERY must be counted.**
853    ///
854    /// Regression test for the defect a review found in the first version of
855    /// this metric: the span sat inside `valid_session`, but `Repo::session`
856    /// runs discovery BEFORE that — only on the stale branch, so discovery is
857    /// part of the refresh — and a failure there propagated via `?` recording
858    /// nothing at all.
859    ///
860    /// That silently excluded the two likeliest refresh outages: an unreachable
861    /// PDS, and the issuer-mismatch check. Those are precisely the events
862    /// `err_count` exists to move on, and they left the metric flat while the
863    /// error surfaced only on whatever repo call happened to trigger it — the
864    /// exact situation this work set out to end.
865    #[tokio::test]
866    async fn a_refresh_that_fails_in_discovery_is_counted() {
867        let state = state_with(Backend::Rust, "https://feather-reader.com")
868            .await
869            .expect("state");
870        let runtime = state.oauth.as_deref().expect("oauth runtime");
871        let did = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
872
873        // EXPIRED, so `Repo::session` takes the stale branch and reaches
874        // discovery. `pds.invalid` cannot resolve, so discovery fails.
875        oauth::store::put_session(
876            &state.db,
877            &runtime.codec,
878            &oauth::store::OAuthSession {
879                sub: did.into(),
880                issuer: "https://auth.invalid".into(),
881                aud: "https://pds.invalid".into(),
882                dpop_key_jwk: oauth::keys::SigningKey::generate("session-dpop")
883                    .to_jwk_json()
884                    .unwrap(),
885                access_token: "at".into(),
886                refresh_token: "rt".into(),
887                token_type: "DPoP".into(),
888                granted_scope: "atproto".into(),
889                expires_at: Some(crate::store::now_unix() - 1),
890            },
891        )
892        .await
893        .unwrap();
894
895        let err = state.repo().session(did).await;
896        assert!(err.is_err(), "an unreachable PDS must fail the refresh");
897
898        let row = state
899            .metrics
900            .snapshot()
901            .into_iter()
902            .find(|r| r.op == "oauth_refresh" && r.backend == Backend::Rust)
903            .expect("a refresh that failed in discovery recorded nothing at all");
904        assert_eq!(row.stats.err_count, 1);
905        assert_eq!(row.stats.ok_count, 0);
906    }
907
908    /// A session that is still fresh must record NO refresh — the property that
909    /// keeps the metric meaningful, since `session()` runs on every repo call.
910    #[tokio::test]
911    async fn a_fresh_session_records_no_refresh() {
912        let state = state_with(Backend::Rust, "https://feather-reader.com")
913            .await
914            .expect("state");
915        let runtime = state.oauth.as_deref().expect("oauth runtime");
916        let did = "did:plc:ewvi7nxzyoun6zhxrhs64oiz";
917        oauth::store::put_session(
918            &state.db,
919            &runtime.codec,
920            &oauth::store::OAuthSession {
921                sub: did.into(),
922                issuer: "https://auth.invalid".into(),
923                aud: "https://pds.invalid".into(),
924                dpop_key_jwk: oauth::keys::SigningKey::generate("session-dpop")
925                    .to_jwk_json()
926                    .unwrap(),
927                access_token: "at".into(),
928                refresh_token: "rt".into(),
929                token_type: "DPoP".into(),
930                granted_scope: "atproto".into(),
931                expires_at: Some(crate::store::now_unix() + 3600),
932            },
933        )
934        .await
935        .unwrap();
936
937        state.repo().session(did).await.expect("fresh session");
938        assert!(
939            state
940                .metrics
941                .snapshot()
942                .iter()
943                .all(|r| r.op != "oauth_refresh"),
944            "a fresh session recorded a refresh it never performed",
945        );
946    }
947}