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