turso-backup 0.8.41

Turso/libSQL WAL-tailing snapshot and S3-compatible backup sink: periodic snapshot_and_upload plus a raw WAL frame streamer via turso_core's conn_raw_api.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
//! Tier 1a — full snapshot sink.
//!
//! `VACUUM INTO` a temp file → upload the whole file to an object store
//! (S3 / R2 / MinIO) → restore = download + open with SQLite. The DB-header
//! `change_counter` gates skip-if-unchanged. Uses turso's public API only.
//!
//! @yah:relay(R003, "Tier 1a — full snapshot sink (VACUUM INTO → S3)")
//! @yah:at(2026-05-26T22:28:30Z)
//! @yah:status(open)
//! @yah:phase(P1)
//! @yah:parent(Q002)
//! @arch:see(.yah/docs/working/turso-s3-backup.md)
//!
//! @yah:ticket(R003-F2, "Full-snapshot sink: VACUUM INTO temp + object_store upload (S3/R2/MinIO, path-style); change_counter skip-if-unchanged")
//! @yah:assignee(bundle-anthropic-ashguard)
//! @yah:at(2026-05-27T00:51:38Z)
//! @yah:status(review)
//! @yah:phase(P1)
//! @yah:parent(R003)
//! @arch:see(.yah/docs/working/turso-s3-backup.md)
//! @yah:depends_on(R003-T1)
//! @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.")
//! @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.")
//! @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.")
//! @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.")
//! @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.")
//!
//! @yah:ticket(R003-F3, "Restore path: download latest snapshot, open with sqlite3, verify")
//! @yah:assignee(agent:claude)
//! @yah:at(2026-05-26T22:30:04Z)
//! @yah:status(review)
//! @yah:phase(P1)
//! @yah:parent(R003)
//! @arch:see(.yah/docs/working/turso-s3-backup.md)
//! @yah:depends_on(R003-F2)
//! @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.")
//! @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).")
//! @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).")
//!
//! @yah:ticket(R003-T4, "Green end-to-end in harness: writer -> snapshot sink -> MinIO -> restore -> verifier checks (5000 rows)")
//! @yah:assignee(bundle-anthropic-ashguard)
//! @yah:at(2026-05-27T08:29:26Z)
//! @yah:status(review)
//! @yah:phase(P1)
//! @yah:parent(R003)
//! @arch:see(.yah/docs/working/turso-s3-backup.md)
//! @yah:depends_on(R003-F3)
//! @yah:verify("bash ci/run.sh (staged: minio -> minio-init -> writer -> verifier); verifier exit 0 = green")
//! @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.")

use anyhow::{Context, Result};
use object_store::path::Path as ObjPath;
use object_store::{ObjectStore, ObjectStoreExt};
use sha2::{Digest, Sha256};
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use crate::stream::validated_against_source;

/// Configuration for an object-store backup target (S3 / R2 / MinIO, path-style).
///
/// Wraps a generic `object_store::ObjectStore` so the caller picks the backend
/// (e.g., `AmazonS3Builder::new().with_endpoint("http://minio:9000").build()`),
/// then hands this crate a `BackupTarget` to use for snapshot operations.
pub struct BackupTarget {
    pub store: Arc<dyn ObjectStore>,
    pub prefix: String,
}

impl BackupTarget {
    /// Object key for a snapshot taken at `unix_nanos` (zero-padded so lexical
    /// order matches chronological order — restore can pick the newest by name
    /// or by last-modified).
    fn snapshot_key(&self, unix_nanos: u128) -> ObjPath {
        join_key(&self.prefix, &format!("snapshots/snapshot-{unix_nanos:020}.db"))
    }

    /// Sidecar recording the hash of the **source files** at the last run —
    /// the cheap gate-1 short-circuit (skip vacuum+upload when bytes match).
    fn source_fingerprint_key(&self) -> ObjPath {
        join_key(&self.prefix, "latest.source-fingerprint")
    }

    /// Sidecar recording the hash of the last uploaded **snapshot** (the
    /// canonical `VACUUM INTO` output) — the deterministic gate-2 dedup.
    fn snapshot_fingerprint_key(&self) -> ObjPath {
        join_key(&self.prefix, "latest.snapshot-fingerprint")
    }
}

/// Result of a [`snapshot_and_upload`] call.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SnapshotOutcome {
    /// A new snapshot was written to `key` (`bytes` long). `snapshot_hash` is
    /// the canonical hash of the uploaded file, now recorded in the sidecar.
    Uploaded {
        key: String,
        bytes: usize,
        snapshot_hash: String,
    },
    /// Gate 1: the source files were byte-identical to the last run, so nothing
    /// was vacuumed or uploaded. `source_hash` is the unchanged source hash.
    Unchanged { source_hash: String },
    /// Gate 2: the source bytes moved (e.g. an auto-checkpoint reshuffled
    /// pages) but the canonical `VACUUM INTO` output matched the last snapshot,
    /// so the upload was skipped. The cheap source fingerprint is refreshed so
    /// the next run short-circuits at gate 1.
    Deduplicated { snapshot_hash: String },
}

/// Take a full `VACUUM INTO` snapshot of the turso database at `db_path` and
/// upload the resulting file to the object store, behind a two-gate
/// skip-if-unchanged guard.
///
/// 1. **Gate 1 (cheap):** hash the source files (`db_path` + `-wal`) and compare
///    to the last run's sidecar. Identical bytes → return [`Unchanged`] without
///    touching the DB. This is the common quiescent path.
/// 2. Otherwise `VACUUM INTO` a unique temp file (point-in-time-consistent,
///    fsync'd; the output is valid vanilla SQLite).
/// 3. **Gate 2 (canonical):** `VACUUM INTO` output is byte-deterministic for
///    identical content, so if its hash matches the last uploaded snapshot the
///    source changed without the content changing → return [`Deduplicated`]
///    (upload skipped, source fingerprint refreshed).
/// 4. Otherwise upload the snapshot and record both fingerprints.
///
/// ## Why hashes, not the SQLite `change_counter`
///
/// The original plan used the DB-header `change_counter` (offset 24) as the
/// "did anything change?" signal. The turso rewrite **does not maintain it** —
/// it is left at the default `1` across every write in rollback, WAL, and
/// WAL+checkpoint mode (verified empirically; `VACUUM INTO` doesn't even copy it
/// through). Using it would skip every snapshot after the first. The two-gate
/// hash design is conservative: it only ever skips on a proven match, so the
/// failure direction is a redundant upload, never a missed backup.
///
/// [`Unchanged`]: SnapshotOutcome::Unchanged
/// [`Deduplicated`]: SnapshotOutcome::Deduplicated
pub async fn snapshot_and_upload(db_path: &str, target: &BackupTarget) -> Result<SnapshotOutcome> {
    // Gate 1 (cheap): source files unchanged since last run -> nothing to do.
    let source_hash = source_fingerprint(db_path)?;
    let prev_source = read_text(&target.store, &target.source_fingerprint_key()).await?;
    if prev_source.as_deref() == Some(source_hash.as_str()) {
        return Ok(SnapshotOutcome::Unchanged { source_hash });
    }

    // Source moved -> take the canonical snapshot. SQLite refuses an existing
    // VACUUM INTO target, so the nanosecond + pid name is unique and removed first.
    let nanos = unix_nanos();
    let temp_path = std::env::temp_dir().join(format!(
        "turso-snapshot-{}-{nanos}.db",
        std::process::id()
    ));
    let _ = std::fs::remove_file(&temp_path);
    vacuum_into(db_path, &temp_path).await?;
    let bytes = std::fs::read(&temp_path)
        .with_context(|| format!("reading vacuum output {}", temp_path.display()))?;
    let _ = std::fs::remove_file(&temp_path);

    // Gate 2 (canonical): the vacuum output is byte-deterministic for identical
    // content, so a match means the bytes moved without the content changing
    // (e.g. an auto-checkpoint reshuffle). Skip the upload, but refresh the
    // cheap source fingerprint so the next run short-circuits at gate 1.
    let snapshot_hash = sha256_hex(&bytes);
    let prev_snapshot = read_text(&target.store, &target.snapshot_fingerprint_key()).await?;
    if prev_snapshot.as_deref() == Some(snapshot_hash.as_str()) {
        put_text(&target.store, &target.source_fingerprint_key(), &source_hash).await?;
        return Ok(SnapshotOutcome::Deduplicated { snapshot_hash });
    }

    // Genuinely new content -> upload, then record both fingerprints.
    let key = target.snapshot_key(nanos);
    let len = bytes.len();
    target
        .store
        .put(&key, bytes.into())
        .await
        .with_context(|| format!("uploading snapshot to {key}"))?;
    put_text(&target.store, &target.snapshot_fingerprint_key(), &snapshot_hash).await?;
    put_text(&target.store, &target.source_fingerprint_key(), &source_hash).await?;

    Ok(SnapshotOutcome::Uploaded {
        key: key.to_string(),
        bytes: len,
        snapshot_hash,
    })
}

/// Publish an already-materialized database image as a tier-1a snapshot,
/// without opening the source database at all.
///
/// [`snapshot_and_upload`] is the right entry point when this process may open
/// `db_path` freely. Two kinds of caller cannot:
///
/// - One that already holds the only permitted handle on the file. roadcase's
///   residency ladder is the example: `turso_core::Database::open_file_with_flags`
///   consults a **process-global registry keyed by file id** before it reads the
///   open flags, so a second opener gets the first one's handle back with its
///   own flags discarded. Such a caller has a live `turso_core::Connection` and
///   must produce the image through it.
/// - One that needs the snapshot to be *page-compatible* with the WAL frames
///   [`crate::stream::tail_frames`] will upload next. `VACUUM INTO` repacks the
///   database, so its output is a valid SQLite file but not necessarily the
///   same page layout the engine will number subsequent frames against. The
///   image a tier-2 base wants is the main database file itself, read after a
///   `PRAGMA wal_checkpoint(TRUNCATE)` has emptied the WAL into it.
///
/// Writes one object under the same `snapshots/snapshot-{unix_nanos:020}.db`
/// layout [`snapshot_and_upload`] uses, so [`restore_latest`] finds it and
/// [`crate::stream::StreamConfig::base_snapshot_key`] can name it. Returns that
/// key.
///
/// Deliberately writes **neither fingerprint sidecar.** Both gate
/// `snapshot_and_upload`, and this path has hashed no source files — recording
/// a `latest.source-fingerprint` here would assert a correspondence between an
/// object and a set of file bytes that was never checked, which is the one
/// thing that could make the two-gate skip drop a real backup.
pub async fn upload_base_snapshot(target: &BackupTarget, image: &[u8]) -> Result<String> {
    anyhow::ensure!(
        image.starts_with(b"SQLite format 3\0"),
        "refusing to publish a base snapshot that is not a SQLite database \
         ({} bytes, prefix {:?})",
        image.len(),
        &image[..image.len().min(16)],
    );
    let key = target.snapshot_key(unix_nanos());
    target
        .store
        .put(&key, image.to_vec().into())
        .await
        .with_context(|| format!("uploading base snapshot to {key}"))?;
    Ok(key.to_string())
}

/// [`upload_base_snapshot`] behind the canonical (gate-2) skip, for a caller
/// that publishes an image *repeatedly*.
///
/// R850-F1. `upload_base_snapshot` is a one-shot: every call writes an object,
/// which is right for publishing a tier-2 base once and wrong for a tail that
/// re-snapshots a mostly-idle database every interval forever. This applies the
/// content hash gate `snapshot_and_upload` calls gate 2 — identical image bytes
/// to the last publish means skip the upload — and returns
/// [`SnapshotOutcome::Deduplicated`] when it fires.
///
/// It records **only** `latest.snapshot-fingerprint`, never
/// `latest.source-fingerprint`, for the reason `upload_base_snapshot` gives at
/// length: this path hashes no source *files*, so claiming a correspondence
/// between an object and a set of file bytes it never read is the one thing that
/// could make `snapshot_and_upload`'s cheap gate 1 skip a real backup later.
/// [`SnapshotOutcome::Unchanged`] is therefore never returned from here.
pub async fn upload_snapshot_image(
    target: &BackupTarget,
    image: &[u8],
) -> Result<SnapshotOutcome> {
    let snapshot_hash = sha256_hex(image);
    let prev = read_text(&target.store, &target.snapshot_fingerprint_key()).await?;
    if prev.as_deref() == Some(snapshot_hash.as_str()) {
        return Ok(SnapshotOutcome::Deduplicated { snapshot_hash });
    }
    let key = upload_base_snapshot(target, image).await?;
    put_text(
        &target.store,
        &target.snapshot_fingerprint_key(),
        &snapshot_hash,
    )
    .await?;
    Ok(SnapshotOutcome::Uploaded {
        key,
        bytes: image.len(),
        snapshot_hash,
    })
}

/// SHA-256 of the source database files: the main `.db` plus its `-wal`
/// sidecar if present (WAL-mode commits live there until a checkpoint). A
/// quiescent database hashes identically on re-read; any committed write
/// changes the result. Returned as a lowercase hex string.
fn source_fingerprint(db_path: &str) -> Result<String> {
    let mut hasher = Sha256::new();
    let main = std::fs::read(db_path).with_context(|| format!("reading source db {db_path}"))?;
    hasher.update(&main);
    if let Ok(wal) = std::fs::read(format!("{db_path}-wal")) {
        hasher.update(b"\0-wal\0");
        hasher.update(&wal);
    }
    Ok(hex::encode(hasher.finalize()))
}

/// Lowercase-hex SHA-256 of an in-memory buffer (used for the canonical
/// snapshot fingerprint).
fn sha256_hex(bytes: &[u8]) -> String {
    let mut hasher = Sha256::new();
    hasher.update(bytes);
    hex::encode(hasher.finalize())
}

/// Read a small text sidecar (a fingerprint), or `None` if it doesn't exist yet.
async fn read_text(store: &Arc<dyn ObjectStore>, key: &ObjPath) -> Result<Option<String>> {
    match store.get(key).await {
        Ok(res) => {
            let bytes = res.bytes().await.context("reading fingerprint object")?;
            Ok(Some(String::from_utf8_lossy(&bytes).trim().to_string()))
        }
        Err(object_store::Error::NotFound { .. }) => Ok(None),
        Err(e) => Err(e).context("fetching fingerprint object"),
    }
}

/// Write a small text sidecar (a fingerprint).
async fn put_text(store: &Arc<dyn ObjectStore>, key: &ObjPath, text: &str) -> Result<()> {
    store
        .put(key, text.to_owned().into_bytes().into())
        .await
        .with_context(|| format!("recording fingerprint {key}"))?;
    Ok(())
}

/// `VACUUM INTO '<temp_path>'` against the database at `db_path`, opened
/// **read-only** and validated against a concurrent foreign writer.
///
/// R858-B18 changed two things here, for one reason: a backup is a reader, and
/// this was the second of the crate's two source opens that did not say so.
///
/// 1. **`turso_core` with `OpenFlags::ReadOnly`, not `turso::Builder`.** The
///    friendly wrapper exposes no `OpenFlags`, so it opened the source writable
///    and took turso's whole-file exclusive `fcntl` lock — which is refused
///    outright when the application that owns the database is running
///    (measured: `examples/foreign_checkpoint_probe.rs` probe S1, `Locking
///    error: File is locked by another process`). `ReadOnly` skips that lock per
///    handle, and `VACUUM INTO` still runs on such a connection (probe S2: 50/50
///    rows, `integrity_check ok`) — note `VACUUM INTO` is not gated by
///    `DatabaseOpts::enable_vacuum`, only bare `VACUUM` is
///    (`turso_core/translate/vacuum.rs:39`).
/// 2. **Wrapped in [`validated_against_source`].** Opening without a lock is not
///    coordination: a foreign checkpoint can still land mid-vacuum, and VACUUM
///    INTO rebuilds the b-tree, so a torn read of the source yields an output
///    that passes `integrity_check` while holding the wrong rows. Same protocol
///    as tier 2's copy — accept only if the source provably held still, else
///    retry, else refuse.
///
/// The destination is removed at the start of every attempt, because SQLite
/// refuses to `VACUUM INTO` a file that already exists and a retried attempt
/// would otherwise fail on its predecessor's output.
async fn vacuum_into(db_path: &str, temp_path: &std::path::Path) -> Result<()> {
    // Single-quote the path SQLite-style (double any embedded quote).
    let escaped = temp_path.to_string_lossy().replace('\'', "''");
    validated_against_source(db_path, "VACUUM INTO snapshot", || async {
        let _ = std::fs::remove_file(temp_path);
        let io: Arc<dyn turso_core::IO> =
            Arc::new(turso_core::PlatformIO::new().context("creating turso_core PlatformIO")?);
        let db = turso_core::Database::open_file_with_flags(
            io,
            db_path,
            turso_core::OpenFlags::ReadOnly,
            turso_core::DatabaseOpts::new(),
            None,
        )
        .with_context(|| format!("opening turso db {db_path} read-only"))?;
        let conn = db.connect().context("connecting to turso db")?;
        conn.execute(format!("VACUUM INTO '{escaped}'"))
            .context("VACUUM INTO failed")?;
        Ok(())
    })
    .await
}

/// Join an object-store prefix and a leaf into a normalized [`ObjPath`],
/// tolerating empty/slash-padded prefixes.
fn join_key(prefix: &str, leaf: &str) -> ObjPath {
    let prefix = prefix.trim_matches('/');
    if prefix.is_empty() {
        ObjPath::from(leaf)
    } else {
        ObjPath::from(format!("{prefix}/{leaf}"))
    }
}

fn unix_nanos() -> u128 {
    SystemTime::now()
        .duration_since(UNIX_EPOCH)
        .map(|d| d.as_nanos())
        .unwrap_or(0)
}

/// Download the most recent snapshot and write it to `dest_path`, returning the
/// object key it was restored from.
///
/// Lists the snapshot objects under `{prefix}/snapshots/` and picks the newest.
/// Snapshot keys are zero-padded nanosecond timestamps (see
/// [`BackupTarget::snapshot_key`]), so the lexically-greatest key is the most
/// recent. We order by the key's baked-in capture time rather than the store's
/// `last_modified` so the choice is deterministic and immune to clock skew
/// between the snapshotting host and the object store.
///
/// The `snapshots/` prefix naturally excludes the `latest.*-fingerprint`
/// sidecars (they live one level up at `{prefix}/`). The downloaded bytes are a
/// valid vanilla-SQLite file written verbatim to `dest_path`, ready to open with
/// `sqlite3` or re-open through turso.
///
/// Errors if no snapshots exist under the prefix yet.
pub async fn restore_latest(target: &BackupTarget, dest_path: &str) -> Result<String> {
    let snapshots_prefix = join_key(&target.prefix, "snapshots");
    let newest = latest_snapshot_key(target)
        .await?
        .with_context(|| format!("no snapshots found under {snapshots_prefix}"))?;

    let location = ObjPath::from(newest.clone());
    let bytes = target
        .store
        .get(&location)
        .await
        .with_context(|| format!("downloading snapshot {location}"))?
        .bytes()
        .await
        .with_context(|| format!("reading snapshot body {location}"))?;

    std::fs::write(dest_path, &bytes)
        .with_context(|| format!("writing restored snapshot to {dest_path}"))?;

    Ok(newest)
}

/// The key [`restore_latest`] would download, without downloading it — or `None`
/// when the prefix holds no snapshot at all.
///
/// R850-F1 split this out for the tail side, which needs the *name* of an
/// existing tier-2 base rather than its bytes: a restarted tail that cannot see
/// that a base already exists re-copies the whole database on every restart, and
/// a `None` here is what tells it there is genuinely nothing to anchor to.
///
/// Absence is a `None` rather than the `Err` [`restore_latest`] returns because
/// the two callers want opposite things from it: a restore with no snapshot has
/// failed, while a first-ever tail with no snapshot is simply first.
pub async fn latest_snapshot_key(target: &BackupTarget) -> Result<Option<String>> {
    let snapshots_prefix = join_key(&target.prefix, "snapshots");
    let listing = target
        .store
        .list_with_delimiter(Some(&snapshots_prefix))
        .await
        .with_context(|| format!("listing snapshots under {snapshots_prefix}"))?;
    Ok(listing
        .objects
        .into_iter()
        .max_by(|a, b| a.location.cmp(&b.location))
        .map(|o| o.location.to_string()))
}

#[cfg(test)]
mod tests {
    use super::*;
    use object_store::memory::InMemory;
    // The library itself no longer opens a source database through the friendly
    // wrapper (R858-B18 moved that to turso_core + OpenFlags::ReadOnly); these
    // tests still use it to *write* their fixtures, which is the right surface
    // for a writer.
    use turso::Builder;

    /// A throwaway db path under the OS temp dir, cleaned up on drop.
    struct TempDb(std::path::PathBuf);
    impl TempDb {
        fn new(tag: &str) -> Self {
            let p = std::env::temp_dir().join(format!(
                "turso-backup-test-{}-{}-{}.db",
                tag,
                std::process::id(),
                unix_nanos()
            ));
            TempDb(p)
        }
        fn path(&self) -> &str {
            self.0.to_str().unwrap()
        }
    }
    impl Drop for TempDb {
        fn drop(&mut self) {
            for suffix in ["", "-wal", "-shm"] {
                let _ = std::fs::remove_file(format!("{}{suffix}", self.0.display()));
            }
        }
    }

    async fn seed_db(path: &str, rows: i64) {
        let db = Builder::new_local(path).build().await.unwrap();
        let conn = db.connect().unwrap();
        conn.execute("CREATE TABLE IF NOT EXISTS t (id INTEGER PRIMARY KEY, v TEXT)", ())
            .await
            .unwrap();
        conn.execute("BEGIN", ()).await.unwrap();
        for i in 0..rows {
            conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (i, format!("v{i}")))
                .await
                .unwrap();
        }
        conn.execute("COMMIT", ()).await.unwrap();
    }

    async fn count_rows(path: &str) -> i64 {
        let db = Builder::new_local(path).build().await.unwrap();
        let conn = db.connect().unwrap();
        let mut r = conn.query("SELECT COUNT(*) FROM t", ()).await.unwrap();
        let row = r.next().await.unwrap().unwrap();
        row.get::<i64>(0).unwrap()
    }

    #[tokio::test]
    async fn uploads_then_skips_unchanged_then_reuploads_on_change() {
        let src = TempDb::new("src");
        seed_db(src.path(), 100).await;

        let target = BackupTarget {
            store: Arc::new(InMemory::new()),
            prefix: "backups".into(),
        };

        // First snapshot: uploaded.
        let out = snapshot_and_upload(src.path(), &target).await.unwrap();
        let (key, snap1) = match out {
            SnapshotOutcome::Uploaded { key, bytes, snapshot_hash } => {
                assert!(bytes > 0, "snapshot should be non-empty");
                assert!(key.starts_with("backups/snapshots/snapshot-"), "key was {key}");
                (key, snapshot_hash)
            }
            other => panic!("expected Uploaded, got {other:?}"),
        };

        // The uploaded object is a real vanilla-SQLite file.
        let bytes = target
            .store
            .get(&ObjPath::from(key.clone()))
            .await
            .unwrap()
            .bytes()
            .await
            .unwrap();
        assert!(
            bytes.starts_with(b"SQLite format 3\0"),
            "uploaded object is not a SQLite database"
        );

        // ...and round-trips through turso with the right row count.
        let restored = TempDb::new("restored");
        std::fs::write(restored.path(), &bytes).unwrap();
        assert_eq!(count_rows(restored.path()).await, 100);

        // Gate 1: no source change -> skipped without vacuuming.
        match snapshot_and_upload(src.path(), &target).await.unwrap() {
            SnapshotOutcome::Unchanged { .. } => {}
            other => panic!("expected Unchanged, got {other:?}"),
        }

        // Mutate the source, then snapshot again: uploaded, new snapshot hash.
        {
            let db = Builder::new_local(src.path()).build().await.unwrap();
            let conn = db.connect().unwrap();
            conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (10_000, "new"))
                .await
                .unwrap();
        }
        match snapshot_and_upload(src.path(), &target).await.unwrap() {
            SnapshotOutcome::Uploaded { snapshot_hash, .. } => assert_ne!(snapshot_hash, snap1),
            other => panic!("expected Uploaded after change, got {other:?}"),
        }
    }

    /// Gate 2: when the source bytes look changed but the canonical VACUUM
    /// output matches the last snapshot, the upload is deduplicated. We simulate
    /// the "source moved without content change" case by deleting the cheap
    /// source-fingerprint sidecar so gate 1 misses and the vacuum runs.
    #[tokio::test]
    async fn deduplicates_when_vacuum_output_matches() {
        let src = TempDb::new("dedup-src");
        seed_db(src.path(), 50).await;

        let target = BackupTarget {
            store: Arc::new(InMemory::new()),
            prefix: String::new(),
        };

        let snap1 = match snapshot_and_upload(src.path(), &target).await.unwrap() {
            SnapshotOutcome::Uploaded { snapshot_hash, .. } => snapshot_hash,
            other => panic!("expected Uploaded, got {other:?}"),
        };

        // Force gate 1 to miss without changing content.
        target
            .store
            .delete(&target.source_fingerprint_key())
            .await
            .unwrap();

        match snapshot_and_upload(src.path(), &target).await.unwrap() {
            SnapshotOutcome::Deduplicated { snapshot_hash } => assert_eq!(snapshot_hash, snap1),
            other => panic!("expected Deduplicated, got {other:?}"),
        }

        // And the refreshed source fingerprint makes the next run short-circuit.
        match snapshot_and_upload(src.path(), &target).await.unwrap() {
            SnapshotOutcome::Unchanged { .. } => {}
            other => panic!("expected Unchanged after dedup refresh, got {other:?}"),
        }
    }

    /// restore_latest errors when empty, then downloads the newest of several
    /// snapshots and round-trips it back through turso with the latest row count.
    #[tokio::test]
    async fn restore_latest_picks_newest_and_round_trips() {
        let src = TempDb::new("restore-src");
        seed_db(src.path(), 100).await;

        let target = BackupTarget {
            store: Arc::new(InMemory::new()),
            prefix: "backups".into(),
        };

        // No snapshots yet -> error (and no file written to dest).
        let dest = TempDb::new("restore-dest");
        assert!(
            restore_latest(&target, dest.path()).await.is_err(),
            "restore should fail when no snapshots exist"
        );

        // First snapshot (100 rows).
        snapshot_and_upload(src.path(), &target).await.unwrap();

        // Mutate, then a second snapshot (101 rows) -> a lexically-newer key.
        {
            let db = Builder::new_local(src.path()).build().await.unwrap();
            let conn = db.connect().unwrap();
            conn.execute("INSERT INTO t (id, v) VALUES (?, ?)", (10_000, "new"))
                .await
                .unwrap();
        }
        let newest_key = match snapshot_and_upload(src.path(), &target).await.unwrap() {
            SnapshotOutcome::Uploaded { key, .. } => key,
            other => panic!("expected Uploaded, got {other:?}"),
        };

        // restore_latest reports it pulled the newest snapshot...
        let restored_from = restore_latest(&target, dest.path()).await.unwrap();
        assert_eq!(restored_from, newest_key, "should restore the newest snapshot");

        // ...the on-disk file is a valid vanilla-SQLite database...
        let on_disk = std::fs::read(dest.path()).unwrap();
        assert!(
            on_disk.starts_with(b"SQLite format 3\0"),
            "restored file is not a SQLite database"
        );

        // ...and re-opens through turso with the post-mutation row count.
        assert_eq!(count_rows(dest.path()).await, 101);
    }
}