Skip to main content

turso_backup/
puller.rs

1//! One-puller-per-box fan-out + warm applier (R574-F3).
2//!
3//! `stream::tail_frames`'s R2 *write* side has no fan-out concept — it's a
4//! 1:1 DB↔sink primitive by design (R005-T1). This module is the *read*
5//! side's answer to the doc's §5 fan-in topology: a [`WalPuller`] tails one
6//! tenant's R2 frame stream **once per box** and fans each newly-pulled
7//! frame to every locally-attached warm-standby applier over an in-process
8//! loop, instead of every replica running its own R2 client. R2 read ops are
9//! O(boxes), not O(replicas).
10//!
11//! Depends on R574-T1's §8 RSS curve (this crate's own
12//! `rss_harness` example): the untrimmed measured slope was ~86
13//! KB/replica, which is what justifies keeping hundreds of `CoreWalSeam`
14//! appliers resident per box in the first place — [`WalPullerConfig`]'s FD
15//! budget defaults straight off that number. Trimming (below) is this
16//! ticket's other half of the sizing story: a warm applier never serves
17//! reads, so every cached page is pure standing RSS with no benefit.
18//!
19//! Implements R574-F3 (relay R574, "WAL→R2 streamer hardening"); the ticket
20//! itself — status, assignee, handoff — is tracked in the W248 working doc,
21//! not duplicated as an in-source annotation here (same convention as
22//! `backpressure.rs`).
23//! @arch:see(.yah/docs/working/W248-wal-streamer-hardening.md)
24//!
25//! ## Design notes (for the reviewer, not the board)
26//!
27//! - **Warm applier = [`WarmApplier`]**, a small supertrait of
28//!   `stream::WalInsertSeam` adding one call: `trim_page_cache`. Kept
29//!   separate from `WalInsertSeam` itself (rather than adding a method to
30//!   it) so R574-T4/F2's already-reviewed trait and its existing
31//!   implementors (`CoreWalSeam`, `stream`'s test `MockInsertSeam`) stay
32//!   untouched — this ticket is purely additive over stream.rs. The one
33//!   `CoreWalSeam` change is a single new inherent method,
34//!   [`crate::stream::CoreWalSeam::trim_page_cache_kb`], which drives the
35//!   trim through `PRAGMA cache_size` (turso_core's own public, documented
36//!   SQL surface) rather than reaching into `turso_core::Pager` directly —
37//!   `Pager` is `pub use`d at turso_core's crate root but its
38//!   `CacheResizeResult` return type is not, so the pragma path is the
39//!   clean way in.
40//! - **v1 scope cut — attach is cold-start-only.** `WalPuller::attach`
41//!   refuses once `pull_once()` has succeeded at least once. Fanning
42//!   *future* frames to a late-joining applier is easy (this module does
43//!   exactly that for every already-attached applier); replaying the
44//!   *backlog* the applier missed is not something this module reinvents —
45//!   that's `stream::restore_latest_stream`'s job, already implemented and
46//!   reviewed (R005-F3). The intended flow for adding a replica mid-stream:
47//!   materialize it cold via `restore_latest_stream`, then hand it to a
48//!   **new** `WalPuller` (or restart the box's puller) rather than grafting
49//!   a partially-caught-up seam onto a puller that's already mid-fan-out.
50//!   Flagged in the ticket handoff as the main thing worth reviewer
51//!   pushback: a longer-lived puller that supports hot-attach would need to
52//!   either replay the backlog itself (duplicating restore's logic) or
53//!   retain every pulled frame in memory to serve late joiners (defeats the
54//!   whole point of a bounded footprint) — neither seemed like the right
55//!   default to ship without a concrete caller shaped by real fan-out
56//!   traffic.
57//! - **A WAL restart mid-fan-out is refused, not patched over** — same
58//!   posture as `restore_latest_stream`'s `validate_generation_chain`
59//!   check. If `pull_once()` cannot prove the chain is still the WAL
60//!   generation it already fanned out, every attached applier's WAL may
61//!   already contain frames from the *old* one, which don't compose with the
62//!   new; the fix is a fresh `restore_latest_stream` + a new `WalPuller`, not
63//!   heroics here. R858-B19: the test is
64//!   [`crate::stream::WalGeneration::is_provably_same_as`] — salt included —
65//!   not `checkpoint_seq`, which a writer-process restart resets to the value
66//!   it already had. An unrecorded salt on either side counts as not proven.
67//! - **FD/disk budget** ([`WalPullerConfig`]) is enforced at `attach()`
68//!   time: `max_appliers` bounds concurrently-open applier connections (each
69//!   is several FDs — raise the box's ulimit accordingly, per §5), and
70//!   `max_disk_bytes` bounds the sum of caller-supplied `dest_size_bytes`
71//!   hints (a warm standby is a full local copy, so disk is Σ tenant DB
72//!   sizes — this module has no way to `stat()` a size that's meaningful
73//!   before the applier exists, hence the caller-supplied hint rather than
74//!   inspecting the filesystem itself).
75//! - The doc left every one of these numbers open, same as F2's backpressure
76//!   knobs. Chosen `WalPullerConfig::default()`: `max_appliers = 1000`
77//!   (directly off R574-T1's measured curve — 1000 replicas/box read as
78//!   "strongly supportive of warm-for-everyone" at ~86 KB/replica
79//!   *untrimmed*; trimming should only improve on that), `max_disk_bytes =
80//!   500 GiB` (an unmeasured, round placeholder for a modern box's local
81//!   NVMe — there is no §8-shaped disk measurement yet, unlike RSS),
82//!   `applier_cache_kb = 64` (SQLite's own historical default is 2 MiB;
83//!   64 KiB is a deliberately hard trim per §5's "trimmed cache" ask for a
84//!   connection that never serves reads — plain struct fields, so a caller
85//!   overrides per-tenant SLA without a code change, same pattern as
86//!   `BackpressureConfig`).
87
88use std::collections::HashMap;
89
90use anyhow::{Context, Result};
91use crate::snapshot::BackupTarget;
92use crate::stream::{
93    for_each_frame_in_generation, list_and_parse_generation_manifests, validate_generation_chain,
94    CoreWalSeam, WalGeneration, WalInsertSeam, Watermark,
95};
96
97/// A [`WalInsertSeam`] plus the one extra call a warm standby needs: trim
98/// its page cache hard immediately on attach, since it never serves reads
99/// (§5: "page cache trimmed hard"). See the module doc for why this is a
100/// separate trait rather than a new `WalInsertSeam` method.
101pub trait WarmApplier: WalInsertSeam {
102    /// Resize the underlying page cache toward `target_kb` kilobytes.
103    /// Best-effort in the sense that `turso_core`'s page cache only evicts
104    /// pages that aren't pinned/dirty — harmless here, since an applier
105    /// only ever holds pages it just wrote and is about to hand off to the
106    /// engine, not mid-query state.
107    fn trim_page_cache(&self, target_kb: i64) -> Result<()>;
108}
109
110impl WarmApplier for CoreWalSeam {
111    fn trim_page_cache(&self, target_kb: i64) -> Result<()> {
112        self.trim_page_cache_kb(target_kb)
113    }
114}
115
116/// FD/disk budget + cache-trim knob for a [`WalPuller`]'s attached
117/// appliers. See the module doc for the chosen `Default` numbers and their
118/// (unmeasured) provenance.
119#[derive(Debug, Clone, Copy, PartialEq, Eq)]
120pub struct WalPullerConfig {
121    /// Max concurrently-attached appliers. Each is a live `CoreWalSeam` —
122    /// several open FDs — so this is the box's FD budget, not just a
123    /// counter; raise the box's ulimit to match before raising this.
124    pub max_appliers: usize,
125    /// Max sum of attached appliers' `dest_size_bytes` hints. A warm
126    /// standby is a full local copy of the tenant DB, so this is the box's
127    /// local-disk budget (Σ tenant DB sizes), per §5.
128    pub max_disk_bytes: u64,
129    /// Page cache size (KB) an applier is trimmed to immediately on attach.
130    pub applier_cache_kb: i64,
131}
132
133impl Default for WalPullerConfig {
134    fn default() -> Self {
135        Self {
136            max_appliers: 1_000,
137            max_disk_bytes: 500 * 1024 * 1024 * 1024,
138            applier_cache_kb: 64,
139        }
140    }
141}
142
143/// What a [`WalPuller::pull_once`] call did.
144#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
145pub struct PullReport {
146    /// Frames downloaded from the object store and fanned out this call.
147    /// Each was downloaded exactly once regardless of `appliers_fanned`.
148    pub frames_pulled: u64,
149    /// The source's `checkpoint_seq` as of this call (0 if nothing has ever
150    /// been pulled and there are still no generation manifests to read).
151    pub checkpoint_seq: u32,
152    /// Attached-applier count at the time of this call.
153    pub appliers_fanned: usize,
154}
155
156struct AttachedApplier {
157    seam: Box<dyn WarmApplier>,
158    dest_size_bytes: u64,
159}
160
161/// Tails one tenant's R2 frame stream once per box and fans each newly
162/// pulled frame to every attached [`WarmApplier`]. See the module doc for
163/// the full design, the v1 cold-start-only `attach` scope cut, and the
164/// restart-refusal posture.
165pub struct WalPuller<'a> {
166    target: &'a BackupTarget,
167    page_size: usize,
168    cfg: WalPullerConfig,
169    appliers: HashMap<String, AttachedApplier>,
170    disk_used_bytes: u64,
171    /// The last `(checkpoint_seq, last_frame)` this puller has already
172    /// fanned out to every currently-attached applier. `None` before the
173    /// first successful `pull_once()`.
174    pulled: Option<Watermark>,
175    /// R858-B19: the WAL generation `pulled.last_frame` is a position within.
176    /// Kept beside the watermark rather than inside it because `Watermark` is
177    /// what `WalSeam::wal_state()` returns, and the engine's `wal_state` does
178    /// not report a salt — see [`crate::stream::WalGeneration`].
179    pulled_generation: Option<WalGeneration>,
180}
181
182impl<'a> WalPuller<'a> {
183    pub fn new(target: &'a BackupTarget, page_size: usize, cfg: WalPullerConfig) -> Self {
184        Self {
185            target,
186            page_size,
187            cfg,
188            appliers: HashMap::new(),
189            disk_used_bytes: 0,
190            pulled: None,
191            pulled_generation: None,
192        }
193    }
194
195    pub fn page_size(&self) -> usize {
196        self.page_size
197    }
198
199    pub fn attached_count(&self) -> usize {
200        self.appliers.len()
201    }
202
203    pub fn disk_used_bytes(&self) -> u64 {
204        self.disk_used_bytes
205    }
206
207    /// `true` once `pull_once()` has succeeded at least once — the point
208    /// past which `attach()` refuses (see module doc: v1 does not replay
209    /// backlog to a late-joining applier).
210    pub fn has_pulled(&self) -> bool {
211        self.pulled.is_some()
212    }
213
214    /// Attach a warm applier: begins its `wal_insert` session and trims its
215    /// page cache to `cfg.applier_cache_kb`. From this call on it receives
216    /// every frame this puller fans out.
217    ///
218    /// `dest_size_bytes` is the caller's own estimate of the applier's
219    /// on-disk footprint, counted against `cfg.max_disk_bytes` — see the
220    /// module doc for why this crate takes a hint instead of `stat`-ing a
221    /// file.
222    ///
223    /// Errors if: an applier is already attached under `name`; the FD
224    /// budget (`max_appliers`) or disk budget (`max_disk_bytes`) would be
225    /// exceeded; or this puller has already completed a `pull_once()` (v1
226    /// scope cut — see module doc).
227    pub fn attach(
228        &mut self,
229        name: impl Into<String>,
230        seam: Box<dyn WarmApplier>,
231        dest_size_bytes: u64,
232    ) -> Result<()> {
233        anyhow::ensure!(
234            self.pulled.is_none(),
235            "WalPuller v1 does not support attaching an applier mid-stream — materialize it via \
236             stream::restore_latest_stream and hand it to a fresh WalPuller, or attach before the \
237             first pull_once()"
238        );
239        let name = name.into();
240        anyhow::ensure!(
241            !self.appliers.contains_key(&name),
242            "applier {name:?} is already attached"
243        );
244        anyhow::ensure!(
245            self.appliers.len() < self.cfg.max_appliers,
246            "FD budget exhausted: {} appliers already attached (max_appliers={})",
247            self.appliers.len(),
248            self.cfg.max_appliers,
249        );
250        let projected = self.disk_used_bytes.saturating_add(dest_size_bytes);
251        anyhow::ensure!(
252            projected <= self.cfg.max_disk_bytes,
253            "disk budget exhausted: attaching {name:?} ({dest_size_bytes} bytes) would use \
254             {projected} bytes total, max_disk_bytes={}",
255            self.cfg.max_disk_bytes,
256        );
257
258        seam.wal_insert_begin()
259            .with_context(|| format!("wal_insert_begin for applier {name:?}"))?;
260        seam.trim_page_cache(self.cfg.applier_cache_kb)
261            .with_context(|| format!("trimming page cache for applier {name:?}"))?;
262
263        self.disk_used_bytes = projected;
264        self.appliers
265            .insert(name, AttachedApplier { seam, dest_size_bytes });
266        Ok(())
267    }
268
269    /// Detach an applier without promoting it: closes its `wal_insert`
270    /// session (`force_commit=false`, the same crash-consistent default
271    /// `restore_latest_stream` uses) and returns ownership, freeing its
272    /// share of the FD/disk budget.
273    pub fn detach(&mut self, name: &str) -> Result<Box<dyn WarmApplier>> {
274        let applier = self
275            .appliers
276            .remove(name)
277            .with_context(|| format!("no applier attached under {name:?}"))?;
278        self.disk_used_bytes = self.disk_used_bytes.saturating_sub(applier.dest_size_bytes);
279        applier
280            .seam
281            .wal_insert_end(false)
282            .with_context(|| format!("wal_insert_end for applier {name:?}"))?;
283        Ok(applier.seam)
284    }
285
286    /// Promote a warm applier to serve live traffic. Mechanically identical
287    /// to [`Self::detach`] today (close the session cleanly, hand back the
288    /// seam — it is already caught up as of the last `pull_once()`); kept as
289    /// a separate name so call sites read intent, matching §5's
290    /// "Promotable fast → SLA-tier RTO".
291    pub fn promote(&mut self, name: &str) -> Result<Box<dyn WarmApplier>> {
292        self.detach(name)
293    }
294
295    /// Pull every frame newly written since the last call (or since this
296    /// puller was created) and fan each one out to every attached applier.
297    /// Each frame is downloaded from the object store **once** regardless
298    /// of how many appliers are attached — R2 reads are O(boxes), not
299    /// O(replicas).
300    ///
301    /// Errors if: the chain's `page_size` doesn't match this puller's; or
302    /// the source's `checkpoint_seq` has changed since the last successful
303    /// pull (a WAL restart — see module doc, refused rather than patched
304    /// over). A missing or wrong-length frame object errors the same way
305    /// `restore_latest_stream`'s replay does.
306    pub async fn pull_once(&mut self) -> Result<PullReport> {
307        let manifests = list_and_parse_generation_manifests(self.target).await?;
308        if manifests.is_empty() {
309            return Ok(PullReport {
310                frames_pulled: 0,
311                checkpoint_seq: self.pulled.map(|p| p.checkpoint_seq).unwrap_or(0),
312                appliers_fanned: self.appliers.len(),
313            });
314        }
315        let chain = validate_generation_chain(&manifests)?;
316        anyhow::ensure!(
317            chain.page_size == self.page_size,
318            "generation chain page_size {} does not match this WalPuller's page_size {}",
319            chain.page_size,
320            self.page_size,
321        );
322        // R858-B19: same widening as `tail_frames`. This used to compare
323        // `checkpoint_seq` alone, which a writer-process restart resets to 0 —
324        // so a puller would happily fan frames from a brand-new WAL onto
325        // appliers holding the old one's. The test is now "provably the same
326        // generation", and an unrecorded salt on either side counts as NOT
327        // proven: a puller that cannot tell which WAL it is reading must stop,
328        // because its appliers are live replicas, not a re-runnable restore.
329        if let Some(prior) = self.pulled_generation {
330            anyhow::ensure!(
331                chain.generation.is_provably_same_as(&prior),
332                "WAL restart since the last pull ({} -> {}) — attached appliers \
333                 hold frames from the old WAL generation and can't be trusted to compose with the \
334                 new one; re-materialize them via a fresh restore_latest_stream and start a new \
335                 WalPuller",
336                prior.describe(),
337                chain.generation.describe(),
338            );
339        }
340        let start_frame = self.pulled.map(|p| p.last_frame + 1).unwrap_or(1);
341        if chain.total_frames < start_frame {
342            return Ok(PullReport {
343                frames_pulled: 0,
344                checkpoint_seq: chain.generation.checkpoint_seq,
345                appliers_fanned: self.appliers.len(),
346            });
347        }
348
349        let mut frames_pulled = 0u64;
350        let mut last_frame_no = start_frame - 1;
351        let appliers = &self.appliers;
352        for m in &manifests {
353            if m.last_frame < start_frame {
354                continue; // fully covered by a prior pull_once() call
355            }
356            // R732-F2: the epoch comes from each manifest, not from the chain,
357            // so a pull spanning an ownership transfer still finds both
358            // owners' frames under their own key prefixes. R761-F2: and the
359            // manifest also says whether they are batched, so a pull spanning
360            // the layout change finds both shapes — all of that is
361            // `for_each_frame_in_generation`'s business, not this loop's.
362            frames_pulled += for_each_frame_in_generation(
363                self.target,
364                m,
365                start_frame,
366                |frame_no, bytes| {
367                    for (applier_name, applier) in appliers.iter() {
368                        applier.seam.wal_insert_frame(frame_no, bytes).with_context(|| {
369                            format!("fanning frame {frame_no} to applier {applier_name:?}")
370                        })?;
371                    }
372                    last_frame_no = frame_no;
373                    Ok(())
374                },
375            )
376            .await?;
377        }
378        self.pulled = Some(Watermark {
379            checkpoint_seq: chain.generation.checkpoint_seq,
380            last_frame: last_frame_no,
381        });
382        self.pulled_generation = Some(chain.generation);
383        Ok(PullReport {
384            frames_pulled,
385            checkpoint_seq: chain.generation.checkpoint_seq,
386            appliers_fanned: self.appliers.len(),
387        })
388    }
389}
390
391#[cfg(test)]
392mod tests {
393    use super::*;
394    use crate::backpressure::BackpressureConfig;
395    use crate::snapshot::BackupTarget;
396    use crate::stream::{tail_frames, FrameInfo, StreamConfig, WalSeam, WAL_FRAME_HEADER_SIZE};
397    use object_store::memory::InMemory;
398    use object_store::path::Path as ObjPath;
399    use object_store::{
400        CopyOptions, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta, ObjectStore,
401        ObjectStoreExt, PutMultipartOptions, PutOptions, PutPayload, PutResult, Result as OsResult,
402    };
403    use std::cell::RefCell;
404    use std::fmt;
405    use std::sync::atomic::{AtomicU64, Ordering};
406    use std::sync::Arc;
407
408    // --- fixtures shared with stream.rs's own conventions ------------------
409
410    fn fresh_target() -> BackupTarget {
411        BackupTarget {
412            store: Arc::new(InMemory::new()),
413            prefix: "backups".into(),
414        }
415    }
416
417    fn stream_cfg() -> StreamConfig<'static> {
418        StreamConfig {
419            base_snapshot_key: "backups/snapshots/snapshot-00000000000000000001.db",
420            page_size: 4096,
421            backpressure: BackpressureConfig::default(),
422            rpo_target: None,
423            epoch: 0,
424            owner: None,
425            pointer_generation: 0,
426        }
427    }
428
429    fn puller_cfg(max_appliers: usize, max_disk_bytes: u64) -> WalPullerConfig {
430        WalPullerConfig {
431            max_appliers,
432            max_disk_bytes,
433            applier_cache_kb: 64,
434        }
435    }
436
437    /// In-memory WAL seam mirroring `stream.rs`'s `MockWal` (duplicated
438    /// rather than exposed from `stream`'s private test module — same
439    /// per-module self-contained fixture convention `dedup.rs`/`snapshot.rs`
440    /// already follow).
441    struct MockWal {
442        state: RefCell<MockState>,
443    }
444    struct MockState {
445        checkpoint_seq: u32,
446        frames: Vec<FrameInfo>,
447    }
448    impl MockWal {
449        fn new() -> Self {
450            Self { state: RefCell::new(MockState { checkpoint_seq: 0, frames: Vec::new() }) }
451        }
452        fn append(&self, page_no: u32, db_size: u32) {
453            self.state.borrow_mut().frames.push(FrameInfo { page_no, db_size });
454        }
455        fn restart(&self) {
456            let mut s = self.state.borrow_mut();
457            s.checkpoint_seq += 1;
458            s.frames.clear();
459        }
460    }
461    impl WalSeam for MockWal {
462        fn wal_state(&self) -> Result<Watermark> {
463            let s = self.state.borrow();
464            Ok(Watermark { checkpoint_seq: s.checkpoint_seq, last_frame: s.frames.len() as u64 })
465        }
466        fn wal_get_frame(&self, frame_no: u64, buf: &mut [u8]) -> Result<FrameInfo> {
467            let s = self.state.borrow();
468            let f = s.frames[frame_no as usize - 1];
469            buf[0..4].copy_from_slice(&f.page_no.to_be_bytes());
470            buf[4..8].copy_from_slice(&f.db_size.to_be_bytes());
471            buf[8..WAL_FRAME_HEADER_SIZE].fill(0);
472            buf[WAL_FRAME_HEADER_SIZE..].fill(frame_no as u8);
473            Ok(f)
474        }
475        fn wal_auto_actions_disable(&self) {}
476    }
477
478    /// Captured applier: records every call (including `trim_page_cache`)
479    /// so tests can assert ordering and content without a real turso
480    /// connection. The event log is `Arc`-shared so a test can keep
481    /// inspecting it after handing `Box<dyn WarmApplier>` ownership to a
482    /// `WalPuller` — mirrors `stream.rs`'s `MockInsertSeam` convention,
483    /// extended with `WarmApplier::trim_page_cache`.
484    #[derive(Clone)]
485    struct MockApplier(std::rc::Rc<RefCell<Vec<MockEvent>>>);
486    #[derive(Debug, Clone, PartialEq, Eq)]
487    enum MockEvent {
488        Begin,
489        TrimCache { target_kb: i64 },
490        Frame { frame_no: u64, fill: u8 },
491        End { force_commit: bool },
492    }
493    impl MockApplier {
494        fn new() -> Self {
495            Self(std::rc::Rc::new(RefCell::new(Vec::new())))
496        }
497        fn events(&self) -> Vec<MockEvent> {
498            self.0.borrow().clone()
499        }
500    }
501    impl WalInsertSeam for MockApplier {
502        fn wal_insert_begin(&self) -> Result<()> {
503            self.0.borrow_mut().push(MockEvent::Begin);
504            Ok(())
505        }
506        fn wal_insert_frame(&self, frame_no: u64, frame: &[u8]) -> Result<()> {
507            self.0
508                .borrow_mut()
509                .push(MockEvent::Frame { frame_no, fill: frame[WAL_FRAME_HEADER_SIZE] });
510            Ok(())
511        }
512        fn wal_insert_end(&self, force_commit: bool) -> Result<()> {
513            self.0.borrow_mut().push(MockEvent::End { force_commit });
514            Ok(())
515        }
516    }
517    impl WarmApplier for MockApplier {
518        fn trim_page_cache(&self, target_kb: i64) -> Result<()> {
519            self.0.borrow_mut().push(MockEvent::TrimCache { target_kb });
520            Ok(())
521        }
522    }
523
524    /// A [`WarmApplier`] whose `trim_page_cache` always fails — for
525    /// exercising `attach`'s error path without touching a real connection.
526    struct FailingTrimApplier;
527    impl WalInsertSeam for FailingTrimApplier {
528        fn wal_insert_begin(&self) -> Result<()> {
529            Ok(())
530        }
531        fn wal_insert_frame(&self, _: u64, _: &[u8]) -> Result<()> {
532            Ok(())
533        }
534        fn wal_insert_end(&self, _: bool) -> Result<()> {
535            Ok(())
536        }
537    }
538    impl WarmApplier for FailingTrimApplier {
539        fn trim_page_cache(&self, _: i64) -> Result<()> {
540            anyhow::bail!("simulated trim failure")
541        }
542    }
543
544    /// Wraps an inner store and counts `get_opts` calls — used to prove
545    /// `pull_once` downloads each frame exactly once regardless of how many
546    /// appliers are attached. Mirrors `backpressure.rs`'s `FaultyStore`
547    /// boilerplate (delegate everything, instrument one method).
548    struct CountingStore {
549        inner: Arc<dyn ObjectStore>,
550        gets: AtomicU64,
551    }
552    impl CountingStore {
553        fn new(inner: Arc<dyn ObjectStore>) -> Self {
554            Self { inner, gets: AtomicU64::new(0) }
555        }
556        fn get_count(&self) -> u64 {
557            self.gets.load(Ordering::SeqCst)
558        }
559    }
560    impl fmt::Display for CountingStore {
561        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
562            write!(f, "CountingStore({})", self.inner)
563        }
564    }
565    impl fmt::Debug for CountingStore {
566        fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
567            write!(f, "CountingStore({:?})", self.inner)
568        }
569    }
570    #[async_trait::async_trait]
571    impl ObjectStore for CountingStore {
572        async fn put_opts(
573            &self,
574            location: &ObjPath,
575            payload: PutPayload,
576            opts: PutOptions,
577        ) -> OsResult<PutResult> {
578            self.inner.put_opts(location, payload, opts).await
579        }
580        async fn put_multipart_opts(
581            &self,
582            location: &ObjPath,
583            opts: PutMultipartOptions,
584        ) -> OsResult<Box<dyn MultipartUpload>> {
585            self.inner.put_multipart_opts(location, opts).await
586        }
587        async fn get_opts(&self, location: &ObjPath, options: GetOptions) -> OsResult<GetResult> {
588            self.gets.fetch_add(1, Ordering::SeqCst);
589            self.inner.get_opts(location, options).await
590        }
591        fn delete_stream(
592            &self,
593            locations: futures_util::stream::BoxStream<'static, OsResult<ObjPath>>,
594        ) -> futures_util::stream::BoxStream<'static, OsResult<ObjPath>> {
595            self.inner.delete_stream(locations)
596        }
597        fn list(
598            &self,
599            prefix: Option<&ObjPath>,
600        ) -> futures_util::stream::BoxStream<'static, OsResult<ObjectMeta>> {
601            self.inner.list(prefix)
602        }
603        async fn list_with_delimiter(&self, prefix: Option<&ObjPath>) -> OsResult<ListResult> {
604            self.inner.list_with_delimiter(prefix).await
605        }
606        async fn copy_opts(
607            &self,
608            from: &ObjPath,
609            to: &ObjPath,
610            options: CopyOptions,
611        ) -> OsResult<()> {
612            self.inner.copy_opts(from, to, options).await
613        }
614    }
615
616    // --- attach: FD/disk budget + cold-start-only scope --------------------
617
618    #[test]
619    fn attach_enforces_fd_budget() {
620        let target = fresh_target();
621        let mut puller = WalPuller::new(&target, 4096, puller_cfg(1, u64::MAX));
622        puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
623        let err = puller.attach("b", Box::new(MockApplier::new()), 0).unwrap_err();
624        assert!(format!("{err}").contains("FD budget"), "err was {err}");
625        assert_eq!(puller.attached_count(), 1);
626    }
627
628    #[test]
629    fn attach_enforces_disk_budget() {
630        let target = fresh_target();
631        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, 100));
632        puller.attach("a", Box::new(MockApplier::new()), 60).unwrap();
633        let err = puller.attach("b", Box::new(MockApplier::new()), 60).unwrap_err();
634        assert!(format!("{err}").contains("disk budget"), "err was {err}");
635        assert_eq!(puller.disk_used_bytes(), 60);
636    }
637
638    #[test]
639    fn attach_rejects_duplicate_name() {
640        let target = fresh_target();
641        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
642        puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
643        let err = puller.attach("a", Box::new(MockApplier::new()), 0).unwrap_err();
644        assert!(format!("{err}").contains("already attached"), "err was {err}");
645    }
646
647    /// `attach` opens the wal_insert session and trims the cache, in that
648    /// order, before the applier is otherwise touched.
649    #[test]
650    fn attach_begins_session_and_trims_cache() {
651        let target = fresh_target();
652        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
653        let a = MockApplier::new();
654        puller.attach("a", Box::new(a.clone()), 0).unwrap();
655        assert_eq!(a.events(), vec![MockEvent::Begin, MockEvent::TrimCache { target_kb: 64 }]);
656    }
657
658    #[test]
659    fn attach_propagates_trim_failure_without_registering_applier() {
660        let target = fresh_target();
661        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
662        let err = puller.attach("a", Box::new(FailingTrimApplier), 0).unwrap_err();
663        assert!(format!("{err}").contains("trimming page cache"), "err was {err}");
664        assert_eq!(puller.attached_count(), 0, "a failed attach must not register");
665    }
666
667    #[tokio::test]
668    async fn attach_refuses_after_first_pull() {
669        let seam = MockWal::new();
670        seam.append(1, 1); // commit
671        let target = fresh_target();
672        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
673
674        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
675        assert!(!puller.has_pulled());
676        let report = puller.pull_once().await.unwrap();
677        assert_eq!(report.frames_pulled, 1, "a real frame must actually be fanned");
678        assert!(puller.has_pulled());
679        let err = puller
680            .attach("late", Box::new(MockApplier::new()), 0)
681            .unwrap_err();
682        assert!(format!("{err}").contains("mid-stream"), "err was {err}");
683    }
684
685    /// An empty `pull_once()` (no generation manifests exist yet) leaves
686    /// `has_pulled()` false — nothing was fanned out, so there is no
687    /// backlog a late-joining applier could miss, and attach must still be
688    /// allowed.
689    #[tokio::test]
690    async fn attach_still_allowed_after_an_empty_pull() {
691        let target = fresh_target();
692        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
693        let report = puller.pull_once().await.unwrap();
694        assert_eq!(report.frames_pulled, 0);
695        assert!(!puller.has_pulled());
696        puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
697        assert_eq!(puller.attached_count(), 1);
698    }
699
700    // --- pull_once: single download, fan-out, resumability -----------------
701
702    /// Core fan-in property: N frames staged, M appliers attached — the
703    /// object store sees a fixed number of `get` calls that does not scale
704    /// with M. Since R761-F2 it does not scale with N either: the whole tail
705    /// call is one batch object, so three frames cost one GET, not three.
706    #[tokio::test]
707    async fn pull_once_downloads_each_frame_once_regardless_of_applier_count() {
708        let seam = MockWal::new();
709        seam.append(1, 0);
710        seam.append(2, 0);
711        seam.append(3, 3); // commit
712        let raw_target = fresh_target();
713        let _ = tail_frames(&seam, &raw_target, &stream_cfg()).await.unwrap();
714
715        let counting = Arc::new(CountingStore::new(raw_target.store.clone()));
716        let counted_target = BackupTarget { store: counting.clone(), prefix: "backups".into() };
717
718        let mut puller = WalPuller::new(&counted_target, 4096, puller_cfg(10, u64::MAX));
719        for name in ["a", "b", "c"] {
720            puller.attach(name, Box::new(MockApplier::new()), 0).unwrap();
721        }
722        let report = puller.pull_once().await.unwrap();
723        assert_eq!(report.frames_pulled, 3);
724        assert_eq!(report.appliers_fanned, 3);
725        assert_eq!(
726            counting.get_count(),
727            2,
728            "1 batch object holding all 3 frames + 1 generation manifest, once each"
729        );
730    }
731
732    /// Every attached applier receives the identical frame content and
733    /// ordering; `wal_insert_begin`/cache-trim precede every frame, and
734    /// `pull_once` itself never calls `wal_insert_end` (that's
735    /// `detach`/`promote`'s job).
736    #[tokio::test]
737    async fn pull_once_fans_identical_frames_to_every_applier() {
738        let seam = MockWal::new();
739        seam.append(1, 0);
740        seam.append(2, 2); // commit
741        let target = fresh_target();
742        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
743
744        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
745        let a = MockApplier::new();
746        let b = MockApplier::new();
747        puller.attach("a", Box::new(a.clone()), 0).unwrap();
748        puller.attach("b", Box::new(b.clone()), 0).unwrap();
749        puller.pull_once().await.unwrap();
750
751        let want = vec![
752            MockEvent::Begin,
753            MockEvent::TrimCache { target_kb: 64 },
754            MockEvent::Frame { frame_no: 1, fill: 1 },
755            MockEvent::Frame { frame_no: 2, fill: 2 },
756        ];
757        assert_eq!(a.events(), want);
758        assert_eq!(b.events(), want);
759    }
760
761    /// A second `pull_once()` after new frames land only fans the delta —
762    /// resumable, same posture as `tail_frames`.
763    #[tokio::test]
764    async fn pull_once_is_incremental_across_calls() {
765        let seam = MockWal::new();
766        seam.append(1, 1); // commit
767        let target = fresh_target();
768        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
769
770        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
771        puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
772        let first = puller.pull_once().await.unwrap();
773        assert_eq!(first.frames_pulled, 1);
774
775        // Nothing new yet.
776        let second = puller.pull_once().await.unwrap();
777        assert_eq!(second.frames_pulled, 0);
778
779        seam.append(2, 2); // commit
780        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
781        let third = puller.pull_once().await.unwrap();
782        assert_eq!(third.frames_pulled, 1, "only the new frame, not a re-fan of frame 1");
783    }
784
785    /// No manifests under the prefix yet: a clean no-op, not an error.
786    #[tokio::test]
787    async fn pull_once_with_no_manifests_is_a_noop() {
788        let target = fresh_target();
789        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
790        let report = puller.pull_once().await.unwrap();
791        assert_eq!(report, PullReport { frames_pulled: 0, checkpoint_seq: 0, appliers_fanned: 0 });
792    }
793
794    /// A page_size mismatch between the puller's config and the chain's
795    /// manifests is refused loudly rather than silently misreading frames.
796    #[tokio::test]
797    async fn pull_once_rejects_page_size_mismatch() {
798        let seam = MockWal::new();
799        seam.append(1, 1);
800        let target = fresh_target();
801        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
802
803        let mut puller = WalPuller::new(&target, 8192, puller_cfg(10, u64::MAX));
804        let err = puller.pull_once().await.unwrap_err();
805        assert!(format!("{err}").contains("page_size"), "err was {err}");
806    }
807
808    /// A WAL restart observed between two `pull_once()` calls is refused —
809    /// attached appliers hold frames from the stale sequence.
810    #[tokio::test]
811    async fn pull_once_refuses_across_a_restart() {
812        let seam = MockWal::new();
813        seam.append(1, 1); // commit under seq 0
814        let target = fresh_target();
815        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
816
817        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
818        puller.attach("a", Box::new(MockApplier::new()), 0).unwrap();
819        puller.pull_once().await.unwrap();
820
821        seam.restart();
822        seam.append(1, 1);
823        let _ = tail_frames(&seam, &target, &stream_cfg()).await.unwrap();
824
825        let err = puller.pull_once().await.unwrap_err();
826        assert!(format!("{err}").contains("WAL restart"), "err was {err}");
827    }
828
829    // --- detach / promote ---------------------------------------------------
830
831    #[test]
832    fn detach_frees_disk_budget_and_ends_session() {
833        let target = fresh_target();
834        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, 100));
835        let a = MockApplier::new();
836        puller.attach("a", Box::new(a.clone()), 60).unwrap();
837        assert_eq!(puller.disk_used_bytes(), 60);
838
839        puller.detach("a").unwrap();
840        assert_eq!(puller.disk_used_bytes(), 0);
841        assert_eq!(puller.attached_count(), 0);
842        assert_eq!(a.events().last(), Some(&MockEvent::End { force_commit: false }));
843
844        // A fresh attach can now reuse the freed budget.
845        puller.attach("b", Box::new(MockApplier::new()), 60).unwrap();
846        assert_eq!(puller.disk_used_bytes(), 60);
847    }
848
849    #[test]
850    fn detach_unknown_name_errors() {
851        let target = fresh_target();
852        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
853        assert!(puller.detach("ghost").is_err());
854    }
855
856    #[test]
857    fn promote_ends_session_with_no_forced_commit() {
858        let target = fresh_target();
859        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
860        let a = MockApplier::new();
861        puller.attach("a", Box::new(a.clone()), 0).unwrap();
862        puller.promote("a").unwrap();
863        assert_eq!(puller.attached_count(), 0);
864        assert_eq!(a.events().last(), Some(&MockEvent::End { force_commit: false }));
865    }
866
867    // --- live end-to-end: real turso + CoreWalSeam --------------------------
868
869    struct TempDb(std::path::PathBuf);
870    impl TempDb {
871        fn new(tag: &str) -> Self {
872            TempDb(std::env::temp_dir().join(format!(
873                "turso-backup-puller-{tag}-{}-{}.db",
874                std::process::id(),
875                std::time::SystemTime::now()
876                    .duration_since(std::time::UNIX_EPOCH)
877                    .unwrap()
878                    .as_nanos(),
879            )))
880        }
881        fn path(&self) -> &str {
882            self.0.to_str().unwrap()
883        }
884    }
885    impl Drop for TempDb {
886        fn drop(&mut self) {
887            for sfx in ["", "-wal", "-shm"] {
888                let _ = std::fs::remove_file(format!("{}{sfx}", self.0.display()));
889            }
890        }
891    }
892
893    async fn seed_rows(path: &str, start: i64, count: i64) {
894        let db = turso::Builder::new_local(path).build().await.unwrap();
895        let conn = db.connect().unwrap();
896        conn.execute("CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT)", ())
897            .await
898            .unwrap();
899        conn.execute("BEGIN", ()).await.unwrap();
900        for i in start..start + count {
901            conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (i, format!("v{i}")))
902                .await
903                .unwrap();
904        }
905        conn.execute("COMMIT", ()).await.unwrap();
906    }
907
908    async fn count_rows(path: &str) -> i64 {
909        let db = turso::Builder::new_local(path).build().await.unwrap();
910        let conn = db.connect().unwrap();
911        let mut r = conn.query("SELECT COUNT(*) FROM t", ()).await.unwrap();
912        let row = r.next().await.unwrap().unwrap();
913        row.get::<i64>(0).unwrap()
914    }
915
916    /// End-to-end against a real turso DB and `CoreWalSeam`: seed, snapshot,
917    /// stream WAL frames to the sink, then attach a real warm-applier seam
918    /// (a bare copy of the base snapshot, no WAL yet) to a `WalPuller` and
919    /// `pull_once()` — the applier's row count must match the source after
920    /// the fan-out, and the cache-trim pragma must not error.
921    #[tokio::test]
922    async fn live_pull_once_replays_onto_a_real_applier() {
923        let src = TempDb::new("live-src");
924        let dest = TempDb::new("live-dest");
925        seed_rows(src.path(), 0, 20).await;
926
927        let target = fresh_target();
928        let base_key = match crate::snapshot::snapshot_and_upload(src.path(), &target)
929            .await
930            .unwrap()
931        {
932            crate::snapshot::SnapshotOutcome::Uploaded { key, .. } => key,
933            other => panic!("expected a fresh Uploaded snapshot, got {other:?}"),
934        };
935
936        // Fresh WAL frames past the base snapshot.
937        seed_rows(src.path(), 1000, 5).await;
938        {
939            let seam = crate::stream::CoreWalSeam::open(src.path()).unwrap();
940            let cfg = StreamConfig {
941                base_snapshot_key: &base_key,
942                page_size: 4096,
943                backpressure: BackpressureConfig::default(),
944                rpo_target: None,
945                epoch: 0,
946                owner: None,
947                pointer_generation: 0,
948            };
949            tail_frames(&seam, &target, &cfg).await.unwrap();
950        }
951
952        // Materialize the applier's destination at the base snapshot only —
953        // no WAL replay yet. `WalPuller::pull_once` supplies the frames.
954        let base_bytes = target
955            .store
956            .get(&ObjPath::from(base_key.clone()))
957            .await
958            .unwrap()
959            .bytes()
960            .await
961            .unwrap();
962        std::fs::write(dest.path(), &base_bytes).unwrap();
963
964        let applier = crate::stream::CoreWalSeam::open(dest.path()).unwrap();
965        let mut puller = WalPuller::new(&target, 4096, puller_cfg(10, u64::MAX));
966        puller.attach("dest", Box::new(applier), 0).unwrap();
967        let report = puller.pull_once().await.unwrap();
968        assert!(report.frames_pulled > 0);
969
970        puller.promote("dest").unwrap();
971        assert_eq!(count_rows(dest.path()).await, 25, "20 base + 5 streamed via the puller");
972    }
973}