Skip to main content

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}