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