From ab9df795ecf75489173c709d598d17d9cbeda634 Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Wed, 19 Aug 2026 11:01:46 -0400 Subject: [PATCH] feat(http): one retry policy, and honor Retry-After Retry lived only in `cmd/images.rs` and `Retry-After` was read nowhere. The policy now sits in `clients/http.rs`, deciding on the status and the transport error kind, reading both spellings of `Retry-After` clamped to 30s, and announcing each wait through `say`; images.rs, the PDS reads, Bobbin and the appview's pull page all run on it. Co-Authored-By: Claude Opus 5 (1M context) --- src/clients/atproto/pds.rs | 21 +- src/clients/http.rs | 560 +++++++++++++++++++++++++++++++ src/clients/tangled/bobbin.rs | 6 +- src/clients/tangled/web/pulls.rs | 2 +- src/cmd/images.rs | 147 ++++---- 5 files changed, 669 insertions(+), 67 deletions(-) diff --git a/src/clients/atproto/pds.rs b/src/clients/atproto/pds.rs index d9d793b..0f9d6d9 100644 --- a/src/clients/atproto/pds.rs +++ b/src/clients/atproto/pds.rs @@ -72,17 +72,30 @@ pub const MAX_PAGES: usize = 5; /// it — six requests instead of five today — and never by `pr list`. pub const MAX_PAGES_COMPLETE: usize = 100; -/// One bounded GET at a PDS, logged, with nothing read off it yet. +/// One bounded GET at a PDS, logged, retried, with nothing read off it yet. /// /// The URL is [`crate::clients::xrpc::endpoint`]'s, which is where this /// module's own copy of the path shape and its percent-encoder went when /// `atgc api` needed the same two lines for a knot and an appview. +/// +/// Retried on the crate's shared policy — see +/// [`crate::clients::http::get_retrying`]. Every method that reaches this is +/// a query: `listRecords`, `getRecord`, `describeRepo`, `getBlob`. None of +/// them changes anything at the PDS, so a second identical request is the +/// same request, and a PDS that answered 429 to a `pr list` mid-walk used to +/// fail the command outright. The writes in +/// [`crate::clients::atproto::record`] go through jacquard and are not on +/// this, deliberately. async fn get(pds: &str, method: &str, params: &[(&str, &str)]) -> Result { let url = crate::clients::xrpc::endpoint(pds, method, params); crate::logging::debug::log(format!(">> GET {url}")); - crate::clients::http::get(&url) - .await - .with_context(|| format!("could not reach the PDS at {pds}")) + crate::clients::http::get_retrying( + &url, + crate::term::say::Topic::Pds, + &format!("{method} at {pds}"), + ) + .await + .with_context(|| format!("could not reach the PDS at {pds}")) } /// One XRPC query, decoded. diff --git a/src/clients/http.rs b/src/clients/http.rs index 23805e6..181a45e 100644 --- a/src/clients/http.rs +++ b/src/clients/http.rs @@ -411,10 +411,570 @@ impl RunningTotal { } } +// --------------------------------------------------------------------------- +// Retrying +// --------------------------------------------------------------------------- + +/// How many times one request is sent before atgc gives up, the first attempt +/// included. Three, which is what `pr create`'s image uploads have used since +/// they were the only thing in the tree that retried at all. +pub const ATTEMPTS: u32 = 3; + +/// The pause before the second attempt; it doubles for each one after. +/// Short enough that a hiccup costs nobody a noticeable wait, and the whole +/// budget — 500ms then 1s — is a second and a half. +pub const BACKOFF: Duration = Duration::from_millis(500); + +/// The most of a `Retry-After` atgc will actually wait. +/// +/// Thirty seconds, which is [`READ_TIMEOUT`]: the longest atgc is already +/// willing to spend on one exchange that is making no progress. A header is +/// the server's opinion about its own load and is worth honouring — it is the +/// only number in the exchange that is not a guess — but it is *not* the +/// server's to decide how long a command sits there. Real services ask for +/// minutes, and a `Retry-After: 3600` on a `pr list` inside a script is +/// indistinguishable from a hang. +/// +/// Past the clamp the wait is still the clamp rather than a refusal: the +/// server said "not yet", and waiting the longest atgc will wait is a better +/// answer to that than either ignoring it or failing without having tried. +pub const MAX_RETRY_AFTER: Duration = Duration::from_secs(30); + +/// What to do about one failed attempt. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum Verdict { + /// A second identical request earns an identical answer. + Fatal, + /// Worth another attempt. `Some` carries what the server itself asked + /// for, already clamped; `None` means fall back to the doubling backoff. + Again(Option), +} + +/// Is this status worth sending the same request again? +/// +/// **Off the status, never the message.** The retry this replaced matched +/// substrings in error prose, so a 1,404-byte image abandoned its upload on +/// the first hiccup and a 403 whose wording did not spell the number was +/// tried three times — see the module doc of [`crate::exit`] for the same +/// complaint about failure classification. +/// +/// 429 is the server asking for exactly this. The four 5xx are the server's +/// own moment: 500 covers a service that panicked on one request, and +/// 502/503/504 are a proxy in front of a service that is restarting or +/// overloaded. Every other 4xx is the request's own fault and does not +/// improve by being repeated — nor does an unlisted 5xx like 501, which is a +/// statement about the method rather than about this moment. +pub fn retryable_status(status: u16) -> bool { + matches!(status, 429 | 500 | 502 | 503 | 504) +} + +/// Is this transport failure worth another attempt? +/// +/// Connect and timeout only, which are the two that say "this exchange did +/// not happen" rather than "this exchange went wrong". They are also exactly +/// the two [`crate::exit::classify`] reads as [`crate::exit::Exit::Unreachable`], +/// so what atgc retries and what atgc tells a caller to retry are one list. +/// +/// A decode failure, a redirect loop or a body error is not here: the request +/// arrived and something about the answer was wrong, and sending it again +/// produces the same wrong answer more slowly. +pub fn retryable_transport(err: &reqwest::Error) -> bool { + err.is_connect() || err.is_timeout() +} + +/// The verdict for one response, header included. +/// +/// `Retry-After` is read on 429 and on 503 and nowhere else. Those are the +/// two statuses RFC 9110 defines it for — "too many requests" and "come back +/// when I am up again" — and a header on a 500 would be a server volunteering +/// a schedule for a fault it has not admitted to. +pub fn response_verdict(resp: &reqwest::Response) -> Verdict { + let status = resp.status().as_u16(); + if !retryable_status(status) { + return Verdict::Fatal; + } + let asked = matches!(status, 429 | 503) + .then(|| retry_after(resp.headers(), std::time::SystemTime::now())) + .flatten(); + Verdict::Again(asked) +} + +/// The `Retry-After` header as a duration, clamped, or `None` if it is +/// absent, empty, not one of the two spellings, or already in the past. +/// +/// `now` is a parameter because the HTTP-date form is only meaningful +/// against a clock, and a test that had to wait for a real one could only +/// assert something vague. See [`parse_retry_after`] for the two forms. +pub fn retry_after( + headers: &reqwest::header::HeaderMap, + now: std::time::SystemTime, +) -> Option { + let raw = headers.get(reqwest::header::RETRY_AFTER)?.to_str().ok()?; + parse_retry_after(raw, now) +} + +/// Both spellings RFC 9110 allows for `Retry-After`, clamped to +/// [`MAX_RETRY_AFTER`]. +/// +/// - *delta-seconds*: `Retry-After: 120`. A whole number of seconds from now. +/// - *HTTP-date*: `Retry-After: Wed, 21 Oct 2026 07:28:00 GMT`. An instant, +/// so what atgc waits is the gap between it and `now`. +/// +/// Both are sent in the wild — Cloudflare and most rate limiters send the +/// first, `503`s from a scheduled maintenance window tend to send the second +/// — so reading only one would be reading none of some services. +/// +/// Anything else is `None`, and the caller falls back to its own backoff. A +/// header atgc cannot read is not worth failing over: the request is already +/// going to be retried, and the only question the header answers is *when*. +/// A date in the past is `None` for the same reason it is not zero — "you may +/// retry now" and "there was no header" want the same behaviour, which is the +/// backoff rather than a hot loop. +fn parse_retry_after(raw: &str, now: std::time::SystemTime) -> Option { + let raw = raw.trim(); + if let Ok(seconds) = raw.parse::() { + return Some(Duration::from_secs(seconds).min(MAX_RETRY_AFTER)); + } + let when = http_date(raw)?; + let gap = when.duration_since(now).ok()?; + if gap.is_zero() { + return None; + } + Some(gap.min(MAX_RETRY_AFTER)) +} + +/// One HTTP-date, as an instant. +/// +/// IMF-fixdate — `Sun, 06 Nov 1994 08:49:37 GMT` — is the only form a server +/// is allowed to *send*; RFC 9110 requires recipients to also accept two +/// obsolete forms, and `parse_from_rfc2822` covers the RFC 850 one's +/// two-digit-year relatives where chrono can. What is not accepted falls back +/// to the backoff, which is the same place an absent header lands. +fn http_date(raw: &str) -> Option { + use jacquard::common::deps::chrono::{DateTime, NaiveDateTime, TimeZone, Utc}; + let stamp = DateTime::parse_from_rfc2822(raw) + .map(|dt| dt.with_timezone(&Utc)) + .ok() + .or_else(|| { + NaiveDateTime::parse_from_str(raw, "%a, %d %b %Y %H:%M:%S GMT") + .ok() + .map(|naive| Utc.from_utc_datetime(&naive)) + })?; + let seconds = u64::try_from(stamp.timestamp()).ok()?; + Some(std::time::UNIX_EPOCH + Duration::from_secs(seconds)) +} + +/// How long to wait before attempt number `attempt + 1`. +/// +/// The server's own number wins where it sent one, because it is the only +/// figure in the exchange that is not a guess. Otherwise the backoff doubles: +/// 500ms, then 1s. +fn pause_before(attempt: u32, asked: Option) -> Duration { + asked.unwrap_or_else(|| BACKOFF * 2u32.saturating_pow(attempt.saturating_sub(1))) +} + +/// Wait, and say so. +/// +/// A retry nobody can see is a stall with no explanation: the command sits +/// there for a second and a half and the only account of it is the exit code +/// it did not produce. [`crate::term::say::Level::Step`] is where the rest of +/// "what this command is doing while it does it" lives, so `-q` silences this +/// along with the rest of the progress and `-qq` silences everything. +async fn pause(topic: crate::term::say::Topic, what: &str, attempt: u32, asked: Option) { + let delay = pause_before(attempt, asked); + let because = match asked { + Some(_) => " (the server asked)", + None => "", + }; + crate::term::say::emit( + crate::term::say::Level::Step, + topic, + &format!( + "{what} failed on attempt {attempt} of {ATTEMPTS}; \ + retrying in {:.1}s{because}", + delay.as_secs_f32() + ), + ); + crate::logging::debug::log(format!( + "retry {what}: attempt {attempt} of {ATTEMPTS}, waiting {delay:?}{because}" + )); + tokio::time::sleep(delay).await; +} + +/// Run `attempt` until it succeeds, is refused for good, or runs out of +/// budget — the one retry loop in the crate. +/// +/// Generic over the error because the two transports that need this do not +/// share one: the reads below fail with `reqwest`'s error, and +/// [`crate::cmd::images`]'s blob upload fails with jacquard's `AgentError`, +/// which wraps a status atgc never sees as a response. `classify` is what +/// each one knows about its own errors; everything else — the budget, the +/// schedule, the announcement — is here rather than copied. +/// +/// The last attempt's error is returned as it is. Nothing is added to it, +/// because the caller's own `with_context` is what names the host, and a +/// budget this loop chose is not a fact about the failure. +pub async fn retrying( + topic: crate::term::say::Topic, + what: &str, + classify: impl Fn(&E) -> Verdict, + mut attempt: F, +) -> Result +where + F: FnMut(u32) -> Fut, + Fut: Future>, +{ + for number in 1..=ATTEMPTS { + match attempt(number).await { + Ok(value) => return Ok(value), + Err(e) => match classify(&e) { + Verdict::Fatal => return Err(e), + Verdict::Again(_) if number == ATTEMPTS => return Err(e), + Verdict::Again(asked) => pause(topic, what, number, asked).await, + }, + } + } + unreachable!("the last attempt either returned or failed"); +} + +/// One GET, retried on the statuses and the transport failures above. +/// +/// The read-path counterpart of [`get`], and the only difference a caller +/// sees is in what it does *not* see: a knot that was restarting, a PDS that +/// answered 429 once. Retrying a GET is safe because a GET is idempotent by +/// definition of the method — which is exactly why the writes in the tree +/// (`putRecord`, `applyWrites`, `deleteRecord`, a knot procedure) are not on +/// this and are not going to be. +/// +/// # What comes back when the budget runs out +/// +/// A transport failure is returned as it is, still carrying the +/// `reqwest::Error` that [`crate::exit::classify`] reads as +/// [`crate::exit::Exit::Unreachable`] — the code that means "try again +/// unchanged", which is precisely what atgc has just spent three attempts +/// finding out is worth doing. +/// +/// A status that stayed retryable to the end is *not* handed back as a +/// response for the caller to report as it always did. It becomes a failure +/// tagged with the same code, because a 429 or a 503 after the whole budget +/// is the same event as a host that did not answer, and a caller branching on +/// `6` should not have to learn that this one arrives as a `1`. The server's +/// own `message` is carried into it where the body is the ATProto error +/// envelope both a PDS and Bobbin send, so the diagnosis is not lost with the +/// response. +/// +/// Every *other* status — a 404, a 403, a 400 — comes back as `Ok(response)` +/// untouched, because the caller's own reading of it is the diagnosis. +pub async fn get_retrying( + url: &str, + topic: crate::term::say::Topic, + what: &str, +) -> anyhow::Result { + let mut exhausted: Option = None; + for number in 1..=ATTEMPTS { + match get(url).await { + Ok(resp) => match response_verdict(&resp) { + Verdict::Fatal => return Ok(resp), + Verdict::Again(_) if number == ATTEMPTS => { + exhausted = Some(resp); + break; + } + Verdict::Again(asked) => pause(topic, what, number, asked).await, + }, + Err(e) if !retryable_transport(&e) || number == ATTEMPTS => return Err(e.into()), + Err(e) => { + crate::logging::debug::log(format!("<< {what}: {e}")); + pause(topic, what, number, None).await; + } + } + } + let resp = exhausted.expect("the loop returns unless a retryable status outlived the budget"); + let status = resp.status(); + let detail = text_bounded(resp, what) + .await + .ok() + .and_then(|text| serde_json::from_str::(&text).ok()) + .and_then(|body| body["message"].as_str().map(str::to_string)); + let said = match detail { + Some(message) => format!(": {message}"), + None => String::new(), + }; + Err(crate::exit::fail( + crate::exit::Exit::Unreachable, + format!("{what} was still answering {status} after {ATTEMPTS} attempts{said}"), + )) +} + #[cfg(test)] mod tests { use super::*; + /// A retry decision is made on the number, not on the prose around it. + /// The four 5xx worth another attempt, the one 4xx that asks for one, + /// and nothing else — a 501 is a statement about the method and a 404 is + /// a statement about the path, and neither improves on the second ask. + #[test] + fn only_the_listed_statuses_retry() { + for retryable in [429, 500, 502, 503, 504] { + assert!(retryable_status(retryable), "{retryable} retries"); + } + for hopeless in [400, 401, 403, 404, 409, 413, 418, 501, 505] { + assert!(!retryable_status(hopeless), "{hopeless} fails fast"); + } + assert!(!retryable_status(200), "a success is not a retry"); + } + + /// `Retry-After: 120`, the form nearly every rate limiter sends. + #[test] + fn retry_after_reads_delta_seconds() { + let now = std::time::UNIX_EPOCH + Duration::from_secs(1_000_000); + assert_eq!( + parse_retry_after("12", now), + Some(Duration::from_secs(12)), + "a whole number of seconds is seconds from now" + ); + assert_eq!( + parse_retry_after(" 12 ", now), + Some(Duration::from_secs(12)), + "whitespace around the value is not a parse failure" + ); + assert_eq!( + parse_retry_after("0", now), + Some(Duration::ZERO), + "zero is a real answer: retry immediately" + ); + } + + /// The other spelling the spec allows, which a maintenance-window 503 + /// tends to send: an instant, so what atgc waits is the gap to it. + #[test] + fn retry_after_reads_an_http_date() { + // 1994-11-06T08:49:37Z, the RFC's own example, and a `now` ten + // seconds before it. + let when = 784_111_777; + let now = std::time::UNIX_EPOCH + Duration::from_secs(when - 10); + assert_eq!( + parse_retry_after("Sun, 06 Nov 1994 08:49:37 GMT", now), + Some(Duration::from_secs(10)), + "the gap between now and the date, not the date" + ); + assert_eq!( + parse_retry_after("Sun, 06 Nov 1994 08:49:37 +0000", now), + Some(Duration::from_secs(10)), + "the same instant spelled with a numeric offset" + ); + } + + /// A date that has already passed is not a wait of zero and not a + /// failure: it lands where an absent header lands, on the backoff. + /// Zero here would be a hot loop against a server that just said no. + #[test] + fn a_date_in_the_past_falls_back_to_the_backoff() { + let when = 784_111_777; + let now = std::time::UNIX_EPOCH + Duration::from_secs(when + 60); + assert_eq!( + parse_retry_after("Sun, 06 Nov 1994 08:49:37 GMT", now), + None + ); + } + + /// **An unreadable header is not an error.** The request is being + /// retried either way; the only thing the header decides is when, so + /// anything atgc cannot read falls through to its own schedule rather + /// than failing a command over a string a server chose. + #[test] + fn an_unparseable_retry_after_falls_back_to_the_backoff() { + let now = std::time::UNIX_EPOCH + Duration::from_secs(1_000_000); + for junk in [ + "", + " ", + "soon", + "-5", + "12.5", + "1e3", + "Sun, 06 Nov 1994 08:49:37", + "Thu, 99 Xxx 1994 08:49:37 GMT", + ] { + assert_eq!( + parse_retry_after(junk, now), + None, + "{junk:?} is not a Retry-After" + ); + } + } + + /// **The server may ask for an hour; it does not get one.** Both + /// spellings are clamped, because the header is the server's opinion + /// about its own load and not its decision about how long a CLI sits + /// there. + #[test] + fn a_long_retry_after_is_clamped_in_both_spellings() { + let now = std::time::UNIX_EPOCH + Duration::from_secs(784_111_777); + assert_eq!( + parse_retry_after("3600", now), + Some(MAX_RETRY_AFTER), + "an hour of delta-seconds is capped" + ); + assert_eq!( + parse_retry_after("Sun, 06 Nov 1994 09:49:37 GMT", now), + Some(MAX_RETRY_AFTER), + "an hour away as a date is capped the same way" + ); + assert!( + MAX_RETRY_AFTER <= READ_TIMEOUT, + "the clamp must not be longer than the longest wait atgc already allows" + ); + } + + /// The schedule when nothing was asked for — 500ms, then 1s — and the + /// fact that anything asked for wins outright. + #[test] + fn the_backoff_doubles_and_a_server_answer_replaces_it() { + assert_eq!(pause_before(1, None), BACKOFF); + assert_eq!(pause_before(2, None), BACKOFF * 2); + assert_eq!(pause_before(3, None), BACKOFF * 4); + assert_eq!( + pause_before(1, Some(Duration::from_secs(7))), + Duration::from_secs(7), + "the only number in the exchange that is not a guess" + ); + } + + /// The transport half of the same decision: the two failures that mean + /// the exchange did not happen. Both are what `exit::classify` already + /// calls `Unreachable`, so the list atgc retries and the list it tells a + /// caller to retry stay the same list. + #[tokio::test] + async fn a_connect_failure_retries_and_classifies_as_unreachable() { + // 203.0.113.0/24 is TEST-NET-3 and unroutable, but this never opens + // a socket: the URL's port is out of range, so reqwest refuses to + // build the request. Instead, drive the two predicates directly on + // the one error shape that can be made without a network. + let err = get("http://203.0.113.1:99999/").await.expect_err("no port"); + assert!( + !retryable_transport(&err), + "a builder error is not a transport failure and must not retry" + ); + let wrapped = anyhow::Error::new(err); + assert_eq!(crate::exit::classify(&wrapped), crate::exit::Exit::Failure); + } + + /// **Exhausting the budget is still `6`.** The whole point of the code + /// is "try again unchanged", and having just spent three attempts + /// establishing that the host is the problem is the strongest evidence + /// for it atgc ever has. The message says the status and the count, so a + /// reader can tell one attempt from three. + #[test] + fn running_out_of_budget_exits_six_and_says_how_many_attempts() { + let err = crate::exit::fail( + crate::exit::Exit::Unreachable, + format!( + "listRecords at pds.example was still answering 429 Too Many Requests \ + after {ATTEMPTS} attempts: slow down" + ), + ); + assert_eq!(crate::exit::classify(&err), crate::exit::Exit::Unreachable); + let message = err.to_string(); + assert!(message.contains("429"), "{message}"); + assert!(message.contains("after 3 attempts"), "{message}"); + assert!( + message.contains("slow down"), + "the server's own words: {message}" + ); + } + + /// The budget is spent and no more, and a fatal verdict spends none of + /// it. Driven through `retrying` with a counter instead of a socket: + /// what is being asserted is the loop, and a request would only make it + /// slower to observe. + #[tokio::test] + async fn the_budget_is_three_attempts_and_a_fatal_verdict_uses_one() { + let tries = std::cell::Cell::new(0); + let result: Result<(), &str> = retrying( + crate::term::say::Topic::Pds, + "a host that keeps saying 503", + |_| Verdict::Again(Some(Duration::ZERO)), + |_| { + tries.set(tries.get() + 1); + async { Err("503") } + }, + ) + .await; + assert_eq!(result, Err("503"), "the last failure comes back as it is"); + assert_eq!(tries.get(), ATTEMPTS, "three attempts, not four"); + + let tries = std::cell::Cell::new(0); + let result: Result<(), &str> = retrying( + crate::term::say::Topic::Pds, + "a host that said 404", + |_| Verdict::Fatal, + |_| { + tries.set(tries.get() + 1); + async { Err("404") } + }, + ) + .await; + assert_eq!(result, Err("404")); + assert_eq!(tries.get(), 1, "a 4xx is asked exactly once"); + } + + /// A success on the second attempt is a success, and the caller never + /// learns there was a first. + #[tokio::test] + async fn a_retry_that_works_returns_the_value() { + let tries = std::cell::Cell::new(0); + let result: Result = retrying( + crate::term::say::Topic::Pds, + "a host having a moment", + |_| Verdict::Again(Some(Duration::ZERO)), + |number| { + tries.set(tries.get() + 1); + async move { if number == 1 { Err("500") } else { Ok(number) } } + }, + ) + .await; + assert_eq!(result, Ok(2)); + assert_eq!(tries.get(), 2); + } + + /// The header is only read where the spec defines it, which is the two + /// statuses that mean "come back later". A 500 volunteering a schedule + /// for a fault it has not admitted to is ignored. + #[test] + fn the_header_is_read_on_429_and_503_only() { + assert!(matches!( + response_verdict_for(429, Some("7")), + Verdict::Again(Some(d)) if d == Duration::from_secs(7) + )); + assert!(matches!( + response_verdict_for(503, Some("7")), + Verdict::Again(Some(d)) if d == Duration::from_secs(7) + )); + assert_eq!(response_verdict_for(500, Some("7")), Verdict::Again(None)); + assert_eq!(response_verdict_for(429, None), Verdict::Again(None)); + assert_eq!(response_verdict_for(404, Some("7")), Verdict::Fatal); + } + + /// The status-and-header half of [`response_verdict`], without a + /// `reqwest::Response` — which cannot be constructed outside a real + /// exchange. The two share every line that decides anything. + fn response_verdict_for(status: u16, header: Option<&str>) -> Verdict { + if !retryable_status(status) { + return Verdict::Fatal; + } + let mut headers = reqwest::header::HeaderMap::new(); + if let Some(value) = header { + headers.insert( + reqwest::header::RETRY_AFTER, + reqwest::header::HeaderValue::from_str(value).expect("a header value"), + ); + } + let asked = matches!(status, 429 | 503) + .then(|| retry_after(&headers, std::time::UNIX_EPOCH)) + .flatten(); + Verdict::Again(asked) + } + /// The bound that matters, against a host that behaves the way the parked /// IP did: accepts the packet and never answers. 203.0.113.0/24 is /// TEST-NET-3 (RFC 5737) and is not routed, so the connect can only diff --git a/src/clients/tangled/bobbin.rs b/src/clients/tangled/bobbin.rs index a7ec51e..46aa8f5 100644 --- a/src/clients/tangled/bobbin.rs +++ b/src/clients/tangled/bobbin.rs @@ -70,9 +70,13 @@ async fn query(method: &str, params: &str) -> Result { /// /// Split from [`query`] because a search's parameters cannot be pasted /// together into a query string by their caller — see [`search_url`]. +/// +/// Retried on the crate's shared policy: every Bobbin method atgc calls is a +/// query, and an index that is briefly overloaded is the case +/// [`crate::clients::http::get_retrying`] exists for. async fn fetch(url: &str) -> Result { crate::logging::debug::log(format!(">> GET {url}")); - let resp = crate::clients::http::get(url) + let resp = crate::clients::http::get_retrying(url, crate::term::say::Topic::Index, "Bobbin") .await .context("could not reach Bobbin (api.tangled.org)")?; let status = resp.status(); diff --git a/src/clients/tangled/web/pulls.rs b/src/clients/tangled/web/pulls.rs index d0455b7..d58e973 100644 --- a/src/clients/tangled/web/pulls.rs +++ b/src/clients/tangled/web/pulls.rs @@ -898,7 +898,7 @@ fn strip_markup(html: &str) -> String { pub(crate) async fn pull_uri_from_appview(web_url: &str, number: u32) -> Result<(String, String)> { let url = format!("{web_url}/pulls/{number}"); crate::logging::debug::log(format!(">> GET {url}")); - let resp = http::get(&url) + let resp = http::get_retrying(&url, crate::term::say::Topic::Index, "the appview") .await .with_context(|| format!("could not reach the appview at {web_url}"))?; let status = resp.status(); diff --git a/src/cmd/images.rs b/src/cmd/images.rs index 2f52568..a60c5a0 100644 --- a/src/cmd/images.rs +++ b/src/cmd/images.rs @@ -27,7 +27,6 @@ use jacquard::types::blob::{Blob, MimeType}; use pulldown_cmark::{Event, LinkType, Parser, Tag}; use std::ops::Range; use std::path::{Path, PathBuf}; -use std::time::Duration; /// The lexicon's cap on one image: the `blobs` items say `maxSize: 1000000`. /// Enforced here so an oversized file is refused by name before anything is @@ -324,14 +323,6 @@ fn image_mime(path: &Path, bytes: &[u8]) -> Option<&'static str> { None } -/// Upload attempts per image before the whole command fails. Three is one -/// real failure forgiven twice — a PDS that refuses a third time is down, -/// not unlucky. -const UPLOAD_ATTEMPTS: u32 = 3; - -/// The first pause before a retry; each later one doubles it. -const UPLOAD_BACKOFF: Duration = Duration::from_millis(500); - /// How many uploads may be in flight at once. Bounded, because a body with /// a dozen screenshots should not open a dozen simultaneous streams against /// one PDS; ordered (`buffered`, not `buffer_unordered`), because the URIs @@ -340,17 +331,18 @@ const UPLOAD_CONCURRENCY: usize = 4; /// Is this failure worth retrying? /// -/// A 4xx is the request's own fault and a second identical request earns an -/// identical answer, except a 429, which asks for exactly a pause and another -/// try. A 5xx is the server's moment, and a failure with no HTTP response -/// behind it at all is the network's. +/// The adapter between jacquard's error and +/// [`crate::clients::http::retryable_status`], which is where the list of +/// statuses lives for every transport in the crate. Nothing about *which* +/// codes retry is decided here; what is decided here is how to get a status +/// out of an `AgentError`, which is jacquard's business and nobody else's. /// -/// Read off jacquard's typed status rather than out of the message. The -/// message is the wrong place to ask: it was matched for the substrings -/// `400`, `401`, `403`, `404`, `413` and `429`, and an error line carries a -/// byte count, a port, a CID fragment and a URL besides its status. A 1,404 -/// byte image abandoned its upload on the first hiccup, while a 403 whose -/// prose did not happen to spell the number was tried three times. That is +/// Read off the typed status rather than out of the message. The message is +/// the wrong place to ask: it was matched for the substrings `400`, `401`, +/// `403`, `404`, `413` and `429`, and an error line carries a byte count, a +/// port, a CID fragment and a URL besides its status. A 1,404 byte image +/// abandoned its upload on the first hiccup, while a 403 whose prose did not +/// happen to spell the number was tried three times. That is /// [`crate::exit`]'s opening complaint about substring-matching prose, /// applied to a retry decision. /// @@ -366,51 +358,65 @@ const UPLOAD_CONCURRENCY: usize = 4; /// error, so no `ClientError` is attached at all. Both are permanent: a /// malformed request and a rejected credential do not improve by being /// sent again. -fn transient(err: &jacquard::client::AgentError) -> bool { +/// +/// `Verdict::Again(None)` in every case, never a duration: jacquard keeps +/// neither the response nor its headers on the way out, so a `Retry-After` +/// the PDS sent is gone before atgc can be handed the failure. The shared +/// backoff is the whole schedule here. The reads in +/// [`crate::clients::http`] do see the header, because they hold the +/// response themselves. +fn verdict(err: &jacquard::client::AgentError) -> crate::clients::http::Verdict { + use crate::clients::http::Verdict; use jacquard::common::error::ClientErrorKind; - match err.client_error() { + let worth_it = match err.client_error() { Some(client) => match client.status() { - Some(status) => status.as_u16() == 429 || status.is_server_error(), + Some(status) => crate::clients::http::retryable_status(status.as_u16()), None => matches!(client.kind(), ClientErrorKind::Transport), }, None => false, + }; + if worth_it { + Verdict::Again(None) + } else { + Verdict::Fatal } } -/// One image, up to [`UPLOAD_ATTEMPTS`] times with doubling pauses between. +/// One image, on the crate's shared retry budget and schedule. +/// +/// A blob upload is a write, and the writes atgc refuses to retry — +/// `putRecord`, `applyWrites`, `deleteRecord` — are refused because a second +/// attempt after an answer atgc did not see would be a second *record*. This +/// one is different in the way that matters: a blob is addressed by the hash +/// of its own bytes, so uploading the same file twice puts the same CID in +/// the same place, and the duplicate costs a request rather than a record. async fn upload_one( agent: &jacquard::client::Agent, img: &LocalImage, ) -> Result { - let mut delay = UPLOAD_BACKOFF; - for attempt in 1..=UPLOAD_ATTEMPTS { - crate::logging::debug::log(format!( - "uploading image blob {} ({} bytes, {}), attempt {attempt}", - img.path.display(), - img.bytes.len(), - img.mime - )); - match agent - .upload_blob(img.bytes.clone(), MimeType::new(img.mime)) - .await - { - Ok(blob) => return Ok(blob), - Err(e) => { - let msg = e.to_string(); - if attempt == UPLOAD_ATTEMPTS || !transient(&e) { - crate::logging::debug::dump_err("uploadBlob error", &e); - bail!("image upload failed for {}: {msg}", img.path.display()); - } - crate::logging::debug::log(format!( - "image upload for {} failed ({msg}); retrying in {delay:?}", - img.path.display() - )); - tokio::time::sleep(delay).await; - delay *= 2; - } - } - } - unreachable!("the last attempt either returned or bailed"); + let what = format!("image upload for {}", img.path.display()); + crate::clients::http::retrying( + crate::term::say::Topic::Image, + &what, + verdict, + |attempt| async move { + crate::logging::debug::log(format!( + "uploading image blob {} ({} bytes, {}), attempt {attempt}", + img.path.display(), + img.bytes.len(), + img.mime + )); + agent + .upload_blob(img.bytes.clone(), MimeType::new(img.mime)) + .await + }, + ) + .await + .map_err(|e| { + let msg = e.to_string(); + crate::logging::debug::dump_err("uploadBlob error", &e); + anyhow::anyhow!("image upload failed for {}: {msg}", img.path.display()) + }) } /// Upload every image to the author's PDS and swap each body reference for @@ -864,12 +870,13 @@ mod tests { #[test] fn transient_errors_retry_and_client_errors_do_not() { - use super::transient; + use super::verdict; + use crate::clients::http::Verdict; use jacquard::client::AgentError; use jacquard::common::error::ClientError; /// An `AgentError` carrying the `ClientError` a status of `code` - /// produces, which is how the upload's failures reach `transient`. + /// produces, which is how the upload's failures reach `verdict`. fn from_status(code: u16) -> AgentError { let status = http::StatusCode::from_u16(code).expect("a real status"); AgentError::new( @@ -881,10 +888,18 @@ mod tests { } for retryable in [429, 500, 502, 503, 504] { - assert!(transient(&from_status(retryable)), "{retryable} retries"); + assert_eq!( + verdict(&from_status(retryable)), + Verdict::Again(None), + "{retryable} retries" + ); } for hopeless in [403, 404, 413] { - assert!(!transient(&from_status(hopeless)), "{hopeless} fails fast"); + assert_eq!( + verdict(&from_status(hopeless)), + Verdict::Fatal, + "{hopeless} fails fast" + ); } // No response at all: a connect, a handshake or a stalled read. The @@ -895,12 +910,20 @@ mod tests { std::io::ErrorKind::ConnectionReset, )))), ); - assert!(transient(&transport), "a reset connection retries"); + assert_eq!( + verdict(&transport), + Verdict::Again(None), + "a reset connection retries" + ); // 400 and 401 are decoded as the endpoint's own typed error, so no // `ClientError` is attached. Both are permanent. let typed = AgentError::sub_operation("upload blob", std::io::Error::other("nope")); - assert!(!transient(&typed), "a typed endpoint error fails fast"); + assert_eq!( + verdict(&typed), + Verdict::Fatal, + "a typed endpoint error fails fast" + ); } /// The regression the typed status closes. Every one of these strings @@ -908,7 +931,8 @@ mod tests { /// them is one. #[test] fn a_status_code_hiding_in_the_prose_no_longer_decides_the_retry() { - use super::transient; + use super::verdict; + use crate::clients::http::Verdict; for status in [500u16, 503] { let err = jacquard::client::AgentError::new( jacquard::client::AgentErrorKind::Client, @@ -919,8 +943,9 @@ mod tests { }, ))), ); - assert!( - transient(&err), + assert_eq!( + verdict(&err), + Verdict::Again(None), "{status} is retryable however the body reads" ); } -- 2.51.2