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}