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}