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