Skip to main content

turso_backup/
backpressure.rs

1//! Explicit R2 backpressure for the tier-2 WAL-frame drain loop (R574-F2).
2//!
3//! `stream::tail_frames`'s upload loop used to have no retry, no bounded
4//! buffer, no shed — an R2 error bubbled immediately (see the W248 gotcha:
5//! "turso-backup deliberately has NO retry/buffer/shed"). That was correct
6//! for a bare primitive; this module is the explicit decision the doc's §5
7//! calls for: a bounded in-memory spill buffer of frames read-but-not-yet-
8//! confirmed-uploaded, an explicit overflow policy (block / shed / fail),
9//! and retry-with-backoff on 429/503-shaped errors *inside* the drain loop
10//! rather than bubbled raw.
11//!
12//! Implements R574-F2 (relay R574, "WAL→R2 streamer hardening"); the ticket
13//! itself — status, assignee, handoff — is tracked in the W248 working doc,
14//! not duplicated as an in-source annotation here.
15//! @arch:see(.yah/docs/working/W248-wal-streamer-hardening.md)
16//!
17//! ## Design notes (for the reviewer, not the board)
18//!
19//! - `BackpressureConfig{spill_buffer_frames, policy, backoff}` is a
20//!   required field on `StreamConfig`. The drain loop
21//!   (`stream::drain_frames_with_backpressure`) reads WAL frames into a
22//!   `VecDeque` bounded at `spill_buffer_frames`; once full it must resolve
23//!   what it has buffered before reading further — since R761-F2 that is one
24//!   `put` of the whole buffer as a ranged batch object, not one put per
25//!   frame, so the bound sizes the largest object too. A throttling `put` error
26//!   retries with backoff (`put_with_backoff`); `Block` never gives up on a
27//!   throttling error (bubbles instantly on any non-throttling error, same
28//!   as `Fail`/`Shed`); `Fail`/`Shed` give up after `backoff.max_retries`.
29//!   On give-up: `Fail` bubbles the error (nothing persists — matches
30//!   pre-F2 behavior, which is why `Fail` is `BackpressurePolicy`'s
31//!   `#[default]`, so every pre-existing `StreamConfig` call site is
32//!   behavior-preserving); `Shed` drops the whole buffered backlog with a
33//!   loud `BackpressureReport` and returns `Ok(StreamOutcome::Shed{..})`
34//!   with no manifest/watermark write, so the next `tail_frames` call
35//!   re-attempts the same range — no corruption, RPO regresses for the
36//!   shed window instead.
37//! - `StreamOutcome::Streamed`/`Restarted` gained a `backpressure:
38//!   BackpressureReport` field. The new `StreamOutcome::Shed` variant
39//!   covers "every buffered frame dropped, nothing persisted" (`Shed`
40//!   policy only) — kept distinct from `Empty` (engine watermark itself
41//!   has no new frames) so callers can tell "nothing new" from "had new
42//!   frames, all shed".
43//! - `is_throttling_error` is a best-effort string scan over the
44//!   `object_store::Error` display + source chain for
45//!   "429"/"503"/"Too Many Requests"/"Service Unavailable"/"SlowDown"/
46//!   "RequestThrottled" — `object_store`'s own status-aware `RetryError`
47//!   (`client::retry`) is `pub(crate)` to that crate and unreachable from
48//!   here, so there is no public typed classifier to hook; this matches how
49//!   a real S3-compatible backend's `Display` text reads once
50//!   `object_store`'s own internal retry budget is exhausted.
51//! - Test fake `fault_injection::FaultyStore` (`cfg(test)` only) wraps an
52//!   inner `Arc<dyn ObjectStore>` and fails a configured queue of
53//!   `put_opts` calls with a throttling-shaped `Error::Generic` before
54//!   delegating through — mirrors this crate's existing
55//!   mock-the-seam-not-the-engine convention (`stream.rs`'s `MockWal`).
56//!   Needed `async-trait` + `futures-util` as dev-dependencies (both
57//!   already resolved transitively via `object_store`'s own aws/http
58//!   client and its own `throttle` module) since `object_store::ObjectStore`
59//!   is itself `#[async_trait]` and its `delete_stream`/`list` methods are
60//!   typed in terms of `futures_util::stream::BoxStream`.
61//! - The doc left the exact `spill_buffer_frames` bound and backoff numbers
62//!   open. Chose `spill_buffer_frames = 256` and
63//!   `BackoffConfig{initial_delay=200ms, max_delay=10s, multiplier=2.0,
64//!   max_retries=5}` as the `Default` — no measured workload to size
65//!   against yet; both are plain struct fields so a caller overrides them
66//!   per-tenant SLA without a code change.
67
68use std::time::Duration;
69
70/// What [`crate::stream::tail_frames`]'s drain loop does once the bounded
71/// spill buffer is full and the head frame's upload keeps failing with a
72/// throttling-shaped error.
73#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
74pub enum BackpressurePolicy {
75    /// Keep retrying the throttled upload with backoff indefinitely —
76    /// RPO-preserving (no frame is ever dropped), at the cost of a
77    /// `tail_frames` call that can run long — or hang — under a sustained
78    /// R2 outage. Right choice for SLA-tier tenants where a slow drain beats
79    /// losing/delaying frames. Retries never give up *on a throttling
80    /// error*; a non-throttling error (auth, not-found, …) still bubbles
81    /// immediately regardless of policy.
82    Block,
83    /// Once backoff is exhausted for the head frame, drop the entire
84    /// buffered backlog (loudly — see [`BackpressureReport::frames_shed`])
85    /// and return whatever prefix already persisted. Nothing is corrupted:
86    /// no manifest/watermark is written for the dropped tail, so the next
87    /// `tail_frames` call resumes the same range from the last real
88    /// watermark. RPO regresses for the shed window instead of blocking.
89    Shed,
90    /// Once backoff is exhausted, fail the whole `tail_frames` call —
91    /// today's implicit "errors bubble immediately" behavior, made
92    /// explicit. Right choice for non-SLA tiers. Chosen as [`Default`] so
93    /// existing callers see no behavior change until they opt in to
94    /// `Block`/`Shed`.
95    #[default]
96    Fail,
97}
98
99/// Exponential backoff parameters for retrying a throttled (429/503-shaped)
100/// R2 `put`. Applied per batch object inside the drain loop (per frame,
101/// before R761-F2 made the two the same thing).
102#[derive(Debug, Clone, Copy, PartialEq)]
103pub struct BackoffConfig {
104    /// Delay before the first retry.
105    pub initial_delay: Duration,
106    /// Hard cap on any single retry's delay.
107    pub max_delay: Duration,
108    /// Multiplier applied per retry (`initial_delay * multiplier^attempt`,
109    /// capped at `max_delay`).
110    pub multiplier: f64,
111    /// Retries allowed before [`BackpressurePolicy::Fail`]/[`BackpressurePolicy::Shed`]
112    /// give up. Ignored by [`BackpressurePolicy::Block`], which never gives
113    /// up on a throttling error.
114    pub max_retries: u32,
115}
116
117impl Default for BackoffConfig {
118    fn default() -> Self {
119        Self {
120            initial_delay: Duration::from_millis(200),
121            max_delay: Duration::from_secs(10),
122            multiplier: 2.0,
123            max_retries: 5,
124        }
125    }
126}
127
128impl BackoffConfig {
129    /// Delay before retry number `attempt` (0-based), exponential with a
130    /// hard `max_delay` cap.
131    pub fn delay_for(&self, attempt: u32) -> Duration {
132        let factor = self.multiplier.powi(attempt.min(32) as i32);
133        let millis = (self.initial_delay.as_millis() as f64 * factor) as u64;
134        Duration::from_millis(millis).min(self.max_delay)
135    }
136}
137
138/// Bounded spill buffer + overflow policy + retry backoff for the frame
139/// drain loop. See `.yah/docs/working/W248-wal-streamer-hardening.md`
140/// ("Explicit backpressure").
141#[derive(Debug, Clone, Copy, PartialEq)]
142pub struct BackpressureConfig {
143    /// Max frames buffered in memory — read from the WAL but not yet
144    /// confirmed uploaded — before the overflow policy applies. Once full,
145    /// the drain loop must resolve the buffered frames before reading further.
146    ///
147    /// R761-F2: those frames go up as ONE batch object, so this also caps the
148    /// largest object the sink writes at `spill_buffer_frames * (24 +
149    /// page_size)` — ~1 MB at this default and a 4 KB page. Raising it trades
150    /// memory and per-object retry cost for fewer Class A ops on a
151    /// heavy-writing tenant; lowering it to `1` restores one PUT per frame
152    /// (under the batch key shape, not the pre-R761-F2 one).
153    pub spill_buffer_frames: usize,
154    pub policy: BackpressurePolicy,
155    pub backoff: BackoffConfig,
156}
157
158impl Default for BackpressureConfig {
159    fn default() -> Self {
160        Self {
161            spill_buffer_frames: 256,
162            policy: BackpressurePolicy::default(),
163            backoff: BackoffConfig::default(),
164        }
165    }
166}
167
168/// Backpressure activity for a single `tail_frames` call — surfaced on
169/// [`crate::stream::StreamOutcome`] so an orchestrator can alert on the
170/// high-water mark or a nonzero shed count.
171#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
172pub struct BackpressureReport {
173    pub policy: BackpressurePolicy,
174    /// Peak number of frames buffered-but-unconfirmed during this call.
175    pub high_water_frames: usize,
176    /// Frames dropped by a `Shed` overflow (always 0 under `Block`/`Fail`).
177    pub frames_shed: u64,
178    /// Throttling retries consumed across the whole call.
179    pub throttle_retries: u32,
180}
181
182/// True if `err` looks like a 429/503-shaped throttling response rather
183/// than a hard failure (auth, not-found, corrupt request, …).
184///
185/// `object_store`'s own retry classification
186/// (`object_store::client::retry::RetryError::status`) is `pub(crate)` to
187/// that crate and unreachable from here, so this is a best-effort string
188/// classifier over the error's `Display` and `source()` chain. Every
189/// observed backend (including object_store's own internal
190/// retry-exhausted error, whose `Display` embeds the HTTP status text)
191/// surfaces the status this way.
192pub(crate) fn is_throttling_error(err: &object_store::Error) -> bool {
193    let mut msg = err.to_string();
194    let mut cause = std::error::Error::source(err);
195    while let Some(c) = cause {
196        msg.push_str(" | ");
197        msg.push_str(&c.to_string());
198        cause = c.source();
199    }
200    const NEEDLES: [&str; 6] = [
201        "429",
202        "503",
203        "Too Many Requests",
204        "Service Unavailable",
205        "SlowDown",
206        "RequestThrottled",
207    ];
208    NEEDLES.iter().any(|n| msg.contains(n))
209}
210
211/// Upload `payload` to `key`, retrying throttling-shaped errors with
212/// [`BackoffConfig`] backoff. Non-throttling errors bubble immediately (no
213/// retry, regardless of policy). Under [`BackpressurePolicy::Block`],
214/// retries never give up on a throttling error; under `Fail`/`Shed` they
215/// stop after `backoff.max_retries` and the caller (the drain loop) decides
216/// what "stop" means for its policy.
217pub(crate) async fn put_with_backoff(
218    store: &std::sync::Arc<dyn object_store::ObjectStore>,
219    key: &object_store::path::Path,
220    payload: object_store::PutPayload,
221    cfg: &BackpressureConfig,
222    retries_used: &mut u32,
223) -> Result<object_store::PutResult, object_store::Error> {
224    use object_store::ObjectStoreExt;
225    let mut attempt: u32 = 0;
226    loop {
227        match store.put(key, payload.clone()).await {
228            Ok(r) => return Ok(r),
229            Err(e) if is_throttling_error(&e) => {
230                let unbounded = matches!(cfg.policy, BackpressurePolicy::Block);
231                if !unbounded && attempt >= cfg.backoff.max_retries {
232                    return Err(e);
233                }
234                tokio::time::sleep(cfg.backoff.delay_for(attempt)).await;
235                attempt += 1;
236                *retries_used += 1;
237            }
238            Err(e) => return Err(e),
239        }
240    }
241}
242
243/// Test-only fake R2 client: an [`object_store::ObjectStore`] wrapper that
244/// fails a configured queue of `put_opts` calls with a throttling-shaped
245/// error before delegating to the wrapped store.
246#[cfg(test)]
247pub(crate) mod fault_injection {
248    use async_trait::async_trait;
249    use futures_util::stream::BoxStream;
250    use object_store::{
251        path::Path as ObjPath, CopyOptions, Error as OsError, GetOptions, GetResult, ListResult,
252        MultipartUpload, ObjectMeta, ObjectStore, PutMultipartOptions, PutOptions, PutPayload,
253        PutResult, Result as OsResult,
254    };
255    use std::collections::VecDeque;
256    use std::fmt;
257    use std::sync::{Arc, Mutex};
258
259    /// A single queued fault: fail the next `put_opts` call with this
260    /// throttling shape, then move on to the next queued fault (or the real
261    /// store, once the queue is empty).
262    #[derive(Debug, Clone, Copy, PartialEq, Eq)]
263    pub(crate) enum Fault {
264        TooManyRequests,
265        ServiceUnavailable,
266        /// Let this one put through untouched. Only useful positionally — it
267        /// is how a test says "fail the SECOND upload", which is the shape
268        /// every partial-progress scenario needs (R761-F2: a drain that lands
269        /// some batch objects and then gets stuck).
270        Pass,
271    }
272
273    #[derive(Debug)]
274    struct SimulatedThrottle(&'static str);
275    impl fmt::Display for SimulatedThrottle {
276        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
277            f.write_str(self.0)
278        }
279    }
280    impl std::error::Error for SimulatedThrottle {}
281
282    pub(crate) struct FaultyStore {
283        inner: Arc<dyn ObjectStore>,
284        faults: Mutex<VecDeque<Fault>>,
285    }
286
287    impl FaultyStore {
288        pub(crate) fn new(
289            inner: Arc<dyn ObjectStore>,
290            faults: impl IntoIterator<Item = Fault>,
291        ) -> Self {
292            Self {
293                inner,
294                faults: Mutex::new(faults.into_iter().collect()),
295            }
296        }
297    }
298
299    impl fmt::Display for FaultyStore {
300        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
301            write!(f, "FaultyStore({})", self.inner)
302        }
303    }
304    impl fmt::Debug for FaultyStore {
305        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
306            write!(f, "FaultyStore({:?})", self.inner)
307        }
308    }
309
310    #[async_trait]
311    impl ObjectStore for FaultyStore {
312        async fn put_opts(
313            &self,
314            location: &ObjPath,
315            payload: PutPayload,
316            opts: PutOptions,
317        ) -> OsResult<PutResult> {
318            let next = self.faults.lock().unwrap().pop_front();
319            if let Some(fault) = next.filter(|f| *f != Fault::Pass) {
320                let msg: &'static str = match fault {
321                    Fault::TooManyRequests => {
322                        "HTTP status client error (429 Too Many Requests) for url"
323                    }
324                    Fault::ServiceUnavailable => {
325                        "HTTP status server error (503 Service Unavailable) for url"
326                    }
327                    Fault::Pass => unreachable!("filtered out above"),
328                };
329                return Err(OsError::Generic {
330                    store: "faulty-test-store",
331                    source: Box::new(SimulatedThrottle(msg)),
332                });
333            }
334            self.inner.put_opts(location, payload, opts).await
335        }
336
337        async fn put_multipart_opts(
338            &self,
339            location: &ObjPath,
340            opts: PutMultipartOptions,
341        ) -> OsResult<Box<dyn MultipartUpload>> {
342            self.inner.put_multipart_opts(location, opts).await
343        }
344
345        async fn get_opts(&self, location: &ObjPath, options: GetOptions) -> OsResult<GetResult> {
346            self.inner.get_opts(location, options).await
347        }
348
349        fn delete_stream(
350            &self,
351            locations: BoxStream<'static, OsResult<ObjPath>>,
352        ) -> BoxStream<'static, OsResult<ObjPath>> {
353            self.inner.delete_stream(locations)
354        }
355
356        fn list(&self, prefix: Option<&ObjPath>) -> BoxStream<'static, OsResult<ObjectMeta>> {
357            self.inner.list(prefix)
358        }
359
360        async fn list_with_delimiter(&self, prefix: Option<&ObjPath>) -> OsResult<ListResult> {
361            self.inner.list_with_delimiter(prefix).await
362        }
363
364        async fn copy_opts(
365            &self,
366            from: &ObjPath,
367            to: &ObjPath,
368            options: CopyOptions,
369        ) -> OsResult<()> {
370            self.inner.copy_opts(from, to, options).await
371        }
372    }
373}
374
375#[cfg(test)]
376mod tests {
377    use super::*;
378
379    #[test]
380    fn backoff_delay_grows_and_caps() {
381        let cfg = BackoffConfig {
382            initial_delay: Duration::from_millis(100),
383            max_delay: Duration::from_millis(500),
384            multiplier: 2.0,
385            max_retries: 10,
386        };
387        assert_eq!(cfg.delay_for(0), Duration::from_millis(100));
388        assert_eq!(cfg.delay_for(1), Duration::from_millis(200));
389        assert_eq!(cfg.delay_for(2), Duration::from_millis(400));
390        assert_eq!(cfg.delay_for(3), Duration::from_millis(500), "capped at max_delay");
391        assert_eq!(cfg.delay_for(10), Duration::from_millis(500));
392    }
393
394    #[test]
395    fn classifies_429_and_503_as_throttling_but_not_other_errors() {
396        let e429 = object_store::Error::Generic {
397            store: "t",
398            source: Box::new(std::io::Error::other("429 Too Many Requests")),
399        };
400        let e503 = object_store::Error::Generic {
401            store: "t",
402            source: Box::new(std::io::Error::other("503 Service Unavailable")),
403        };
404        let e_not_found = object_store::Error::NotFound {
405            path: "x".into(),
406            source: Box::new(std::io::Error::other("nope")),
407        };
408        assert!(is_throttling_error(&e429));
409        assert!(is_throttling_error(&e503));
410        assert!(!is_throttling_error(&e_not_found));
411    }
412}