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}