turso_backup/snapshot.rs
1//! Tier 1a — full snapshot sink.
2//!
3//! `VACUUM INTO` a temp file → upload the whole file to an object store
4//! (S3 / R2 / MinIO) → restore = download + open with SQLite. The DB-header
5//! `change_counter` gates skip-if-unchanged. Uses turso's public API only.
6//!
7//! @yah:relay(R003, "Tier 1a — full snapshot sink (VACUUM INTO → S3)")
8//! @yah:at(2026-05-26T22:28:30Z)
9//! @yah:status(open)
10//! @yah:phase(P1)
11//! @yah:parent(Q002)
12//! @arch:see(.yah/docs/working/turso-s3-backup.md)
13//!
14//! @yah:ticket(R003-F2, "Full-snapshot sink: VACUUM INTO temp + object_store upload (S3/R2/MinIO, path-style); change_counter skip-if-unchanged")
15//! @yah:assignee(bundle-anthropic-ashguard)
16//! @yah:at(2026-05-27T00:51:38Z)
17//! @yah:status(review)
18//! @yah:phase(P1)
19//! @yah:parent(R003)
20//! @arch:see(.yah/docs/working/turso-s3-backup.md)
21//! @yah:depends_on(R003-T1)
22//! @yah:handoff("Implemented snapshot_and_upload with a TWO-GATE skip guard: gate 1 = source-file hash (main + -wal) skips vacuum+upload on the quiescent path; gate 2 = canonical VACUUM INTO output hash (byte-deterministic, verified) dedups the upload when a checkpoint reshuffled bytes without changing content. VACUUM INTO unique temp -> object_store.put + two sidecars (latest.source-fingerprint, latest.snapshot-fingerprint). Returns SnapshotOutcome::{Uploaded,Unchanged,Deduplicated}. Backend-agnostic (Arc<dyn ObjectStore>) -> path-style S3/R2/MinIO is the caller's builder config.")
23//! @yah:handoff("DEVIATION from ticket title: change_counter is DEAD in turso — stays 1 across rollback/WAL/WAL+checkpoint (verified; VACUUM INTO doesn't copy it). Hashes replace it. Two-gate design follows user's steer (determinism worth the lift when not expensive); gate 1 keeps the idle path off the vacuum. Still in review for final sign-off. Working doc + Cargo.toml + rustdoc updated.")
24//! @yah:verify("cargo test -p turso-backup: 2 tests green on InMemory store — (1) upload -> SQLite-magic + turso 100-row readback -> gate-1 Unchanged -> mutate -> re-upload; (2) gate-2 Deduplicated path + fingerprint refresh. cargo build/clippy clean.")
25//! @yah:next("R003-F3: restore_latest (list snapshots/ by last-modified, get newest, write file). Layout: {prefix}/snapshots/snapshot-{unix_nanos:020}.db + {prefix}/latest.fingerprint.")
26//! @yah:next("R003-T4: MinIO/path-style e2e in harness; thread experimental_multiprocess_wal through vacuum_into if the writer's WAL format requires it.")
27//!
28//! @yah:ticket(R003-F3, "Restore path: download latest snapshot, open with sqlite3, verify")
29//! @yah:assignee(agent:claude)
30//! @yah:at(2026-05-26T22:30:04Z)
31//! @yah:status(review)
32//! @yah:phase(P1)
33//! @yah:parent(R003)
34//! @arch:see(.yah/docs/working/turso-s3-backup.md)
35//! @yah:depends_on(R003-F2)
36//! @yah:verify("cargo test -p turso-backup: 3 green on InMemory store (new restore_latest_picks_newest_and_round_trips). cargo clippy --all-targets clean.")
37//! @yah:handoff("Implemented restore_latest(target, dest_path) -> Result<String>: list_with_delimiter under {prefix}/snapshots/ (excludes the latest.* sidecars by construction), pick the lexically-greatest key (zero-padded nanos = newest; clock-skew-immune vs last_modified), get -> std::fs::write to dest. Returns the restored key. Errors if no snapshots. Moved the fn above mod tests (clippy items_after_test_module).")
38//! @yah:next("R003-T4: wire writer -> snapshot_and_upload -> MinIO -> restore_latest -> verifier (5000 rows, sqlite3 integrity_check). The literal 'open with sqlite3' verify lives in the harness; F3 covers it at the lib level (SQLite-magic + turso readback).")
39//!
40//! @yah:ticket(R003-T4, "Green end-to-end in harness: writer -> snapshot sink -> MinIO -> restore -> verifier checks (5000 rows)")
41//! @yah:assignee(bundle-anthropic-ashguard)
42//! @yah:at(2026-05-27T08:29:26Z)
43//! @yah:status(review)
44//! @yah:phase(P1)
45//! @yah:parent(R003)
46//! @arch:see(.yah/docs/working/turso-s3-backup.md)
47//! @yah:depends_on(R003-F3)
48//! @yah:verify("bash ci/run.sh (staged: minio -> minio-init -> writer -> verifier); verifier exit 0 = green")
49//! @yah:handoff("GREEN e2e at 5000 rows. Rewired the harness off Litestream onto turso-backup: new harness/ crate (writer + restore bins + shared backup_target lib, path-dep on turso-backup, object_store aws feature for MinIO). One multi-stage Dockerfile (shared builder, writer/verifier targets). Deleted dead Litestream writer/ + verifier/ dirs. Writer uploaded 446KB snapshot; verifier restored + passed sqlite3 integrity/rowcount/boundary/hash (matched byte-for-byte), exit 0.")
50
51use anyhow::{Context, Result};
52use object_store::path::Path as ObjPath;
53use object_store::{ObjectStore, ObjectStoreExt};
54use sha2::{Digest, Sha256};
55use std::sync::Arc;
56use std::time::{SystemTime, UNIX_EPOCH};
57use crate::stream::validated_against_source;
58
59/// Configuration for an object-store backup target (S3 / R2 / MinIO, path-style).
60///
61/// Wraps a generic `object_store::ObjectStore` so the caller picks the backend
62/// (e.g., `AmazonS3Builder::new().with_endpoint("http://minio:9000").build()`),
63/// then hands this crate a `BackupTarget` to use for snapshot operations.
64pub struct BackupTarget {
65 pub store: Arc<dyn ObjectStore>,
66 pub prefix: String,
67}
68
69impl BackupTarget {
70 /// Object key for a snapshot taken at `unix_nanos` (zero-padded so lexical
71 /// order matches chronological order — restore can pick the newest by name
72 /// or by last-modified).
73 fn snapshot_key(&self, unix_nanos: u128) -> ObjPath {
74 join_key(&self.prefix, &format!("snapshots/snapshot-{unix_nanos:020}.db"))
75 }
76
77 /// Sidecar recording the hash of the **source files** at the last run —
78 /// the cheap gate-1 short-circuit (skip vacuum+upload when bytes match).
79 fn source_fingerprint_key(&self) -> ObjPath {
80 join_key(&self.prefix, "latest.source-fingerprint")
81 }
82
83 /// Sidecar recording the hash of the last uploaded **snapshot** (the
84 /// canonical `VACUUM INTO` output) — the deterministic gate-2 dedup.
85 fn snapshot_fingerprint_key(&self) -> ObjPath {
86 join_key(&self.prefix, "latest.snapshot-fingerprint")
87 }
88}
89
90/// Result of a [`snapshot_and_upload`] call.
91#[derive(Debug, Clone, PartialEq, Eq)]
92pub enum SnapshotOutcome {
93 /// A new snapshot was written to `key` (`bytes` long). `snapshot_hash` is
94 /// the canonical hash of the uploaded file, now recorded in the sidecar.
95 Uploaded {
96 key: String,
97 bytes: usize,
98 snapshot_hash: String,
99 },
100 /// Gate 1: the source files were byte-identical to the last run, so nothing
101 /// was vacuumed or uploaded. `source_hash` is the unchanged source hash.
102 Unchanged { source_hash: String },
103 /// Gate 2: the source bytes moved (e.g. an auto-checkpoint reshuffled
104 /// pages) but the canonical `VACUUM INTO` output matched the last snapshot,
105 /// so the upload was skipped. The cheap source fingerprint is refreshed so
106 /// the next run short-circuits at gate 1.
107 Deduplicated { snapshot_hash: String },
108}
109
110/// Take a full `VACUUM INTO` snapshot of the turso database at `db_path` and
111/// upload the resulting file to the object store, behind a two-gate
112/// skip-if-unchanged guard.
113///
114/// 1. **Gate 1 (cheap):** hash the source files (`db_path` + `-wal`) and compare
115/// to the last run's sidecar. Identical bytes → return [`Unchanged`] without
116/// touching the DB. This is the common quiescent path.
117/// 2. Otherwise `VACUUM INTO` a unique temp file (point-in-time-consistent,
118/// fsync'd; the output is valid vanilla SQLite).
119/// 3. **Gate 2 (canonical):** `VACUUM INTO` output is byte-deterministic for
120/// identical content, so if its hash matches the last uploaded snapshot the
121/// source changed without the content changing → return [`Deduplicated`]
122/// (upload skipped, source fingerprint refreshed).
123/// 4. Otherwise upload the snapshot and record both fingerprints.
124///
125/// ## Why hashes, not the SQLite `change_counter`
126///
127/// The original plan used the DB-header `change_counter` (offset 24) as the
128/// "did anything change?" signal. The turso rewrite **does not maintain it** —
129/// it is left at the default `1` across every write in rollback, WAL, and
130/// WAL+checkpoint mode (verified empirically; `VACUUM INTO` doesn't even copy it
131/// through). Using it would skip every snapshot after the first. The two-gate
132/// hash design is conservative: it only ever skips on a proven match, so the
133/// failure direction is a redundant upload, never a missed backup.
134///
135/// [`Unchanged`]: SnapshotOutcome::Unchanged
136/// [`Deduplicated`]: SnapshotOutcome::Deduplicated
137pub async fn snapshot_and_upload(db_path: &str, target: &BackupTarget) -> Result<SnapshotOutcome> {
138 // Gate 1 (cheap): source files unchanged since last run -> nothing to do.
139 let source_hash = source_fingerprint(db_path)?;
140 let prev_source = read_text(&target.store, &target.source_fingerprint_key()).await?;
141 if prev_source.as_deref() == Some(source_hash.as_str()) {
142 return Ok(SnapshotOutcome::Unchanged { source_hash });
143 }
144
145 // Source moved -> take the canonical snapshot. SQLite refuses an existing
146 // VACUUM INTO target, so the nanosecond + pid name is unique and removed first.
147 let nanos = unix_nanos();
148 let temp_path = std::env::temp_dir().join(format!(
149 "turso-snapshot-{}-{nanos}.db",
150 std::process::id()
151 ));
152 let _ = std::fs::remove_file(&temp_path);
153 vacuum_into(db_path, &temp_path).await?;
154 let bytes = std::fs::read(&temp_path)
155 .with_context(|| format!("reading vacuum output {}", temp_path.display()))?;
156 let _ = std::fs::remove_file(&temp_path);
157
158 // Gate 2 (canonical): the vacuum output is byte-deterministic for identical
159 // content, so a match means the bytes moved without the content changing
160 // (e.g. an auto-checkpoint reshuffle). Skip the upload, but refresh the
161 // cheap source fingerprint so the next run short-circuits at gate 1.
162 let snapshot_hash = sha256_hex(&bytes);
163 let prev_snapshot = read_text(&target.store, &target.snapshot_fingerprint_key()).await?;
164 if prev_snapshot.as_deref() == Some(snapshot_hash.as_str()) {
165 put_text(&target.store, &target.source_fingerprint_key(), &source_hash).await?;
166 return Ok(SnapshotOutcome::Deduplicated { snapshot_hash });
167 }
168
169 // Genuinely new content -> upload, then record both fingerprints.
170 let key = target.snapshot_key(nanos);
171 let len = bytes.len();
172 target
173 .store
174 .put(&key, bytes.into())
175 .await
176 .with_context(|| format!("uploading snapshot to {key}"))?;
177 put_text(&target.store, &target.snapshot_fingerprint_key(), &snapshot_hash).await?;
178 put_text(&target.store, &target.source_fingerprint_key(), &source_hash).await?;
179
180 Ok(SnapshotOutcome::Uploaded {
181 key: key.to_string(),
182 bytes: len,
183 snapshot_hash,
184 })
185}
186
187/// Publish an already-materialized database image as a tier-1a snapshot,
188/// without opening the source database at all.
189///
190/// [`snapshot_and_upload`] is the right entry point when this process may open
191/// `db_path` freely. Two kinds of caller cannot:
192///
193/// - One that already holds the only permitted handle on the file. roadcase's
194/// residency ladder is the example: `turso_core::Database::open_file_with_flags`
195/// consults a **process-global registry keyed by file id** before it reads the
196/// open flags, so a second opener gets the first one's handle back with its
197/// own flags discarded. Such a caller has a live `turso_core::Connection` and
198/// must produce the image through it.
199/// - One that needs the snapshot to be *page-compatible* with the WAL frames
200/// [`crate::stream::tail_frames`] will upload next. `VACUUM INTO` repacks the
201/// database, so its output is a valid SQLite file but not necessarily the
202/// same page layout the engine will number subsequent frames against. The
203/// image a tier-2 base wants is the main database file itself, read after a
204/// `PRAGMA wal_checkpoint(TRUNCATE)` has emptied the WAL into it.
205///
206/// Writes one object under the same `snapshots/snapshot-{unix_nanos:020}.db`
207/// layout [`snapshot_and_upload`] uses, so [`restore_latest`] finds it and
208/// [`crate::stream::StreamConfig::base_snapshot_key`] can name it. Returns that
209/// key.
210///
211/// Deliberately writes **neither fingerprint sidecar.** Both gate
212/// `snapshot_and_upload`, and this path has hashed no source files — recording
213/// a `latest.source-fingerprint` here would assert a correspondence between an
214/// object and a set of file bytes that was never checked, which is the one
215/// thing that could make the two-gate skip drop a real backup.
216pub async fn upload_base_snapshot(target: &BackupTarget, image: &[u8]) -> Result<String> {
217 anyhow::ensure!(
218 image.starts_with(b"SQLite format 3\0"),
219 "refusing to publish a base snapshot that is not a SQLite database \
220 ({} bytes, prefix {:?})",
221 image.len(),
222 &image[..image.len().min(16)],
223 );
224 let key = target.snapshot_key(unix_nanos());
225 target
226 .store
227 .put(&key, image.to_vec().into())
228 .await
229 .with_context(|| format!("uploading base snapshot to {key}"))?;
230 Ok(key.to_string())
231}
232
233/// [`upload_base_snapshot`] behind the canonical (gate-2) skip, for a caller
234/// that publishes an image *repeatedly*.
235///
236/// R850-F1. `upload_base_snapshot` is a one-shot: every call writes an object,
237/// which is right for publishing a tier-2 base once and wrong for a tail that
238/// re-snapshots a mostly-idle database every interval forever. This applies the
239/// content hash gate `snapshot_and_upload` calls gate 2 — identical image bytes
240/// to the last publish means skip the upload — and returns
241/// [`SnapshotOutcome::Deduplicated`] when it fires.
242///
243/// It records **only** `latest.snapshot-fingerprint`, never
244/// `latest.source-fingerprint`, for the reason `upload_base_snapshot` gives at
245/// length: this path hashes no source *files*, so claiming a correspondence
246/// between an object and a set of file bytes it never read is the one thing that
247/// could make `snapshot_and_upload`'s cheap gate 1 skip a real backup later.
248/// [`SnapshotOutcome::Unchanged`] is therefore never returned from here.
249pub async fn upload_snapshot_image(
250 target: &BackupTarget,
251 image: &[u8],
252) -> Result<SnapshotOutcome> {
253 let snapshot_hash = sha256_hex(image);
254 let prev = read_text(&target.store, &target.snapshot_fingerprint_key()).await?;
255 if prev.as_deref() == Some(snapshot_hash.as_str()) {
256 return Ok(SnapshotOutcome::Deduplicated { snapshot_hash });
257 }
258 let key = upload_base_snapshot(target, image).await?;
259 put_text(
260 &target.store,
261 &target.snapshot_fingerprint_key(),
262 &snapshot_hash,
263 )
264 .await?;
265 Ok(SnapshotOutcome::Uploaded {
266 key,
267 bytes: image.len(),
268 snapshot_hash,
269 })
270}
271
272/// SHA-256 of the source database files: the main `.db` plus its `-wal`
273/// sidecar if present (WAL-mode commits live there until a checkpoint). A
274/// quiescent database hashes identically on re-read; any committed write
275/// changes the result. Returned as a lowercase hex string.
276fn source_fingerprint(db_path: &str) -> Result<String> {
277 let mut hasher = Sha256::new();
278 let main = std::fs::read(db_path).with_context(|| format!("reading source db {db_path}"))?;
279 hasher.update(&main);
280 if let Ok(wal) = std::fs::read(format!("{db_path}-wal")) {
281 hasher.update(b"\0-wal\0");
282 hasher.update(&wal);
283 }
284 Ok(hex::encode(hasher.finalize()))
285}
286
287/// Lowercase-hex SHA-256 of an in-memory buffer (used for the canonical
288/// snapshot fingerprint).
289fn sha256_hex(bytes: &[u8]) -> String {
290 let mut hasher = Sha256::new();
291 hasher.update(bytes);
292 hex::encode(hasher.finalize())
293}
294
295/// Read a small text sidecar (a fingerprint), or `None` if it doesn't exist yet.
296async fn read_text(store: &Arc<dyn ObjectStore>, key: &ObjPath) -> Result<Option<String>> {
297 match store.get(key).await {
298 Ok(res) => {
299 let bytes = res.bytes().await.context("reading fingerprint object")?;
300 Ok(Some(String::from_utf8_lossy(&bytes).trim().to_string()))
301 }
302 Err(object_store::Error::NotFound { .. }) => Ok(None),
303 Err(e) => Err(e).context("fetching fingerprint object"),
304 }
305}
306
307/// Write a small text sidecar (a fingerprint).
308async fn put_text(store: &Arc<dyn ObjectStore>, key: &ObjPath, text: &str) -> Result<()> {
309 store
310 .put(key, text.to_owned().into_bytes().into())
311 .await
312 .with_context(|| format!("recording fingerprint {key}"))?;
313 Ok(())
314}
315
316/// `VACUUM INTO '<temp_path>'` against the database at `db_path`, opened
317/// **read-only** and validated against a concurrent foreign writer.
318///
319/// R858-B18 changed two things here, for one reason: a backup is a reader, and
320/// this was the second of the crate's two source opens that did not say so.
321///
322/// 1. **`turso_core` with `OpenFlags::ReadOnly`, not `turso::Builder`.** The
323/// friendly wrapper exposes no `OpenFlags`, so it opened the source writable
324/// and took turso's whole-file exclusive `fcntl` lock — which is refused
325/// outright when the application that owns the database is running
326/// (measured: `examples/foreign_checkpoint_probe.rs` probe S1, `Locking
327/// error: File is locked by another process`). `ReadOnly` skips that lock per
328/// handle, and `VACUUM INTO` still runs on such a connection (probe S2: 50/50
329/// rows, `integrity_check ok`) — note `VACUUM INTO` is not gated by
330/// `DatabaseOpts::enable_vacuum`, only bare `VACUUM` is
331/// (`turso_core/translate/vacuum.rs:39`).
332/// 2. **Wrapped in [`validated_against_source`].** Opening without a lock is not
333/// coordination: a foreign checkpoint can still land mid-vacuum, and VACUUM
334/// INTO rebuilds the b-tree, so a torn read of the source yields an output
335/// that passes `integrity_check` while holding the wrong rows. Same protocol
336/// as tier 2's copy — accept only if the source provably held still, else
337/// retry, else refuse.
338///
339/// The destination is removed at the start of every attempt, because SQLite
340/// refuses to `VACUUM INTO` a file that already exists and a retried attempt
341/// would otherwise fail on its predecessor's output.
342async fn vacuum_into(db_path: &str, temp_path: &std::path::Path) -> Result<()> {
343 // Single-quote the path SQLite-style (double any embedded quote).
344 let escaped = temp_path.to_string_lossy().replace('\'', "''");
345 validated_against_source(db_path, "VACUUM INTO snapshot", || async {
346 let _ = std::fs::remove_file(temp_path);
347 let io: Arc<dyn turso_core::IO> =
348 Arc::new(turso_core::PlatformIO::new().context("creating turso_core PlatformIO")?);
349 let db = turso_core::Database::open_file_with_flags(
350 io,
351 db_path,
352 turso_core::OpenFlags::ReadOnly,
353 turso_core::DatabaseOpts::new(),
354 None,
355 )
356 .with_context(|| format!("opening turso db {db_path} read-only"))?;
357 let conn = db.connect().context("connecting to turso db")?;
358 conn.execute(format!("VACUUM INTO '{escaped}'"))
359 .context("VACUUM INTO failed")?;
360 Ok(())
361 })
362 .await
363}
364
365/// Join an object-store prefix and a leaf into a normalized [`ObjPath`],
366/// tolerating empty/slash-padded prefixes.
367fn join_key(prefix: &str, leaf: &str) -> ObjPath {
368 let prefix = prefix.trim_matches('/');
369 if prefix.is_empty() {
370 ObjPath::from(leaf)
371 } else {
372 ObjPath::from(format!("{prefix}/{leaf}"))
373 }
374}
375
376fn unix_nanos() -> u128 {
377 SystemTime::now()
378 .duration_since(UNIX_EPOCH)
379 .map(|d| d.as_nanos())
380 .unwrap_or(0)
381}
382
383/// Download the most recent snapshot and write it to `dest_path`, returning the
384/// object key it was restored from.
385///
386/// Lists the snapshot objects under `{prefix}/snapshots/` and picks the newest.
387/// Snapshot keys are zero-padded nanosecond timestamps (see
388/// [`BackupTarget::snapshot_key`]), so the lexically-greatest key is the most
389/// recent. We order by the key's baked-in capture time rather than the store's
390/// `last_modified` so the choice is deterministic and immune to clock skew
391/// between the snapshotting host and the object store.
392///
393/// The `snapshots/` prefix naturally excludes the `latest.*-fingerprint`
394/// sidecars (they live one level up at `{prefix}/`). The downloaded bytes are a
395/// valid vanilla-SQLite file written verbatim to `dest_path`, ready to open with
396/// `sqlite3` or re-open through turso.
397///
398/// Errors if no snapshots exist under the prefix yet.
399pub async fn restore_latest(target: &BackupTarget, dest_path: &str) -> Result<String> {
400 let snapshots_prefix = join_key(&target.prefix, "snapshots");
401 let newest = latest_snapshot_key(target)
402 .await?
403 .with_context(|| format!("no snapshots found under {snapshots_prefix}"))?;
404
405 let location = ObjPath::from(newest.clone());
406 let bytes = target
407 .store
408 .get(&location)
409 .await
410 .with_context(|| format!("downloading snapshot {location}"))?
411 .bytes()
412 .await
413 .with_context(|| format!("reading snapshot body {location}"))?;
414
415 std::fs::write(dest_path, &bytes)
416 .with_context(|| format!("writing restored snapshot to {dest_path}"))?;
417
418 Ok(newest)
419}
420
421/// The key [`restore_latest`] would download, without downloading it — or `None`
422/// when the prefix holds no snapshot at all.
423///
424/// R850-F1 split this out for the tail side, which needs the *name* of an
425/// existing tier-2 base rather than its bytes: a restarted tail that cannot see
426/// that a base already exists re-copies the whole database on every restart, and
427/// a `None` here is what tells it there is genuinely nothing to anchor to.
428///
429/// Absence is a `None` rather than the `Err` [`restore_latest`] returns because
430/// the two callers want opposite things from it: a restore with no snapshot has
431/// failed, while a first-ever tail with no snapshot is simply first.
432pub async fn latest_snapshot_key(target: &BackupTarget) -> Result<Option<String>> {
433 let snapshots_prefix = join_key(&target.prefix, "snapshots");
434 let listing = target
435 .store
436 .list_with_delimiter(Some(&snapshots_prefix))
437 .await
438 .with_context(|| format!("listing snapshots under {snapshots_prefix}"))?;
439 Ok(listing
440 .objects
441 .into_iter()
442 .max_by(|a, b| a.location.cmp(&b.location))
443 .map(|o| o.location.to_string()))
444}
445
446#[cfg(test)]
447mod tests {
448 use super::*;
449 use object_store::memory::InMemory;
450 // The library itself no longer opens a source database through the friendly
451 // wrapper (R858-B18 moved that to turso_core + OpenFlags::ReadOnly); these
452 // tests still use it to *write* their fixtures, which is the right surface
453 // for a writer.
454 use turso::Builder;
455
456 /// A throwaway db path under the OS temp dir, cleaned up on drop.
457 struct TempDb(std::path::PathBuf);
458 impl TempDb {
459 fn new(tag: &str) -> Self {
460 let p = std::env::temp_dir().join(format!(
461 "turso-backup-test-{}-{}-{}.db",
462 tag,
463 std::process::id(),
464 unix_nanos()
465 ));
466 TempDb(p)
467 }
468 fn path(&self) -> &str {
469 self.0.to_str().unwrap()
470 }
471 }
472 impl Drop for TempDb {
473 fn drop(&mut self) {
474 for suffix in ["", "-wal", "-shm"] {
475 let _ = std::fs::remove_file(format!("{}{suffix}", self.0.display()));
476 }
477 }
478 }
479
480 async fn seed_db(path: &str, rows: i64) {
481 let db = Builder::new_local(path).build().await.unwrap();
482 let conn = db.connect().unwrap();
483 conn.execute("CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT)", ())
484 .await
485 .unwrap();
486 conn.execute("BEGIN", ()).await.unwrap();
487 for i in 0..rows {
488 conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (i, format!("v{i}")))
489 .await
490 .unwrap();
491 }
492 conn.execute("COMMIT", ()).await.unwrap();
493 }
494
495 async fn count_rows(path: &str) -> i64 {
496 let db = Builder::new_local(path).build().await.unwrap();
497 let conn = db.connect().unwrap();
498 let mut r = conn.query("SELECT COUNT(*) FROM t", ()).await.unwrap();
499 let row = r.next().await.unwrap().unwrap();
500 row.get::<i64>(0).unwrap()
501 }
502
503 #[tokio::test]
504 async fn uploads_then_skips_unchanged_then_reuploads_on_change() {
505 let src = TempDb::new("src");
506 seed_db(src.path(), 100).await;
507
508 let target = BackupTarget {
509 store: Arc::new(InMemory::new()),
510 prefix: "backups".into(),
511 };
512
513 // First snapshot: uploaded.
514 let out = snapshot_and_upload(src.path(), &target).await.unwrap();
515 let (key, snap1) = match out {
516 SnapshotOutcome::Uploaded { key, bytes, snapshot_hash } => {
517 assert!(bytes > 0, "snapshot should be non-empty");
518 assert!(key.starts_with("backups/snapshots/snapshot-"), "key was {key}");
519 (key, snapshot_hash)
520 }
521 other => panic!("expected Uploaded, got {other:?}"),
522 };
523
524 // The uploaded object is a real vanilla-SQLite file.
525 let bytes = target
526 .store
527 .get(&ObjPath::from(key.clone()))
528 .await
529 .unwrap()
530 .bytes()
531 .await
532 .unwrap();
533 assert!(
534 bytes.starts_with(b"SQLite format 3\0"),
535 "uploaded object is not a SQLite database"
536 );
537
538 // ...and round-trips through turso with the right row count.
539 let restored = TempDb::new("restored");
540 std::fs::write(restored.path(), &bytes).unwrap();
541 assert_eq!(count_rows(restored.path()).await, 100);
542
543 // Gate 1: no source change -> skipped without vacuuming.
544 match snapshot_and_upload(src.path(), &target).await.unwrap() {
545 SnapshotOutcome::Unchanged { .. } => {}
546 other => panic!("expected Unchanged, got {other:?}"),
547 }
548
549 // Mutate the source, then snapshot again: uploaded, new snapshot hash.
550 {
551 let db = Builder::new_local(src.path()).build().await.unwrap();
552 let conn = db.connect().unwrap();
553 conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (10_000, "new"))
554 .await
555 .unwrap();
556 }
557 match snapshot_and_upload(src.path(), &target).await.unwrap() {
558 SnapshotOutcome::Uploaded { snapshot_hash, .. } => assert_ne!(snapshot_hash, snap1),
559 other => panic!("expected Uploaded after change, got {other:?}"),
560 }
561 }
562
563 /// Gate 2: when the source bytes look changed but the canonical VACUUM
564 /// output matches the last snapshot, the upload is deduplicated. We simulate
565 /// the "source moved without content change" case by deleting the cheap
566 /// source-fingerprint sidecar so gate 1 misses and the vacuum runs.
567 #[tokio::test]
568 async fn deduplicates_when_vacuum_output_matches() {
569 let src = TempDb::new("dedup-src");
570 seed_db(src.path(), 50).await;
571
572 let target = BackupTarget {
573 store: Arc::new(InMemory::new()),
574 prefix: String::new(),
575 };
576
577 let snap1 = match snapshot_and_upload(src.path(), &target).await.unwrap() {
578 SnapshotOutcome::Uploaded { snapshot_hash, .. } => snapshot_hash,
579 other => panic!("expected Uploaded, got {other:?}"),
580 };
581
582 // Force gate 1 to miss without changing content.
583 target
584 .store
585 .delete(&target.source_fingerprint_key())
586 .await
587 .unwrap();
588
589 match snapshot_and_upload(src.path(), &target).await.unwrap() {
590 SnapshotOutcome::Deduplicated { snapshot_hash } => assert_eq!(snapshot_hash, snap1),
591 other => panic!("expected Deduplicated, got {other:?}"),
592 }
593
594 // And the refreshed source fingerprint makes the next run short-circuit.
595 match snapshot_and_upload(src.path(), &target).await.unwrap() {
596 SnapshotOutcome::Unchanged { .. } => {}
597 other => panic!("expected Unchanged after dedup refresh, got {other:?}"),
598 }
599 }
600
601 /// restore_latest errors when empty, then downloads the newest of several
602 /// snapshots and round-trips it back through turso with the latest row count.
603 #[tokio::test]
604 async fn restore_latest_picks_newest_and_round_trips() {
605 let src = TempDb::new("restore-src");
606 seed_db(src.path(), 100).await;
607
608 let target = BackupTarget {
609 store: Arc::new(InMemory::new()),
610 prefix: "backups".into(),
611 };
612
613 // No snapshots yet -> error (and no file written to dest).
614 let dest = TempDb::new("restore-dest");
615 assert!(
616 restore_latest(&target, dest.path()).await.is_err(),
617 "restore should fail when no snapshots exist"
618 );
619
620 // First snapshot (100 rows).
621 snapshot_and_upload(src.path(), &target).await.unwrap();
622
623 // Mutate, then a second snapshot (101 rows) -> a lexically-newer key.
624 {
625 let db = Builder::new_local(src.path()).build().await.unwrap();
626 let conn = db.connect().unwrap();
627 conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (10_000, "new"))
628 .await
629 .unwrap();
630 }
631 let newest_key = match snapshot_and_upload(src.path(), &target).await.unwrap() {
632 SnapshotOutcome::Uploaded { key, .. } => key,
633 other => panic!("expected Uploaded, got {other:?}"),
634 };
635
636 // restore_latest reports it pulled the newest snapshot...
637 let restored_from = restore_latest(&target, dest.path()).await.unwrap();
638 assert_eq!(restored_from, newest_key, "should restore the newest snapshot");
639
640 // ...the on-disk file is a valid vanilla-SQLite database...
641 let on_disk = std::fs::read(dest.path()).unwrap();
642 assert!(
643 on_disk.starts_with(b"SQLite format 3\0"),
644 "restored file is not a SQLite database"
645 );
646
647 // ...and re-opens through turso with the post-mutation row count.
648 assert_eq!(count_rows(dest.path()).await, 101);
649 }
650}