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}