Skip to main content

turso_backup/
tail.rs

1//! R850-F1 — the **backup half** of hydrate-on-place: keep a running
2//! appliance's declared databases shipped to the object store, under the same
3//! [`claim`] fence [`hydrate`](crate::hydrate) reads.
4//!
5//! [`hydrate`](crate::hydrate) fills an empty volume from the store before a
6//! workload starts. Until something puts bytes *into* the store, that is a
7//! restore path with nothing to restore from — every hydrate returns
8//! `nothing_in_the_store` forever. This module is that something.
9//!
10//! # The constraint that shapes everything here: the application holds the file
11//!
12//! A tenant's local database is a replica this fleet owns outright, so
13//! `yubaba/crates/tenant-streamer` opens one writable [`CoreWalSeam`] per tenant
14//! and holds it for the process lifetime. An appliance's database is the
15//! opposite: the *application* opened it, and turso takes a whole-file exclusive
16//! `fcntl` lock at a writable open. Three consequences, all measured by
17//! `examples/appliance_tail_probe.rs` rather than reasoned about:
18//!
19//! 1. **`CoreWalSeam::open` is refused outright** (probe A1) while the app runs.
20//!    Everything here uses [`CoreWalSeam::open_reader`], which takes no lock.
21//! 2. **A held reader never sees the application's later commits** (A4a); a
22//!    reader opened fresh does (A4b). So each round opens its own seam and drops
23//!    it — and it must drop the old one *first*, because turso's process-global
24//!    `DATABASE_MANAGER` hands a second open of the same path the first handle
25//!    back.
26//! 3. **`VACUUM INTO` and `PRAGMA wal_checkpoint(TRUNCATE)` are unavailable**,
27//!    because both need an ordinary writable connection. That rules out
28//!    [`snapshot::snapshot_and_upload`] and [`dedup::snapshot_dedup`] as written.
29//!    All three tiers here take their image from
30//!    [`stream::raw_consistent_copy_live`] instead, and hand it to the
31//!    image-shaped entry points ([`snapshot::upload_snapshot_image`],
32//!    [`dedup::snapshot_dedup_image`]) that exist for this caller.
33//!
34//! That measurement is also why this is a separate process at all. It was
35//! reported — from the first consumer, and read rather than measured — that a
36//! turso database cannot be tailed from outside the application that holds it,
37//! which would have killed every out-of-process design. It is true of
38//! `CoreWalSeam::open` and false of the crate: probe A5 takes a full
39//! snapshot + tail from a second process and restores it byte-exact.
40//!
41//! # Ownership: acquire once, assert every round
42//!
43//! [`start`] takes the claim; [`round`] re-verifies it. The claim is not a
44//! lease and has no TTL — as [`claim`]'s module doc says, whether a takeover is
45//! *allowed* is placement's call, and what the store guarantees is only that
46//! takeovers are totally ordered and the loser finds out.
47//!
48//! Here the placement call is delegated to the supervisor, and the delegation is
49//! the whole safety argument: **a tail is started only by the node that is
50//! actually running the workload.** If two nodes run it — the split brain the
51//! fence exists for — both tails acquire, the later acquire wins, and the
52//! earlier one's next [`round`] returns [`RoundOutcome::Fenced`]. That is
53//! decisive rather than a coin flip only because the supervisor is obliged to
54//! act on it: `Fenced` means **stop the workload**, not "log and carry on".
55//! Acquiring happens once, at start, so the two nodes converge instead of
56//! ping-ponging.
57//!
58//! [`claim`]: crate::claim
59//! [`CoreWalSeam`]: crate::stream::CoreWalSeam
60//! [`CoreWalSeam::open_reader`]: crate::stream::CoreWalSeam::open_reader
61
62use std::path::{Path, PathBuf};
63use std::sync::Arc;
64use std::time::{Duration, Instant};
65
66use anyhow::{Context, Result};
67use object_store::{ObjectStore, ObjectStoreExt};
68
69use crate::claim::{self, ClaimLost, ClaimOutcome, ClaimRecord};
70use crate::dedup::{self, DedupOutcome};
71use crate::hydrate::{subject_path, Tier};
72use crate::snapshot::{self, BackupTarget, SnapshotOutcome};
73use crate::stream::{
74    self, CoreWalSeam, PreconditionSupport, PreflightStage, SourceFingerprint, StreamConfig,
75    StreamOutcome,
76};
77
78/// How many times [`round`] will re-take a tier-2 base whose WAL generation
79/// moved before the first tail could anchor to it.
80///
81/// Three, and the bound matters more than the number: each attempt is a full
82/// copy of the database, so an application checkpointing faster than we can copy
83/// must give up and report rather than spin. See [`SubjectOutcome::Stream`].
84const MAX_BASE_ATTEMPTS: u32 = 3;
85
86/// Everything [`start`] needs. Mirrors
87/// [`HydrateRequest`](crate::hydrate::HydrateRequest) field for field where the
88/// meaning is the same, because they are read from the same declaration.
89pub struct TailRequest<'a> {
90    /// The object store both the claim and the data live in.
91    pub store: Arc<dyn ObjectStore>,
92    /// Key prefix for this *workload*. The claim sits here; each subject hangs
93    /// one level below — the identical layout [`crate::hydrate`] restores from,
94    /// and not a parallel one, so a tail and a hydrate of the same workload
95    /// cannot disagree about where the bytes are.
96    pub store_prefix: &'a str,
97    /// Host directory the named volume is bound from.
98    pub volume_root: &'a Path,
99    /// Volume-relative database paths, in declaration order.
100    pub subjects: &'a [String],
101    pub tier: Tier,
102    /// Label recorded in the claim. Diagnostic; see [`ClaimRecord::owner`].
103    pub owner: &'a str,
104    /// Page size of the source databases. Required rather than sniffed for the
105    /// reason [`StreamConfig::page_size`] is: a wrong guess silently produces
106    /// frames nothing can replay.
107    pub page_size: usize,
108    /// RPO the caller's scheduler promises, carried into the watermark so a
109    /// missed round is visible as a breach rather than as silence.
110    pub rpo_target: Option<Duration>,
111}
112
113/// Why a tail could not start. Every variant means **do not stream**, and the
114/// caller must decide separately whether the workload may still run.
115#[derive(Debug, Clone, PartialEq, Eq)]
116pub enum TailRefusal {
117    /// The sink does not honour conditional puts, so the fence is not real.
118    /// Same refusal [`crate::hydrate`] makes, for the same reason: pointed at a
119    /// store that ignores `If-Match`, two nodes' claims both succeed and every
120    /// test still passes.
121    SinkNotFenced { stage: PreflightStage },
122    /// Somebody else's acquire landed between our read and our compare-and-swap.
123    /// Not "somebody else owns this" — a sequentially-later acquire always wins;
124    /// this is the genuinely concurrent case, and retrying it is the caller's
125    /// call, not ours.
126    ClaimLost { current: ClaimRecord },
127}
128
129impl TailRefusal {
130    /// One line an operator can act on.
131    pub fn headline(&self) -> String {
132        match self {
133            TailRefusal::SinkNotFenced { stage } => format!(
134                "the sink does not enforce conditional puts ({}), so an ownership claim on it \
135                 cannot fence anybody; refusing to stream rather than ship bytes behind a fence \
136                 that is not there",
137                stage.as_str()
138            ),
139            TailRefusal::ClaimLost { current } => format!(
140                "lost the ownership claim race to epoch {} (owner {})",
141                current.epoch, current.owner
142            ),
143        }
144    }
145}
146
147/// A running tail. Holds the epoch it acquired and one target per subject.
148pub struct TailSession {
149    /// The fencing token every write in this session is stamped with.
150    epoch: u64,
151    /// Whoever we displaced at [`start`], if anybody. Diagnostic, and the only
152    /// record that a takeover happened at all.
153    displaced: Option<ClaimRecord>,
154    workload_target: BackupTarget,
155    subjects: Vec<SubjectState>,
156    tier: Tier,
157    owner: String,
158    page_size: usize,
159    rpo_target: Option<Duration>,
160}
161
162/// Per-subject state carried across rounds.
163struct SubjectState {
164    /// Volume-relative name, as declared.
165    name: String,
166    /// Absolute path on this node.
167    path: PathBuf,
168    target: BackupTarget,
169    /// Tier 2 only: the base snapshot the frames are anchored to. `None` until
170    /// one is published (or discovered) — see [`SubjectOutcome::Stream`].
171    base_snapshot_key: Option<String>,
172}
173
174/// What one subject's turn in a round did.
175#[derive(Debug, Clone, PartialEq, Eq)]
176pub enum SubjectOutcome {
177    /// The declared database does not exist on the volume. Not an error: a
178    /// workload may declare a subject it creates lazily, and the first rounds
179    /// after a first placement legitimately find nothing. Reported rather than
180    /// skipped silently so "why is this subject not backed up" has an answer.
181    Absent,
182    /// The source database was moving under every attempt to copy it, so
183    /// nothing was published this round.
184    ///
185    /// R858-B18 made [`stream::raw_consistent_copy_live`] **refuse** rather than
186    /// return a torn image when a checkpoint lands inside the read window, after
187    /// [`stream::COPY_VALIDATION_ATTEMPTS`] tries. That refusal is a transient
188    /// property of the writer, not a failure of this tail — @Ashguard:polaris
189    /// measured 9 of 12 copies refused under a hammering writer with a 16-page
190    /// autocheckpoint, and 12 of 12 accepted at a realistic rate. Treating it as
191    /// fatal would take a busy appliance's backup down and leave it down.
192    ///
193    /// It is a *reported* non-event rather than a silent one because a subject
194    /// that reports this every round forever is not backed up, and the operator
195    /// needs to see that as a longer `wal_autocheckpoint`, not as silence.
196    SourceTooHot { detail: String },
197    /// Tier 1a.
198    Snapshot(SnapshotOutcome),
199    /// Tier 1b.
200    Dedup(DedupOutcome),
201    /// Tier 2. `base_published` is set on the round that established the base
202    /// this stream is anchored to; `base_attempts` is how many copies that took
203    /// (see [`MAX_BASE_ATTEMPTS`]).
204    Stream {
205        base_snapshot_key: String,
206        base_published: bool,
207        base_attempts: u32,
208        outcome: StreamOutcome,
209    },
210}
211
212/// One subject's turn, with what it cost.
213#[derive(Debug, Clone, PartialEq)]
214pub struct SubjectReport {
215    pub subject: String,
216    pub outcome: SubjectOutcome,
217    pub seconds: f64,
218}
219
220/// What one [`round`] did.
221#[derive(Debug, Clone, PartialEq)]
222pub enum RoundOutcome {
223    /// Every subject got its turn.
224    Backed(Vec<SubjectReport>),
225    /// **This node is no longer the owner and wrote nothing.** The supervisor
226    /// must stop the workload: a fenced node cannot ship a byte, so every write
227    /// its application accepts from here on is a write nobody will ever be able
228    /// to recover.
229    ///
230    /// Reached two ways, deliberately collapsed into one outcome because they
231    /// oblige the caller identically: the pre-round [`claim::assert_holds`] said
232    /// somebody minted a higher epoch, or a tier-2 sink bounced the write on the
233    /// same evidence.
234    Fenced { detail: String },
235}
236
237/// Acquire the claim and prepare a session, or say why not.
238///
239/// The order is [`crate::hydrate`]'s: prove the sink can fence *before* leaning
240/// on the fence, then take the claim. Unlike hydrate there is no readiness check
241/// to run first — a tail is only ever started for a workload the supervisor is
242/// running, so it always intends to take ownership.
243pub async fn start(req: TailRequest<'_>) -> Result<std::result::Result<TailSession, TailRefusal>> {
244    let workload_target = BackupTarget {
245        store: req.store.clone(),
246        prefix: req.store_prefix.to_string(),
247    };
248
249    if let PreconditionSupport::Degraded { stage } =
250        stream::probe_conditional_puts(&workload_target).await?
251    {
252        return Ok(Err(TailRefusal::SinkNotFenced { stage }));
253    }
254
255    let (epoch, displaced) = match claim::acquire(&workload_target, req.owner).await? {
256        ClaimOutcome::Granted { claim, previous } => (claim.epoch, previous),
257        ClaimOutcome::Lost { current } => return Ok(Err(TailRefusal::ClaimLost { current })),
258    };
259
260    let prefix = req.store_prefix.trim_matches('/');
261    let mut subjects = Vec::with_capacity(req.subjects.len());
262    for name in req.subjects {
263        subjects.push(SubjectState {
264            path: subject_path(req.volume_root, name)?,
265            target: BackupTarget {
266                store: req.store.clone(),
267                prefix: if prefix.is_empty() {
268                    name.clone()
269                } else {
270                    format!("{prefix}/{name}")
271                },
272            },
273            name: name.clone(),
274            base_snapshot_key: None,
275        });
276    }
277
278    Ok(Ok(TailSession {
279        epoch,
280        displaced,
281        workload_target,
282        subjects,
283        tier: req.tier,
284        owner: req.owner.to_string(),
285        page_size: req.page_size,
286        rpo_target: req.rpo_target,
287    }))
288}
289
290impl TailSession {
291    /// The fencing token this session streams under.
292    pub fn epoch(&self) -> u64 {
293        self.epoch
294    }
295
296    /// Whoever this session's [`start`] displaced, if anybody.
297    pub fn displaced(&self) -> Option<&ClaimRecord> {
298        self.displaced.as_ref()
299    }
300
301    pub fn subjects(&self) -> impl Iterator<Item = &str> {
302        self.subjects.iter().map(|s| s.name.as_str())
303    }
304}
305
306/// One pass over every subject.
307///
308/// The claim is re-verified first, and that check is not redundant with the
309/// sink-side fence: only tier 2 stamps its writes with an epoch the sink can
310/// bounce. A tier-1a or tier-1b tail writes snapshots and manifests
311/// unconditionally, so without this read a fenced node would keep overwriting
312/// the real owner's backups with its own stale ones — the exact corruption the
313/// fence exists to prevent, arriving through the door tier 2 happens to have
314/// closed and the other two do not.
315///
316/// A subject that fails is propagated, not swallowed: unlike
317/// `tenant-streamer`'s per-tenant isolation (one sick tenant must not stop a box
318/// full of healthy ones), every subject here belongs to the same workload, and a
319/// workload whose accounts database is backed up but whose sessions database is
320/// not is not backed up.
321pub async fn round(session: &mut TailSession) -> Result<RoundOutcome> {
322    if let Err(lost) = claim::assert_holds(&session.workload_target, session.epoch).await? {
323        return Ok(RoundOutcome::Fenced {
324            detail: fenced_detail(&lost),
325        });
326    }
327
328    let tier = session.tier;
329    let page_size = session.page_size;
330    let owner = session.owner.clone();
331    let epoch = session.epoch;
332    let rpo_target = session.rpo_target;
333
334    // Cloned rather than borrowed: the loop below takes `session.subjects`
335    // mutably, and `rebase` needs the WORKLOAD prefix — the claim lives there,
336    // one level above each subject's. Getting that wrong is not a compile error
337    // and not a wrong number; it is a rebase that asserts a claim nobody ever
338    // wrote and refuses forever, which is how the test that found it failed.
339    let workload_target = BackupTarget {
340        store: session.workload_target.store.clone(),
341        prefix: session.workload_target.prefix.clone(),
342    };
343
344    let mut reports = Vec::with_capacity(session.subjects.len());
345    for subject in &mut session.subjects {
346        let started = Instant::now();
347        if !subject.path.exists() {
348            reports.push(SubjectReport {
349                subject: subject.name.clone(),
350                outcome: SubjectOutcome::Absent,
351                seconds: started.elapsed().as_secs_f64(),
352            });
353            continue;
354        }
355        let outcome = back_up_subject(
356            subject,
357            &workload_target,
358            tier,
359            page_size,
360            &owner,
361            epoch,
362            rpo_target,
363        )
364            .await
365            .with_context(|| format!("backing up subject {}", subject.name))?;
366        if let SubjectOutcome::Stream {
367            outcome: StreamOutcome::Fenced { current_epoch, our_epoch, .. },
368            ..
369        } = &outcome
370        {
371            return Ok(RoundOutcome::Fenced {
372                detail: format!(
373                    "the sink bounced subject {}: it is stamped with epoch {current_epoch} and we \
374                     hold {our_epoch}",
375                    subject.name
376                ),
377            });
378        }
379        reports.push(SubjectReport {
380            subject: subject.name.clone(),
381            outcome,
382            seconds: started.elapsed().as_secs_f64(),
383        });
384    }
385    Ok(RoundOutcome::Backed(reports))
386}
387
388fn fenced_detail(lost: &ClaimLost) -> String {
389    match lost {
390        ClaimLost::Superseded { .. } => format!("the ownership claim moved: {lost}"),
391        // Not a state this crate produces. Treated as fenced anyway rather than
392        // as "carry on": a claim we cannot verify is a claim we do not hold,
393        // which is the same rule `claim::assert_holds`'s own doc states for an
394        // unreachable store.
395        ClaimLost::Vanished { .. } => format!("{lost}; treating that as fenced"),
396    }
397}
398
399async fn back_up_subject(
400    subject: &mut SubjectState,
401    workload_target: &BackupTarget,
402    tier: Tier,
403    page_size: usize,
404    owner: &str,
405    epoch: u64,
406    rpo_target: Option<Duration>,
407) -> Result<SubjectOutcome> {
408    // Owned, so the `Tier::Stream` arm can take `subject` mutably below.
409    let path = subject
410        .path
411        .to_str()
412        .with_context(|| format!("subject path {} is not UTF-8", subject.path.display()))?
413        .to_string();
414    let path = path.as_str();
415
416    match tier {
417        Tier::Snapshot => {
418            let image = match live_image(path, page_size).await? {
419                Ok(image) => image,
420                Err(too_hot) => return Ok(too_hot),
421            };
422            Ok(SubjectOutcome::Snapshot(
423                snapshot::upload_snapshot_image(&subject.target, &image).await?,
424            ))
425        }
426        Tier::Dedup => {
427            let image = match live_image(path, page_size).await? {
428                Ok(image) => image,
429                Err(too_hot) => return Ok(too_hot),
430            };
431            Ok(SubjectOutcome::Dedup(
432                dedup::snapshot_dedup_image(&image, &subject.target).await?,
433            ))
434        }
435        Tier::Stream => {
436            stream_subject(subject, workload_target, path, page_size, owner, epoch, rpo_target)
437                .await
438        }
439    }
440}
441
442/// [`stream::raw_consistent_copy_live`] with its one *expected* failure split
443/// out of the error channel.
444///
445/// `Ok(Ok(image))` is a validated point-in-time copy; `Ok(Err(SourceTooHot))` is
446/// "the writer never held still, try again next round"; `Err` is a real failure.
447///
448/// The discrimination is a string match on the refusal's own sentence, which is
449/// ugly and is the honest option available: that path returns `anyhow::Error`
450/// with no typed variant, and adding one means editing `stream.rs`, which
451/// @Ashguard:polaris holds live under R858-B18. The phrase is pinned by
452/// [`tests::a_source_that_never_holds_still_is_too_hot_not_a_failure`], which
453/// provokes the real error rather than asserting on a copy of the string — so a
454/// reword on their side fails a test here instead of silently turning a busy
455/// database into a dead backup.
456async fn live_image(
457    path: &str,
458    page_size: usize,
459) -> Result<std::result::Result<Vec<u8>, SubjectOutcome>> {
460    match stream::raw_consistent_copy_live(path, page_size).await {
461        Ok(image) => Ok(Ok(image)),
462        Err(e) if is_source_too_hot(&e) => Ok(Err(SubjectOutcome::SourceTooHot {
463            detail: format!("{e:#}"),
464        })),
465        Err(e) => Err(e),
466    }
467}
468
469/// Whether an error is R858-B18's "the source moved under every attempt"
470/// refusal. See [`live_image`].
471fn is_source_too_hot(e: &anyhow::Error) -> bool {
472    format!("{e:#}").contains("the source moved under every one of")
473}
474
475/// Tier 2: anchor to a base, then tail frames onto it.
476///
477/// ## Why the base and the first tail are one operation
478///
479/// A tier-2 restore is *base image + replayed frames*. The base this module
480/// publishes is a [`stream::raw_consistent_copy_live`] image, which already has
481/// every frame committed up to the copy point folded in; the first tail then
482/// uploads frames `1..=max` of the live WAL, and replaying those over the base
483/// is idempotent because it rewrites the same pages with the same content.
484///
485/// What breaks that is a **checkpoint landing between the copy and the tail**.
486/// The fold moves those frames into the main file and resets the WAL, so the
487/// tail uploads frames `1..=max` of a *new* generation — and anything the
488/// application committed after our copy point but before the fold is in neither
489/// the base nor the uploaded frames. Nothing downstream can detect this: the
490/// chain validates, the restore succeeds, and the result is a database missing a
491/// window of writes.
492///
493/// So the generation is sampled around the whole base-plus-first-tail and the
494/// pair is redone if it moved, up to [`MAX_BASE_ATTEMPTS`]. This costs nothing
495/// on every subsequent round, which is the overwhelming majority: once a base
496/// exists, a fold is just [`StreamOutcome::Restarted`], which re-uploads from
497/// frame 1 against a base that is already correct.
498async fn stream_subject(
499    subject: &mut SubjectState,
500    workload_target: &BackupTarget,
501    path: &str,
502    page_size: usize,
503    owner: &str,
504    epoch: u64,
505    rpo_target: Option<Duration>,
506) -> Result<SubjectOutcome> {
507    // A base may already exist from an earlier incarnation of this workload —
508    // including the one whose bytes `hydrate` just restored onto this volume.
509    // Re-publishing then would be a second full copy for nothing.
510    if subject.base_snapshot_key.is_none() {
511        subject.base_snapshot_key = snapshot::latest_snapshot_key(&subject.target).await?;
512    }
513
514    if let Some(base) = subject.base_snapshot_key.clone() {
515        let outcome = tail_once(subject, path, &base, page_size, owner, epoch, rpo_target).await?;
516        // R858-B19: a WAL recreate is not an ordinary round. `tail_frames` has
517        // re-uploaded frames 1..N under a NEW generation, and the prefix now
518        // holds manifests from two — which `validate_generation_chain` refuses
519        // ("WAL restart between generations, restore needs a fresh tier-1a
520        // snapshot"). Folding this into the Streamed arm and logging "streamed"
521        // leaves a prefix that looks healthy and cannot be restored, which is
522        // the worst of the available outcomes.
523        if matches!(outcome, StreamOutcome::Restarted { .. }) {
524            return rebase(
525                subject,
526                workload_target,
527                path,
528                page_size,
529                owner,
530                epoch,
531                rpo_target,
532            )
533            .await;
534        }
535        return Ok(SubjectOutcome::Stream {
536            base_snapshot_key: base,
537            base_published: false,
538            base_attempts: 0,
539            outcome,
540        });
541    }
542
543    let mut last_generation = None;
544    for attempt in 1..=MAX_BASE_ATTEMPTS {
545        let before = SourceFingerprint::read(path)?;
546        let image = match live_image(path, page_size).await? {
547            Ok(image) => image,
548            Err(too_hot) => return Ok(too_hot),
549        };
550        let base = snapshot::upload_base_snapshot(&subject.target, &image).await?;
551        let outcome = tail_once(subject, path, &base, page_size, owner, epoch, rpo_target).await?;
552        let after = SourceFingerprint::read(path)?;
553        if before.stable_across(&after) {
554            subject.base_snapshot_key = Some(base.clone());
555            return Ok(SubjectOutcome::Stream {
556                base_snapshot_key: base,
557                base_published: true,
558                base_attempts: attempt,
559                outcome,
560            });
561        }
562        last_generation = Some((before, after));
563    }
564
565    let (before, after) = last_generation.expect("MAX_BASE_ATTEMPTS is non-zero");
566    anyhow::bail!(
567        "could not anchor a tier-2 base for {path}: the source WAL was folded during each of \
568         {MAX_BASE_ATTEMPTS} attempts (last: {before:?} -> {after:?}). The base and the frames \
569         would describe different points in time, so nothing was anchored; the next round retries. \
570         An application checkpointing faster than its database can be copied needs a longer \
571         wal_autocheckpoint, not a longer retry loop."
572    );
573}
574
575/// Re-anchor a subject whose WAL was recreated: publish a fresh tier-1a base,
576/// drop the chain that can no longer validate, and tail onto the new base.
577///
578/// R858-B19 made a generation's identity `(checkpoint_seq, salt)` and taught
579/// restore to REFUSE a chain that spans a recreate rather than splice it. That
580/// turned a silent wrong restore into a loud one, and left somebody owing the
581/// repair. This is it.
582///
583/// ## The order is chosen for what a crash in the middle leaves behind
584///
585/// 1. **Publish the new base.** A crash here leaves two bases and the old
586///    manifests, so a restore *refuses* — loud, and the next round rebases
587///    again and clears it.
588/// 2. **Delete every generation manifest.** A crash here leaves the new base
589///    and no chain, which `hydrate` restores as a bare tier-1a snapshot at the
590///    point the base was taken. Correct, just behind.
591/// 3. **Delete the watermark sidecar**, so the next tail starts a chain at
592///    frame 1 — `validate_generation_chain` requires a chain to start there,
593///    so resuming mid-range would produce another unrestorable prefix.
594/// 4. **Tail onto the new base.**
595///
596/// Deleting them the other way round — chain first — would leave the *old* base
597/// as the newest snapshot with no manifests, and a restore would silently
598/// succeed at a stale point in time. Loud beats silent, so the base goes first.
599///
600/// ## What this leaves behind, deliberately
601///
602/// Two things, not one — R850-T3 corrected this doc while building the sweep
603/// that reclaims them:
604///
605/// - The old generation's frame objects, orphaned under
606///   `frames/{old_checkpoint_seq}/`. They are keyed by sequence so they collide
607///   with nothing and are invisible to restore once their manifests are gone.
608/// - **The old base snapshot**, which step 1 supersedes and nothing deletes:
609///   [`snapshot::upload_base_snapshot`] is the non-deduplicating one-shot
610///   variant, so every rebase leaves a complete extra copy of the database
611///   behind. For any database past a few megabytes this is the *larger* of the
612///   two leaks, and this doc did not mention it until R850-T3.
613///
614/// Reclaiming either is GC's job, not a rebase's — deleting data as part of a
615/// recovery path is how a recovery path becomes the outage. The sweep is
616/// [`crate::stream::gc_stream`]; it is explicitly invoked, dry-run by default,
617/// and deletes only what no restore can reach.
618///
619/// ## The fence during the gap
620///
621/// Step 3 removes the sidecar that carries the epoch `tail_frames` bounces a
622/// stale writer on, so between it and step 4 the sink is briefly unfenced. The
623/// claim is re-asserted immediately before step 1 to narrow that window, and the
624/// claim — not the watermark — is this design's real fence: `round` verifies it
625/// every pass, and a node that lost it stops its workload.
626async fn rebase(
627    subject: &mut SubjectState,
628    workload_target: &BackupTarget,
629    path: &str,
630    page_size: usize,
631    owner: &str,
632    epoch: u64,
633    rpo_target: Option<Duration>,
634) -> Result<SubjectOutcome> {
635    if let Err(lost) = claim::assert_holds(workload_target, epoch).await? {
636        anyhow::bail!(
637            "refusing to re-anchor {path} after a WAL restart: {lost}. Rebasing rewrites the \
638             chain, and doing that without the claim would destroy the real owner's"
639        );
640    }
641
642    let image = match live_image(path, page_size).await? {
643        Ok(image) => image,
644        Err(too_hot) => return Ok(too_hot),
645    };
646    let base = snapshot::upload_base_snapshot(&subject.target, &image).await?;
647    delete_generation_manifests(&subject.target).await?;
648    subject
649        .target
650        .store
651        .delete(&subject.target.watermark_key())
652        .await
653        .or_else(ignore_absent)
654        .with_context(|| format!("clearing the stream watermark for {path}"))?;
655
656    let outcome = tail_once(subject, path, &base, page_size, owner, epoch, rpo_target).await?;
657    subject.base_snapshot_key = Some(base.clone());
658    Ok(SubjectOutcome::Stream {
659        base_snapshot_key: base,
660        base_published: true,
661        base_attempts: 1,
662        outcome,
663    })
664}
665
666async fn delete_generation_manifests(target: &BackupTarget) -> Result<()> {
667    let prefix = target.prefix.trim_matches('/');
668    let dir = object_store::path::Path::from(if prefix.is_empty() {
669        "generations".to_string()
670    } else {
671        format!("{prefix}/generations")
672    });
673    let listing = target
674        .store
675        .list_with_delimiter(Some(&dir))
676        .await
677        .with_context(|| format!("listing generation manifests under {dir}"))?;
678    for object in listing.objects {
679        target
680            .store
681            .delete(&object.location)
682            .await
683            .or_else(ignore_absent)
684            .with_context(|| format!("deleting stale generation manifest {}", object.location))?;
685    }
686    Ok(())
687}
688
689/// A delete of something already gone is the outcome we wanted, not a failure —
690/// two rebases racing, or a retry after a partial one, both land here.
691fn ignore_absent(e: object_store::Error) -> std::result::Result<(), object_store::Error> {
692    match e {
693        object_store::Error::NotFound { .. } => Ok(()),
694        other => Err(other),
695    }
696}
697
698/// One `tail_frames` through a seam opened for this call only.
699///
700/// The seam's lifetime is exactly this function, and that is load-bearing rather
701/// than tidy — see this module's doc, point 2: a held reader never observes the
702/// application's later commits, and reopening while still holding one returns
703/// the same handle from turso's process-global registry.
704async fn tail_once(
705    subject: &SubjectState,
706    path: &str,
707    base_snapshot_key: &str,
708    page_size: usize,
709    owner: &str,
710    epoch: u64,
711    rpo_target: Option<Duration>,
712) -> Result<StreamOutcome> {
713    let cfg = StreamConfig {
714        base_snapshot_key,
715        page_size,
716        backpressure: Default::default(),
717        rpo_target,
718        epoch,
719        owner: Some(owner),
720        // R736-T2: an appliance is not in a cell-move protocol, so there is no
721        // pointer generation to track. `0` is the documented unfenced default
722        // for that level; the epoch above is the level that is real here.
723        pointer_generation: 0,
724    };
725    let seam = CoreWalSeam::open_reader(path)
726        .with_context(|| format!("opening a read-only WAL seam on {path}"))?;
727    stream::tail_frames(&seam, &subject.target, &cfg)
728        .await
729        .with_context(|| format!("tailing {path} at epoch {epoch}"))
730}
731
732#[cfg(test)]
733mod tests {
734    use super::*;
735    use object_store::memory::InMemory;
736
737    /// A scratch volume that cleans up after itself.
738    struct Volume(PathBuf);
739
740    impl Volume {
741        fn new(tag: &str) -> Self {
742            let dir = std::env::temp_dir().join(format!(
743                "turso-backup-tail-{tag}-{}-{:?}",
744                std::process::id(),
745                std::thread::current().id()
746            ));
747            let _ = std::fs::remove_dir_all(&dir);
748            std::fs::create_dir_all(&dir).unwrap();
749            Self(dir)
750        }
751        fn path(&self) -> &Path {
752            &self.0
753        }
754        fn subject(&self, name: &str) -> String {
755            self.0.join(name).to_str().unwrap().to_string()
756        }
757    }
758
759    impl Drop for Volume {
760        fn drop(&mut self) {
761            let _ = std::fs::remove_dir_all(&self.0);
762        }
763    }
764
765    async fn seed(path: &str, start: i64, count: i64) {
766        let db = turso::Builder::new_local(path).build().await.unwrap();
767        let conn = db.connect().unwrap();
768        conn.execute("CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT)", ())
769            .await
770            .unwrap();
771        conn.execute("BEGIN", ()).await.unwrap();
772        for i in start..start + count {
773            conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (i, format!("v{i}")))
774                .await
775                .unwrap();
776        }
777        conn.execute("COMMIT", ()).await.unwrap();
778    }
779
780    async fn count_rows(path: &str) -> i64 {
781        let db = turso::Builder::new_local(path).build().await.unwrap();
782        let conn = db.connect().unwrap();
783        let mut r = conn.query("SELECT COUNT(*) FROM t", ()).await.unwrap();
784        r.next().await.unwrap().unwrap().get::<i64>(0).unwrap()
785    }
786
787    fn store() -> Arc<dyn ObjectStore> {
788        Arc::new(InMemory::new())
789    }
790
791    fn request<'a>(
792        store: Arc<dyn ObjectStore>,
793        vol: &'a Volume,
794        subjects: &'a [String],
795        tier: Tier,
796        owner: &'a str,
797    ) -> TailRequest<'a> {
798        TailRequest {
799            store,
800            store_prefix: "workloads/acct",
801            volume_root: vol.path(),
802            subjects,
803            tier,
804            owner,
805            page_size: 4096,
806            rpo_target: None,
807        }
808    }
809
810    /// The end-to-end this ticket's remaining half is for: a tail writes to the
811    /// store and a hydrate reads it back. Before this module a hydrate against a
812    /// real camp returned `nothing_in_the_store` forever, because nothing wrote.
813    #[tokio::test]
814    async fn a_tail_then_a_hydrate_round_trips_the_declared_state() {
815        let vol = Volume::new("roundtrip");
816        let subjects = vec!["accounts.db".to_string()];
817        seed(&vol.subject("accounts.db"), 0, 120).await;
818
819        let store = store();
820        let mut session =
821            match start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
822                .await
823                .unwrap()
824            {
825                Ok(s) => s,
826                Err(r) => panic!("start refused: {}", r.headline()),
827            };
828        assert_eq!(session.epoch(), claim::FIRST_EPOCH);
829        assert!(session.displaced().is_none(), "a virgin prefix displaces nobody");
830
831        let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
832            panic!("the first round should not be fenced");
833        };
834        assert_eq!(reports.len(), 1);
835        assert!(
836            matches!(&reports[0].outcome, SubjectOutcome::Stream { base_published: true, .. }),
837            "the first round publishes the base: {:?}",
838            reports[0].outcome
839        );
840
841        // More writes, then a second round: the base is reused, frames advance.
842        seed(&vol.subject("accounts.db"), 120, 80).await;
843        let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
844            panic!("the second round should not be fenced");
845        };
846        assert!(
847            matches!(&reports[0].outcome, SubjectOutcome::Stream { base_published: false, .. }),
848            "the base is published once, not per round: {:?}",
849            reports[0].outcome
850        );
851
852        // Now the consumer: hydrate an empty volume from what the tail wrote.
853        let dest = Volume::new("roundtrip-dest");
854        let outcome = crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
855            store,
856            store_prefix: "workloads/acct",
857            volume_root: dest.path(),
858            subjects: &subjects,
859            tier: Tier::Stream,
860            owner: "node-b",
861        })
862        .await
863        .unwrap();
864        assert!(
865            matches!(outcome, crate::hydrate::HydrateOutcome::Hydrated { .. }),
866            "the tail's output must be hydratable: {outcome:?}"
867        );
868        assert_eq!(
869            count_rows(&dest.subject("accounts.db")).await,
870            200,
871            "every row the application committed before the last round must come back"
872        );
873    }
874
875    /// The driving shape: three databases in one volume, all of them declared.
876    /// A workload backed up on two of three is not backed up.
877    #[tokio::test]
878    async fn every_declared_subject_is_backed_up_not_just_the_first() {
879        let vol = Volume::new("three");
880        let subjects = vec![
881            "accounts.db".to_string(),
882            "passkeys.db".to_string(),
883            "sessions.db".to_string(),
884        ];
885        for (i, s) in subjects.iter().enumerate() {
886            seed(&vol.subject(s), 0, 10 * (i as i64 + 1)).await;
887        }
888
889        let store = store();
890        let mut session = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
891            .await
892            .unwrap()
893            .expect("start");
894        let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
895            panic!("not fenced");
896        };
897        assert_eq!(reports.len(), 3);
898        for r in &reports {
899            assert!(
900                matches!(r.outcome, SubjectOutcome::Stream { .. }),
901                "{} was not streamed: {:?}",
902                r.subject,
903                r.outcome
904            );
905        }
906
907        let dest = Volume::new("three-dest");
908        crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
909            store,
910            store_prefix: "workloads/acct",
911            volume_root: dest.path(),
912            subjects: &subjects,
913            tier: Tier::Stream,
914            owner: "node-b",
915        })
916        .await
917        .unwrap();
918        for (i, s) in subjects.iter().enumerate() {
919            assert_eq!(
920                count_rows(&dest.subject(s)).await,
921                10 * (i as i64 + 1),
922                "subject {s} did not come back"
923            );
924        }
925    }
926
927    /// R858-B18's refusal must reach [`round`] as a reported non-event, not as a
928    /// dead backup. Provokes the REAL error rather than asserting on a copy of
929    /// its text, so a reword in `stream.rs` fails here instead of silently
930    /// reclassifying a busy database as a failure.
931    #[tokio::test]
932    async fn a_source_that_never_holds_still_is_too_hot_not_a_failure() {
933        let vol = Volume::new("hot");
934        let path = vol.subject("accounts.db");
935        seed(&path, 0, 10).await;
936
937        // A `take` that grows the source between the two fingerprint samples —
938        // every attempt sees movement, which is exactly the hammering-writer
939        // shape, without needing a second process to hammer.
940        let err = stream::validated_against_source(&path, "live copy", || async {
941            let mut f = std::fs::OpenOptions::new().append(true).open(&path).unwrap();
942            std::io::Write::write_all(&mut f, &vec![0u8; 4096]).unwrap();
943            Ok(Vec::<u8>::new())
944        })
945        .await
946        .expect_err("a source that moves under every attempt must be refused");
947        assert!(
948            is_source_too_hot(&err),
949            "the too-hot refusal was not recognised, so a busy appliance's tail would die \
950             instead of retrying next round: {err:#}"
951        );
952
953        // The negative half: an ordinary failure must NOT be classified as
954        // transient, or a genuinely broken subject retries forever in silence.
955        let missing = stream::raw_consistent_copy_live(&vol.subject("nope.db"), 4096)
956            .await
957            .expect_err("a missing database is a failure");
958        assert!(!is_source_too_hot(&missing), "{missing:#}");
959    }
960
961    /// Fold the WAL under a running tail and the restore must still work.
962    ///
963    /// This is the failure @Ashguard:hydra caught in review, and it is silent
964    /// without the rebase: `tail_frames` reports `Restarted`, the prefix ends up
965    /// holding manifests from two WAL generations, and
966    /// `validate_generation_chain` refuses the whole chain — so a hydrate of a
967    /// prefix that has been happily "streaming" for weeks fails, or (if the
968    /// refusal is ever relaxed) restores spliced frames. Both halves are
969    /// asserted: the chain restores AND it carries the post-fold rows.
970    #[tokio::test]
971    async fn a_wal_restart_re_anchors_the_chain_instead_of_breaking_the_restore() {
972        let vol = Volume::new("refold");
973        let subjects = vec!["accounts.db".to_string()];
974        let db = vol.subject("accounts.db");
975        seed(&db, 0, 50).await;
976
977        let store = store();
978        let mut session = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
979            .await
980            .unwrap()
981            .expect("start");
982        assert!(matches!(round(&mut session).await.unwrap(), RoundOutcome::Backed(_)));
983
984        // Fold the WAL into the main file and reset it — a checkpoint, which is
985        // what SQLite does on its own every `wal_autocheckpoint` pages. The test
986        // owns this database, so it can drive one directly; a real appliance's
987        // does it unprompted, which is why this cannot be designed away.
988        {
989            let d = turso::Builder::new_local(&db).build().await.unwrap();
990            let c = d.connect().unwrap();
991            let mut rows = c.query("PRAGMA wal_checkpoint(TRUNCATE)", ()).await.unwrap();
992            while rows.next().await.unwrap().is_some() {}
993        }
994        seed(&db, 50, 25).await;
995
996        let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
997            panic!("not fenced")
998        };
999        match &reports[0].outcome {
1000            SubjectOutcome::Stream { base_published, .. } => assert!(
1001                base_published,
1002                "a WAL restart must re-anchor onto a fresh base, not keep streaming onto the \
1003                 old one: {:?}",
1004                reports[0].outcome
1005            ),
1006            other => panic!("{other:?}"),
1007        }
1008
1009        let dest = Volume::new("refold-dest");
1010        crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
1011            store,
1012            store_prefix: "workloads/acct",
1013            volume_root: dest.path(),
1014            subjects: &subjects,
1015            tier: Tier::Stream,
1016            owner: "node-b",
1017        })
1018        .await
1019        .expect("a prefix that survived a WAL restart must still be restorable");
1020        assert_eq!(
1021            count_rows(&dest.subject("accounts.db")).await,
1022            75,
1023            "the restore came back at the wrong point in time"
1024        );
1025    }
1026
1027    /// A subject the application has not created yet is reported, not fatal.
1028    #[tokio::test]
1029    async fn an_absent_subject_is_reported_rather_than_failing_the_round() {
1030        let vol = Volume::new("absent");
1031        let subjects = vec!["accounts.db".to_string(), "not-yet.db".to_string()];
1032        seed(&vol.subject("accounts.db"), 0, 5).await;
1033
1034        let store = store();
1035        let mut session = start(request(store, &vol, &subjects, Tier::Stream, "node-a"))
1036            .await
1037            .unwrap()
1038            .expect("start");
1039        let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
1040            panic!("not fenced");
1041        };
1042        assert!(matches!(reports[0].outcome, SubjectOutcome::Stream { .. }));
1043        assert_eq!(reports[1].outcome, SubjectOutcome::Absent);
1044    }
1045
1046    /// The fence, from the losing side. A second node's `start` displaces the
1047    /// first, and the first finds out on its next round rather than continuing
1048    /// to write into somebody else's prefix.
1049    #[tokio::test]
1050    async fn a_displaced_tail_is_fenced_on_its_next_round() {
1051        let vol = Volume::new("fence");
1052        let subjects = vec!["accounts.db".to_string()];
1053        seed(&vol.subject("accounts.db"), 0, 20).await;
1054
1055        let store = store();
1056        let mut first = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
1057            .await
1058            .unwrap()
1059            .expect("start a");
1060        assert!(matches!(round(&mut first).await.unwrap(), RoundOutcome::Backed(_)));
1061
1062        let second = start(request(store, &vol, &subjects, Tier::Stream, "node-b"))
1063            .await
1064            .unwrap()
1065            .expect("start b");
1066        assert_eq!(second.epoch(), first.epoch() + 1, "a takeover is a monotonic bump");
1067        assert_eq!(
1068            second.displaced().map(|c| c.owner.as_str()),
1069            Some("node-a"),
1070            "the takeover records who it displaced"
1071        );
1072
1073        match round(&mut first).await.unwrap() {
1074            RoundOutcome::Fenced { detail } => {
1075                assert!(detail.contains("claim"), "unhelpful fence detail: {detail}")
1076            }
1077            other => panic!("the displaced tail kept writing: {other:?}"),
1078        }
1079    }
1080
1081    /// The reason [`round`] re-reads the claim rather than relying on the
1082    /// sink-side fence: tiers 1a and 1b stamp nothing, so *only* this check
1083    /// stops a fenced node from overwriting the real owner's backups.
1084    #[tokio::test]
1085    async fn a_tier_1_tail_is_fenced_too_even_though_its_writes_carry_no_epoch() {
1086        let vol = Volume::new("fence-t1");
1087        let subjects = vec!["accounts.db".to_string()];
1088        seed(&vol.subject("accounts.db"), 0, 20).await;
1089
1090        let store = store();
1091        let mut first = start(request(store.clone(), &vol, &subjects, Tier::Snapshot, "node-a"))
1092            .await
1093            .unwrap()
1094            .expect("start a");
1095        assert!(matches!(round(&mut first).await.unwrap(), RoundOutcome::Backed(_)));
1096        let _second = start(request(store, &vol, &subjects, Tier::Snapshot, "node-b"))
1097            .await
1098            .unwrap()
1099            .expect("start b");
1100        assert!(
1101            matches!(round(&mut first).await.unwrap(), RoundOutcome::Fenced { .. }),
1102            "a tier-1a tail that lost the claim must stop writing"
1103        );
1104    }
1105
1106    /// Tier 1a on an idle database must not re-upload the whole file every
1107    /// round — that is what [`snapshot::upload_snapshot_image`]'s gate is for.
1108    #[tokio::test]
1109    async fn an_idle_tier_1a_subject_stops_uploading() {
1110        let vol = Volume::new("idle");
1111        let subjects = vec!["accounts.db".to_string()];
1112        seed(&vol.subject("accounts.db"), 0, 30).await;
1113
1114        let store = store();
1115        let mut session = start(request(store, &vol, &subjects, Tier::Snapshot, "node-a"))
1116            .await
1117            .unwrap()
1118            .expect("start");
1119        let RoundOutcome::Backed(first) = round(&mut session).await.unwrap() else {
1120            panic!("not fenced")
1121        };
1122        assert!(
1123            matches!(first[0].outcome, SubjectOutcome::Snapshot(SnapshotOutcome::Uploaded { .. })),
1124            "{:?}",
1125            first[0].outcome
1126        );
1127        let RoundOutcome::Backed(second) = round(&mut session).await.unwrap() else {
1128            panic!("not fenced")
1129        };
1130        assert!(
1131            matches!(
1132                second[0].outcome,
1133                SubjectOutcome::Snapshot(SnapshotOutcome::Deduplicated { .. })
1134            ),
1135            "an untouched database was re-uploaded: {:?}",
1136            second[0].outcome
1137        );
1138    }
1139
1140    /// Tier 1b round-trips through the image-shaped entry point this ticket
1141    /// split out of [`dedup::snapshot_dedup`].
1142    #[tokio::test]
1143    async fn tier_1b_backs_up_a_live_database_and_hydrate_reads_it() {
1144        let vol = Volume::new("dedup");
1145        let subjects = vec!["accounts.db".to_string()];
1146        seed(&vol.subject("accounts.db"), 0, 60).await;
1147
1148        let store = store();
1149        let mut session = start(request(store.clone(), &vol, &subjects, Tier::Dedup, "node-a"))
1150            .await
1151            .unwrap()
1152            .expect("start");
1153        let RoundOutcome::Backed(reports) = round(&mut session).await.unwrap() else {
1154            panic!("not fenced")
1155        };
1156        assert!(
1157            matches!(reports[0].outcome, SubjectOutcome::Dedup(DedupOutcome::Snapshotted { .. })),
1158            "{:?}",
1159            reports[0].outcome
1160        );
1161
1162        let dest = Volume::new("dedup-dest");
1163        crate::hydrate::hydrate(crate::hydrate::HydrateRequest {
1164            store,
1165            store_prefix: "workloads/acct",
1166            volume_root: dest.path(),
1167            subjects: &subjects,
1168            tier: Tier::Dedup,
1169            owner: "node-b",
1170        })
1171        .await
1172        .unwrap();
1173        assert_eq!(count_rows(&dest.subject("accounts.db")).await, 60);
1174    }
1175
1176    /// A second incarnation on the same node — the ordinary restart — must not
1177    /// re-upload a base it already has. The hydrate that put the bytes there
1178    /// left a base in the prefix; a new session discovers it.
1179    #[tokio::test]
1180    async fn a_restarted_tail_adopts_the_existing_base_instead_of_republishing() {
1181        let vol = Volume::new("restart");
1182        let subjects = vec!["accounts.db".to_string()];
1183        seed(&vol.subject("accounts.db"), 0, 40).await;
1184
1185        let store = store();
1186        let mut first = start(request(store.clone(), &vol, &subjects, Tier::Stream, "node-a"))
1187            .await
1188            .unwrap()
1189            .expect("start a");
1190        assert!(matches!(round(&mut first).await.unwrap(), RoundOutcome::Backed(_)));
1191        drop(first);
1192
1193        let mut second = start(request(store, &vol, &subjects, Tier::Stream, "node-a"))
1194            .await
1195            .unwrap()
1196            .expect("restart a");
1197        let RoundOutcome::Backed(reports) = round(&mut second).await.unwrap() else {
1198            panic!("not fenced")
1199        };
1200        match &reports[0].outcome {
1201            SubjectOutcome::Stream { base_published, base_snapshot_key, .. } => {
1202                assert!(!base_published, "a restart re-copied the whole database for nothing");
1203                assert!(base_snapshot_key.contains("snapshots/"), "{base_snapshot_key}");
1204            }
1205            other => panic!("{other:?}"),
1206        }
1207    }
1208}