feather_reader/oauth/request.rs
1//! DPoP-authenticated form POSTs, with nonce persistence and a bounded retry.
2//!
3//! Every OAuth POST this client makes goes through here: PAR, token exchange,
4//! refresh. Three behaviours matter and all three are easy to get subtly wrong:
5//!
6//! * **The nonce is persisted per origin and harvested from EVERY response**,
7//! including successes. Using one only for an immediate retry means every
8//! request pays a wasted round trip.
9//! * **The retry is bounded at one.** A server that answers every request with
10//! `use_dpop_nonce` would otherwise spin forever.
11//! * **The endpoint kind is passed, not inferred.** The authorization server
12//! signals a nonce requirement with `400` + a JSON body; a resource server
13//! uses `401` + `WWW-Authenticate`. Reading only one of those misses every
14//! challenge from the other.
15
16use anyhow::{Context as _, Result};
17use reqwest::Client;
18use sqlx::SqlitePool;
19
20use super::discovery::origin_of;
21use super::dpop::{self, Endpoint};
22use super::keys::SigningKey;
23use super::store;
24use crate::net;
25
26/// A completed request: what the server said, and what it said it with.
27pub struct PostOutcome {
28 pub status: u16,
29 /// The raw body. **Public, so the `Debug` impl below is a default rather
30 /// than a barrier**: `{:?}` on this field prints everything. On a success
31 /// this holds the access and refresh tokens, so it must not be logged,
32 /// formatted into an error, or echoed on any path that is not already known
33 /// to be a failure response.
34 pub body: Vec<u8>,
35}
36
37/// Hand-written, NOT derived. A token-endpoint body holds the access and refresh
38/// tokens; a derived `Debug` would put them into any log line, panic message or
39/// `{:?}` that ever touches this value.
40impl std::fmt::Debug for PostOutcome {
41 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
42 f.debug_struct("PostOutcome")
43 .field("status", &self.status)
44 .field(
45 "body",
46 &format_args!("<{} bytes redacted>", self.body.len()),
47 )
48 .finish()
49 }
50}
51
52/// The most nodes a response on THIS funnel may build.
53///
54/// **`MAX_LIST_STRUCTURAL_CHARS` is the wrong number here, and a review was right about
55/// that.** 640 000 was sized by the densest page the lexicons permit — a full
56/// `readState` listing, 403 003 counted characters — and at the measured 210 B
57/// per counted character it admits about 134 MB. Nothing on this funnel is a
58/// listing: it carries the token response, the PAR response, the session refresh,
59/// and the repo writers' results.
60///
61/// Measured, so the headroom is a number rather than a feeling. The largest
62/// legitimate body here is an `applyWrites` result set, and the batch is bounded
63/// by `max_subs_per_did` (default 500) or by the reader's subscribed feeds:
64///
65/// | body | wire | counted |
66/// |---|---|---|
67/// | token response | 1.2 kB | 23 |
68/// | `applyWrites`, 1 result | 348 B | 31 |
69/// | `applyWrites`, 500 results | 120 kB | 8 514 |
70/// | `applyWrites`, 2 000 results | 480 kB | 34 014 |
71///
72/// 64 000 leaves 7.5x over a 500-op batch, still covers a 2 000-op one, and is
73/// ten times tighter than the listing cap — about 13 MB instead of 134.
74const MAX_RESPONSE_NODES: usize = 64_000;
75
76/// Refuse a response body that would build more than [`MAX_RESPONSE_NODES`].
77///
78/// Message carries the status and the counts only — never the body, not even an
79/// excerpt, because a token response passes through here.
80fn refuse_a_response_explosion(body: &[u8], status: u16) -> Result<()> {
81 let nodes = crate::atproto::count_structural_chars(body);
82 anyhow::ensure!(
83 nodes <= MAX_RESPONSE_NODES,
84 "the response body (status {status}) counts at least {nodes} structural \
85 characters, over the {MAX_RESPONSE_NODES} cap for this endpoint — refusing \
86 before parsing it"
87 );
88 Ok(())
89}
90
91impl PostOutcome {
92 /// Parse the body as JSON.
93 ///
94 /// The body is NEVER echoed into the error, not even an excerpt. This is
95 /// called on the SUCCESSFUL token response, so a 200 that fails to parse —
96 /// truncated by a proxy, a WAF interstitial appended to JSON — would
97 /// otherwise put the access and refresh tokens into whatever logs the error.
98 /// Status and length carry the same diagnostic value.
99 pub fn json(&self) -> Result<serde_json::Value> {
100 // **The one funnel where a body from outside becomes a `Value`, so the
101 // node guard belongs here.** `listRecords` is bounded by
102 // `parse_list_records`, but every OTHER body this client turns into a
103 // `Value` arrives through this method: the repo writers' responses
104 // (`createRecord`, `putRecord`, `applyWrites`), the PAR response, the
105 // token response and the session refresh. None of them had a bound
106 // before the parse — `read_capped` bounds the WIRE at 8 MB, which is the
107 // input to the amplification, not a limit on it, and 8 MB of the cheapest
108 // node shape measured 824 MB retained on a 512 MB box.
109 //
110 // Guarding here rather than at the four call sites is deliberate: the
111 // listing guard had to be fitted to three clients one round at a time,
112 // twice, which is the whole argument for a single choke point.
113 //
114 // The message carries node counts only — never the body, not even an
115 // excerpt, for the same reason the parse error below does not.
116 refuse_a_response_explosion(&self.body, self.status)?;
117 serde_json::from_slice(&self.body).with_context(|| {
118 format!(
119 "response (status {}, {} bytes) is not valid JSON",
120 self.status,
121 self.body.len()
122 )
123 })
124 }
125
126 pub fn is_success(&self) -> bool {
127 (200..300).contains(&self.status)
128 }
129}
130
131/// Whether this request may be repeated if the server demands a nonce.
132#[derive(Debug, Clone, Copy, PartialEq, Eq)]
133pub enum Retry {
134 /// Safe to repeat: PAR, refresh, resource reads.
135 Allowed,
136 /// **Must not be repeated**, because the body cannot be sent again — a
137 /// consumed stream, say.
138 ///
139 /// This is NOT the right setting for the authorization-code exchange, which
140 /// an earlier revision assumed. A `use_dpop_nonce` challenge means the
141 /// server rejected the request BEFORE processing the grant, so the code was
142 /// never consumed and resending it is safe. Refusing the retry there cost a
143 /// real login against a live PDS: the server rotated its nonce between PAR
144 /// and the token request (nonces last at most five minutes, and user
145 /// approval can take longer), and the flow died with the user already
146 /// approved. The reference declines a retry only when the request body has
147 /// been consumed, which a buffered form body never is.
148 Forbidden,
149}
150
151/// The nonce to retry with, or `None` to stop.
152///
153/// Pure so the policy is testable in isolation. The round trip itself is now
154/// exercised too — `login::tests::a_nonce_challenge_on_the_token_endpoint_is_retried`
155/// drives a real challenge-then-retry against a TLS test server, which the SSRF
156/// guard's loopback refusal used to make impossible.
157fn next_nonce(
158 attempt: usize,
159 retry: Retry,
160 challenge: Option<String>,
161 already_sent: Option<&str>,
162) -> Option<String> {
163 if attempt != 0 || retry == Retry::Forbidden {
164 return None;
165 }
166 let fresh = challenge?;
167 // Retrying with the nonce we already sent is a guaranteed-wasted round trip.
168 if Some(fresh.as_str()) == already_sent {
169 return None;
170 }
171 Some(fresh)
172}
173
174/// The headers for one attempt.
175///
176/// A resource request carries BOTH the proof and the token. The proof's `ath`
177/// binds to an access token the server never sees otherwise, so omitting the
178/// `Authorization` header makes the request unauthenticated — and the resulting
179/// 401 carries no `use_dpop_nonce`, so the retry cannot recover it either. The
180/// scheme is `DPoP`, not `Bearer`: presenting a DPoP-bound token as a bearer
181/// token discards the binding.
182fn request_headers(
183 proof: &str,
184 access_token: Option<&str>,
185) -> Result<Vec<(reqwest::header::HeaderName, reqwest::header::HeaderValue)>> {
186 let mut headers = vec![(
187 reqwest::header::HeaderName::from_static("dpop"),
188 reqwest::header::HeaderValue::from_str(proof)
189 .context("DPoP proof is not a valid header value")?,
190 )];
191 if let Some(token) = access_token {
192 headers.push((
193 reqwest::header::AUTHORIZATION,
194 reqwest::header::HeaderValue::from_str(&format!("DPoP {token}"))
195 .context("access token is not a valid header value")?,
196 ));
197 }
198 Ok(headers)
199}
200
201/// What this request carries, and therefore which method it uses.
202///
203/// The method is DERIVED rather than passed alongside the body. A DPoP proof
204/// binds `htm` to the HTTP method, so a mismatch between the two produces a
205/// proof the server rejects — and keeping them in one value makes that
206/// mismatch unrepresentable.
207pub enum DpopBody<'a> {
208 /// A GET; any parameters are already in the URL.
209 Query,
210 /// `application/x-www-form-urlencoded` — the OAuth endpoints.
211 Form(&'a [(&'a str, &'a str)]),
212 /// `application/json` — the XRPC write endpoints.
213 Json(Vec<u8>),
214}
215
216impl DpopBody<'_> {
217 pub(crate) fn method(&self) -> &'static str {
218 match self {
219 DpopBody::Query => "GET",
220 DpopBody::Form(_) | DpopBody::Json(_) => "POST",
221 }
222 }
223}
224
225/// One DPoP-authenticated request.
226///
227/// Grouped rather than passed positionally so a call site reads as a
228/// description of the request — `Retry::Forbidden` for a request whose body cannot be replayed is
229/// the kind of thing that should be visible at the call, not buried in an
230/// argument list.
231pub struct DpopRequest<'a> {
232 pub endpoint: Endpoint,
233 pub url: &'a str,
234 /// The session's DPoP key — the same one used from PAR onward.
235 pub key: &'a SigningKey,
236 /// Binds the proof via `ath` and is sent as `Authorization: DPoP …`.
237 /// `None` for the authorization-server endpoints.
238 pub access_token: Option<&'a str>,
239 pub body: DpopBody<'a>,
240 pub retry: Retry,
241}
242
243/// Send a request with a DPoP proof, retrying once if the server demands a
244/// nonce and [`Retry`] permits it.
245pub async fn send_with_dpop(
246 client: &Client,
247 pool: &SqlitePool,
248 request: &DpopRequest<'_>,
249) -> Result<PostOutcome> {
250 let DpopRequest {
251 endpoint,
252 url,
253 key,
254 access_token,
255 retry,
256 ref body,
257 } = *request;
258 let method = body.method();
259 let origin = origin_of(url)?;
260 let mut nonce = store::get_nonce(pool, &origin).await?;
261
262 for attempt in 0..2 {
263 let proof = dpop::proof(key, method, url, access_token, nonce.as_deref())?;
264 let headers = request_headers(&proof, access_token)?;
265
266 let response = match body {
267 // NO REDIRECTS on a DPoP GET. A proof's `htu` names the URL it was
268 // minted for, so the request cannot succeed after a cross-origin
269 // redirect anyway — following one buys nothing and costs two things:
270 // the `DPoP-Nonce` is harvested from the FINAL response and stored
271 // under the ORIGINAL origin, letting a redirect poison the nonce for
272 // the real PDS; and the proof itself is re-sent to the redirect
273 // target, since `net`'s cross-host header stripping covers
274 // `Authorization` but not `DPoP`.
275 DpopBody::Query => net::guarded_get_no_redirect(client, url, &headers).await,
276 DpopBody::Form(params) => net::guarded_post_form(client, url, &headers, params).await,
277 DpopBody::Json(bytes) => {
278 net::guarded_post_json(client, url, &headers, bytes.clone()).await
279 }
280 }
281 .with_context(|| format!("{method} {url}"))?;
282
283 let status = response.status().as_u16();
284 let www_authenticate = response
285 .headers()
286 .get(reqwest::header::WWW_AUTHENTICATE)
287 .and_then(|v| v.to_str().ok())
288 .map(str::to_string);
289 // Harvest from EVERY response, success included: the server rotates
290 // nonces, and carrying the newest one forward is what keeps the retry
291 // exceptional rather than routine.
292 let offered = response
293 .headers()
294 .get("DPoP-Nonce")
295 .and_then(|v| v.to_str().ok())
296 .map(str::to_string);
297
298 let body = net::read_capped(response)
299 .await
300 .with_context(|| format!("reading the response from {url}"))?;
301
302 if let Some(offered) = &offered {
303 if Some(offered) != nonce.as_ref() {
304 store::put_nonce(pool, &origin, offered, chrono::Utc::now().timestamp()).await?;
305 }
306 }
307
308 let challenge = dpop::nonce_challenge(
309 endpoint,
310 status,
311 www_authenticate.as_deref(),
312 &body,
313 offered.as_deref(),
314 );
315 if let Some(fresh) = next_nonce(attempt, retry, challenge, nonce.as_deref()) {
316 // Logged because a nonce challenge DOUBLES the round trips for that
317 // request, and nothing else makes that visible: the call is recorded
318 // once by the metrics either way. A server that challenges every
319 // request would halve throughput silently.
320 tracing::debug!(%url, status, "DPoP nonce challenge; retrying once with a fresh nonce");
321 nonce = Some(fresh);
322 continue;
323 }
324 return Ok(PostOutcome { status, body });
325 }
326 unreachable!("the loop returns on its second pass")
327}
328
329#[cfg(test)]
330mod tests {
331 use super::*;
332
333 /// Like every other outbound path, this must fail closed on an internal
334 /// target — asserted on the guard's own error, not merely `is_err()`.
335 #[tokio::test]
336 async fn a_dpop_post_fails_closed_on_an_internal_target() {
337 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
338 store::init_schema(&pool).await.unwrap();
339 let key = SigningKey::generate("k");
340
341 let err = send_with_dpop(
342 &Client::new(),
343 &pool,
344 &DpopRequest {
345 endpoint: Endpoint::AuthorizationServer,
346 url: "http://127.0.0.1/oauth/token",
347 key: &key,
348 access_token: None,
349 body: DpopBody::Form(&[("grant_type", "refresh_token")]),
350 retry: Retry::Allowed,
351 },
352 )
353 .await
354 .expect_err("must refuse a loopback token endpoint");
355 let rendered = format!("{err:#}");
356 assert!(
357 rendered.contains("forbidden (internal) address"),
358 "failed for the wrong reason: {rendered}"
359 );
360 }
361
362 /// **The cap on this funnel is the endpoint's, not the listing's.**
363 ///
364 /// `MAX_LIST_STRUCTURAL_CHARS` (640 000) was sized by the densest page the
365 /// lexicons permit. Nothing here is a listing — token, PAR, refresh, and the
366 /// repo writers' results — so at 210 B per counted character that cap admitted
367 /// about 134 MB per call on paths whose largest legitimate body is 120 kB. A
368 /// review called that out; the constant was borrowed from an unrelated floor.
369 ///
370 /// Both directions, because a cap too tight breaks a real OPML import: the
371 /// largest legitimate body is an `applyWrites` result set, measured at 8 514
372 /// counted characters for 500 results (`max_subs_per_did`'s default).
373 #[test]
374 fn the_response_funnel_has_its_own_cap_sized_for_its_own_traffic() {
375 // A 500-op applyWrites result: the largest thing a real deployment sends
376 // through here, and it must pass.
377 let results: Vec<serde_json::Value> = (0..500)
378 .map(|i| {
379 serde_json::json!({
380 "$type": "com.atproto.repo.applyWrites#createResult",
381 "uri": format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/community.lexicon.rss.subscription/3lab{i:08}"),
382 "cid": "bafyreibaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
383 "validationStatus": "valid",
384 })
385 })
386 .collect();
387 let big_but_legitimate = serde_json::json!({
388 "commit": { "cid": "bafyreibaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", "rev": "3labcdefghijk" },
389 "results": results,
390 })
391 .to_string();
392 let counted = crate::atproto::count_structural_chars(big_but_legitimate.as_bytes());
393 assert!(
394 counted > 8_000,
395 "the probe body counts only {counted}, so it is not the large legitimate \
396 case it is meant to be",
397 );
398 let outcome = PostOutcome {
399 status: 200,
400 body: big_but_legitimate.into_bytes(),
401 };
402 outcome
403 .json()
404 .expect("a 500-op applyWrites result is legitimate traffic on this path");
405
406 // And an explosion is refused — one that the LISTING cap would have
407 // admitted, which is the whole point of a separate constant.
408 let mut body = String::from(r#"{"results":["#);
409 for _ in 0..70_000 {
410 body.push_str("{},");
411 }
412 body.push_str("{}]}");
413 let counted = crate::atproto::count_structural_chars(body.as_bytes());
414 assert!(
415 counted > MAX_RESPONSE_NODES && counted < crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
416 "this probe must sit BETWEEN the two caps to prove they differ: counted \
417 {counted}, endpoint cap {MAX_RESPONSE_NODES}, listing cap {}",
418 crate::atproto::MAX_LIST_STRUCTURAL_CHARS,
419 );
420 let err = PostOutcome {
421 status: 200,
422 body: body.into_bytes(),
423 }
424 .json()
425 .expect_err("a body the listing cap would admit was accepted here");
426 let rendered = format!("{err:#}");
427 assert!(
428 rendered.contains("structural characters"),
429 "failed for the wrong reason: {rendered}"
430 );
431 }
432
433 /// **The body must never reach an error message.** `json()` is called on the
434 /// SUCCESSFUL token response, so a 200 whose body fails to parse — truncated
435 /// by a proxy, a WAF interstitial appended to JSON — would put
436 /// `{"access_token":"eyJ…` into whatever logs the error. That is exactly the
437 /// disclosure the hand-written `Debug` exists to prevent, and echoing an
438 /// "excerpt" reopens it through a different door.
439 #[test]
440 fn a_non_json_body_is_never_echoed_into_the_error() {
441 let outcome = PostOutcome {
442 status: 200,
443 body: br#"{"access_token":"eyJhbGciOiJFUzI1NiJ9.SECRET-TOKEN-VALUE"#.to_vec(),
444 };
445 let rendered = format!("{:#}", outcome.json().unwrap_err());
446 assert!(rendered.contains("200"), "status is the useful diagnostic");
447 assert!(
448 !rendered.contains("SECRET-TOKEN-VALUE") && !rendered.contains("eyJhbGciOiJ"),
449 "the body leaked into the error: {rendered}"
450 );
451 }
452
453 // ── the retry decision ───────────────────────────────────────────────────
454
455 #[test]
456 fn a_nonce_challenge_on_the_first_attempt_is_retried() {
457 assert_eq!(
458 next_nonce(0, Retry::Allowed, Some("fresh".into()), None).as_deref(),
459 Some("fresh")
460 );
461 }
462
463 /// Bounded at one. A server answering every request with `use_dpop_nonce`
464 /// must not make this spin.
465 #[test]
466 fn a_second_attempt_never_retries() {
467 assert!(next_nonce(1, Retry::Allowed, Some("fresh".into()), None).is_none());
468 }
469
470 /// **A request marked `Forbidden` is never retried, whatever the server
471 /// says.**
472 ///
473 /// NOTE: nothing currently passes `Forbidden`, and the authorization-code
474 /// exchange deliberately does NOT — a nonce challenge is rejected before the
475 /// grant is processed, so the code is not consumed and the request is safe
476 /// to resend. An earlier round of review marked that call site `Forbidden`
477 /// on the reasoning below and it killed a real login against a live PDS.
478 ///
479 /// The reference draws the line at whether the request BODY can be re-read,
480 /// not at what the request means; a buffered form always can. This variant
481 /// is kept for a body that cannot be replayed, and the original reasoning is
482 /// preserved here only so it is not rediscovered and re-applied: re-POSTing `grant_type=authorization_code` can burn the
483 /// authorization code, and the login then dies AFTER the user approved,
484 /// presenting as intermittent "login just doesn't work".
485 #[test]
486 fn a_request_that_must_not_repeat_is_never_retried() {
487 assert!(next_nonce(0, Retry::Forbidden, Some("fresh".into()), None).is_none());
488 }
489
490 /// Retrying with the nonce we already sent is a guaranteed-wasted round
491 /// trip; the reference short-circuits it too.
492 #[test]
493 fn an_unchanged_nonce_is_not_worth_retrying() {
494 assert!(next_nonce(0, Retry::Allowed, Some("same".into()), Some("same")).is_none());
495 assert_eq!(
496 next_nonce(0, Retry::Allowed, Some("new".into()), Some("old")).as_deref(),
497 Some("new")
498 );
499 }
500
501 #[test]
502 fn no_challenge_means_no_retry() {
503 assert!(next_nonce(0, Retry::Allowed, None, None).is_none());
504 }
505
506 // ── request shape ────────────────────────────────────────────────────────
507
508 /// The DPoP proof's `htm` must match the method actually sent, so the method
509 /// is derived from the body rather than passed alongside it — there is no
510 /// way for the two to disagree.
511 #[test]
512 fn the_method_follows_the_body_kind() {
513 assert_eq!(DpopBody::Query.method(), "GET");
514 assert_eq!(DpopBody::Form(&[("a", "b")]).method(), "POST");
515 assert_eq!(DpopBody::Json(b"{}".to_vec()).method(), "POST");
516 }
517
518 // ── headers ──────────────────────────────────────────────────────────────
519
520 /// **`ath` without the token is useless.** The proof binds to an access
521 /// token the server never receives, so a resource request arrives
522 /// unauthenticated — and the resulting 401 carries no `use_dpop_nonce`, so
523 /// even the retry cannot recover it.
524 #[test]
525 fn a_resource_request_carries_the_token_as_well_as_the_proof() {
526 let headers = request_headers("the-proof", Some("the-token")).unwrap();
527 let names: Vec<String> = headers.iter().map(|(n, _)| n.to_string()).collect();
528 assert!(names.contains(&"dpop".to_string()));
529 assert!(names.contains(&"authorization".to_string()));
530
531 let auth = headers
532 .iter()
533 .find(|(n, _)| n.as_str() == "authorization")
534 .map(|(_, v)| v.to_str().unwrap().to_string())
535 .unwrap();
536 // The scheme is DPoP, not Bearer: a DPoP-bound token presented as a
537 // bearer token is a downgrade the server should reject.
538 assert_eq!(auth, "DPoP the-token");
539 }
540
541 #[test]
542 fn an_authorization_server_request_carries_only_the_proof() {
543 let headers = request_headers("the-proof", None).unwrap();
544 assert_eq!(headers.len(), 1);
545 assert_eq!(headers[0].0.as_str(), "dpop");
546 }
547
548 /// A token with a newline would otherwise split the header.
549 #[test]
550 fn a_malformed_token_is_rejected_rather_than_injected() {
551 assert!(request_headers("proof", Some("tok\r\nX-Evil: 1")).is_err());
552 assert!(request_headers("pro\nof", None).is_err());
553 }
554}