use anyhow::{Context, Result};
use object_store::memory::InMemory;
use std::io::{BufRead, BufReader, Read, Write};
use std::process::{Command, Stdio};
use std::sync::Arc;
use std::time::{Duration, Instant};
use turso_backup::backpressure::BackpressureConfig;
use turso_backup::snapshot::{snapshot_and_upload, upload_base_snapshot, BackupTarget};
use turso_backup::stream::{
raw_consistent_copy_live, restore_latest_stream, tail_frames, CoreWalSeam, SourceFingerprint,
StreamConfig, StreamOutcome, WalSeam, Watermark,
};
const PAGE_SIZE: usize = 4096;
fn sqlite3_bin() -> String {
std::env::var("PROBE_SQLITE3").unwrap_or_else(|_| "/usr/bin/sqlite3".to_string())
}
#[derive(Debug)]
#[allow(dead_code)]
enum ForeignRun {
Ran { elapsed: Duration },
Blocked { waited: Duration, stderr: String },
Died { stderr: String },
}
fn sqlite3_kill9(db: &str, sql: &str, deadline: Duration) -> Result<ForeignRun> {
let started = Instant::now();
let mut child = Command::new(sqlite3_bin())
.arg(db)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.with_context(|| format!("spawning {}", sqlite3_bin()))?;
let pid = child.id();
let mut sin = child.stdin.take().expect("piped stdin");
let sout = child.stdout.take().expect("piped stdout");
let wd = std::thread::spawn(move || {
std::thread::sleep(deadline);
let _ = Command::new("kill").arg("-9").arg(pid.to_string()).status();
});
writeln!(sin, "{sql}")?;
writeln!(sin, "SELECT 'PROBE-SENTINEL';")?;
sin.flush()?;
let mut rd = BufReader::new(sout);
let mut line = String::new();
let mut saw_sentinel = false;
loop {
line.clear();
if rd.read_line(&mut line)? == 0 {
break;
}
if line.contains("PROBE-SENTINEL") {
saw_sentinel = true;
break;
}
}
let elapsed = started.elapsed();
let _ = child.kill();
let _ = child.wait();
drop(sin);
let mut stderr = String::new();
if let Some(mut e) = child.stderr.take() {
let _ = e.read_to_string(&mut stderr);
}
let _ = wd.join();
Ok(if saw_sentinel {
ForeignRun::Ran { elapsed }
} else if elapsed >= deadline.saturating_sub(Duration::from_millis(250)) {
ForeignRun::Blocked { waited: elapsed, stderr: stderr.trim().to_string() }
} else {
ForeignRun::Died { stderr: stderr.trim().to_string() }
})
}
fn sqlite3_clean(db: &str, sql: &str) -> Result<String> {
let out = Command::new(sqlite3_bin())
.arg(db)
.arg(sql)
.output()
.with_context(|| format!("running {} {db}", sqlite3_bin()))?;
anyhow::ensure!(
out.status.success(),
"sqlite3 failed on {sql:?}: {}",
String::from_utf8_lossy(&out.stderr)
);
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
}
fn row_count(db: &str) -> Result<i64> {
Ok(sqlite3_clean(db, "SELECT count(*) FROM t;")?.parse()?)
}
fn rows_via_copy(db: &str, tag: &str) -> Result<i64> {
let c = format!("{db}.rd-{tag}");
for sfx in ["", "-wal", "-shm"] {
let _ = std::fs::copy(format!("{db}{sfx}"), format!("{c}{sfx}"));
}
let n = row_count(&c);
for sfx in ["", "-wal", "-shm"] {
let _ = std::fs::remove_file(format!("{c}{sfx}"));
}
n
}
#[derive(Debug, PartialEq, Eq, Clone, Copy)]
struct WalHdr {
pgsz: u32,
ckpt_seq: u32,
salt1: u32,
salt2: u32,
}
fn wal_hdr(db: &str) -> Option<WalHdr> {
let b = std::fs::read(format!("{db}-wal")).ok()?;
if b.len() < 32 {
return None;
}
let be = |o: usize| u32::from_be_bytes([b[o], b[o + 1], b[o + 2], b[o + 3]]);
Some(WalHdr { pgsz: be(8), ckpt_seq: be(12), salt1: be(16), salt2: be(20) })
}
fn hdr_str(db: &str) -> String {
match wal_hdr(db) {
Some(h) => format!(
"pgsz {} ckpt_seq {} salt {:08x}/{:08x}",
h.pgsz, h.ckpt_seq, h.salt1, h.salt2
),
None => "<no WAL header (file absent or truncated)>".to_string(),
}
}
fn file_len(p: &str) -> u64 {
std::fs::metadata(p).map(|m| m.len()).unwrap_or(0)
}
fn turso_watermark(db: &str) -> Result<Watermark> {
CoreWalSeam::open(db)?.wal_state()
}
fn state_line(db: &str, tag: &str) -> Result<String> {
let w = turso_watermark(db)?;
Ok(format!(
"{tag}: turso (seq {}, max_frame {}) | wal-hdr {} | wal {}B main {}B",
w.checkpoint_seq,
w.last_frame,
hdr_str(db),
file_len(&format!("{db}-wal")),
file_len(db)
))
}
struct Probe {
dir: std::path::PathBuf,
}
impl Probe {
fn new(tag: &str) -> Result<Self> {
let dir = std::env::temp_dir().join(format!(
"turso-backup-foreign-probe-{tag}-{}",
std::process::id()
));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir)?;
Ok(Self { dir })
}
fn path(&self, name: &str) -> String {
self.dir.join(name).to_str().unwrap().to_string()
}
}
fn cfg(base_snapshot_key: &str) -> StreamConfig<'_> {
StreamConfig {
base_snapshot_key,
page_size: PAGE_SIZE,
backpressure: BackpressureConfig::default(),
rpo_target: None,
epoch: 0,
owner: Some("R858-S6-probe"),
pointer_generation: 0,
}
}
fn target() -> BackupTarget {
BackupTarget { store: Arc::new(InMemory::new()), prefix: "probe".into() }
}
fn seed_clean(db: &str, n: i64) -> Result<()> {
let mut sql = String::from(
"PRAGMA journal_mode=WAL;\nCREATE TABLE IF NOT EXISTS t(id INTEGER PRIMARY KEY, v TEXT);\nBEGIN;\n",
);
for i in 0..n {
sql.push_str(&format!("INSERT INTO t(id,v) VALUES({i},'v{i}');\n"));
}
sql.push_str("COMMIT;\n");
sqlite3_clean(db, &sql)?;
Ok(())
}
fn append_dirty(db: &str, start: i64, n: i64, txns: i64, pad: usize) -> Result<ForeignRun> {
let mut sql = String::new();
let per = (n / txns).max(1);
let mut i = start;
while i < start + n {
sql.push_str("BEGIN;\n");
for j in i..(i + per).min(start + n) {
sql.push_str(&format!("INSERT INTO t(id,v) VALUES({j},'{}');\n", "x".repeat(pad)));
}
sql.push_str("COMMIT;\n");
i += per;
}
sqlite3_kill9(db, &sql, Duration::from_secs(30))
}
fn probe_a() -> Result<()> {
println!("\n=== PROBE A: does turso observe a foreign checkpoint? (process-restart regime) ===");
println!("NOTE: the checkpointing connection here CLOSES, which makes SQLite delete the WAL.");
println!(" The mode label is therefore not the variable under test — see probe E for the");
println!(" in-process autocheckpoint regime, where checkpoint_seq DOES advance.");
for mode in ["TRUNCATE", "RESTART", "FULL"] {
let p = Probe::new(&format!("a-{mode}"))?;
let db = p.path("a.db");
seed_clean(&db, 3)?;
append_dirty(&db, 100, 3, 1, 8)?;
let before = turso_watermark(&db)?;
let hb = wal_hdr(&db);
println!("{}", state_line(&db, &format!("A1[{mode}] pre-fold "))?);
let r = sqlite3_clean(&db, &format!("PRAGMA wal_checkpoint({mode});"))?;
println!("A2[{mode}] foreign `wal_checkpoint({mode})` returned {r}");
println!("{}", state_line(&db, &format!("A3[{mode}] post-fold"))?);
append_dirty(&db, 200, 3, 1, 8)?;
let after = turso_watermark(&db)?;
let ha = wal_hdr(&db);
println!("{}", state_line(&db, &format!("A4[{mode}] post-fold+write"))?);
println!(
"A[{mode}] VERDICT: turso checkpoint_seq {} -> {} (advanced: {}) | on-disk ckpt_seq {:?} -> {:?} | salt changed: {}",
before.checkpoint_seq,
after.checkpoint_seq,
after.checkpoint_seq > before.checkpoint_seq,
hb.map(|h| h.ckpt_seq),
ha.map(|h| h.ckpt_seq),
hb.map(|h| (h.salt1, h.salt2)) != ha.map(|h| (h.salt1, h.salt2)),
);
}
Ok(())
}
async fn tail_across_fold(
tag: &str,
pre_fold_rows: i64,
pre_fold_txns: i64,
post_fold_rows: i64,
post_fold_txns: i64,
orphan_early: bool,
) -> Result<()> {
let p = Probe::new(tag)?;
let db = p.path("s.db");
let tgt = target();
seed_clean(&db, 3)?;
let base = upload_base_snapshot(&tgt, &raw_consistent_copy_live(&db, PAGE_SIZE).await?).await?;
append_dirty(&db, 100, pre_fold_rows, pre_fold_txns, 400)?;
println!("{}", state_line(&db, " pre-tail#1")?);
let o1 = tail_frames(&CoreWalSeam::open(&db)?, &tgt, &cfg(&base)).await?;
println!(" tail#1 -> {}", short(&o1));
sqlite3_clean(&db, "PRAGMA wal_checkpoint(TRUNCATE);")?;
if orphan_early {
let mut sql =
String::from("CREATE TABLE IF NOT EXISTS u(id INTEGER PRIMARY KEY, v TEXT);\nBEGIN;\n");
for i in 0..40 {
sql.push_str(&format!("INSERT INTO u(id,v) VALUES({i},'{}');\n", "u".repeat(400)));
}
sql.push_str("COMMIT;\n");
sqlite3_kill9(&db, &sql, Duration::from_secs(30))?;
println!("{}", state_line(&db, " post-fold, after the orphan-prefix write")?);
}
append_dirty(&db, 200, post_fold_rows, post_fold_txns, 400)?;
println!("{}", state_line(&db, " post-fold")?);
let o2 = tail_frames(&CoreWalSeam::open(&db)?, &tgt, &cfg(&base)).await?;
println!(" tail#2 (across the fold) -> {}", short(&o2));
println!(
" restart signalled? {}",
matches!(o2, StreamOutcome::Restarted { .. })
);
let src = rows_via_copy(&db, "final")?;
let src_u = if orphan_early {
sqlite3_clean(&p.path("s.db.rd-u"), "SELECT 1").ok();
let c = format!("{db}.rd-u");
for sfx in ["", "-wal", "-shm"] {
let _ = std::fs::copy(format!("{db}{sfx}"), format!("{c}{sfx}"));
}
let n = sqlite3_clean(&c, "SELECT count(*) FROM u;").ok();
for sfx in ["", "-wal", "-shm"] {
let _ = std::fs::remove_file(format!("{c}{sfx}"));
}
n
} else {
None
};
let dest = p.path("restored.db");
match restore_latest_stream(&tgt, &dest).await {
Ok(out) => {
let integrity = sqlite3_clean(&dest, "PRAGMA integrity_check;");
let n = row_count(&dest);
let n_u = if orphan_early {
sqlite3_clean(&dest, "SELECT count(*) FROM u;").ok()
} else {
None
};
println!(
" restore SUCCEEDED (frames_replayed {}, generations {}); integrity_check {:?}; t rows {:?} vs source {src}; u rows {:?} vs source {:?}",
out.frames_replayed, out.generation_count, integrity, n, n_u, src_u
);
let t_ok = matches!(n, Ok(n) if n == src);
let u_ok = !orphan_early || (n_u.is_some() && n_u == src_u);
match (&n, t_ok && u_ok && matches!(integrity.as_deref(), Ok("ok"))) {
(Ok(_), true) => println!(" VERDICT[{tag}]: CORRECT"),
(Ok(n), false) => println!(
" VERDICT[{tag}]: SILENT WRONG IMAGE — restored t={n}/{src}, u={n_u:?}/{src_u:?}, integrity {integrity:?}. No error was raised anywhere."
),
(Err(e), _) => {
println!(" VERDICT[{tag}]: restored image UNREADABLE by upstream sqlite3: {e}")
}
}
}
Err(e) => println!(" restore REFUSED (fail-loud): {e:#}\n VERDICT[{tag}]: safe-but-stalled"),
}
Ok(())
}
fn short(o: &StreamOutcome) -> String {
match o {
StreamOutcome::Empty { watermark, .. } => format!(
"Empty (seq {}, last_frame {})",
watermark.checkpoint_seq, watermark.last_frame
),
StreamOutcome::Streamed { first_frame, last_frame, checkpoint_seq, frame_count, .. } => {
format!("Streamed seq {checkpoint_seq} frames {first_frame}..={last_frame} ({frame_count})")
}
StreamOutcome::Restarted {
previous_generation, new_generation, first_frame, last_frame, ..
} => format!(
"Restarted gen (seq {}, salt {})->(seq {}, salt {}) frames {first_frame}..={last_frame}",
previous_generation.checkpoint_seq,
previous_generation.salt.map_or("<unknown>".to_string(), |s| s.to_string()),
new_generation.checkpoint_seq,
new_generation.salt.map_or("<unknown>".to_string(), |s| s.to_string()),
),
other => format!("{other:?}"),
}
}
async fn probe_c() -> Result<()> {
println!("\n=== PROBE C: hydration via raw_consistent_copy_live ===");
let p = Probe::new("c")?;
let db = p.path("c.db");
seed_clean(&db, 5)?;
append_dirty(&db, 100, 4, 2, 400)?;
println!("{}", state_line(&db, "C1 foreign source")?);
let src = rows_via_copy(&db, "pre")?;
let img = raw_consistent_copy_live(&db, PAGE_SIZE).await?;
let out = p.path("c-image.db");
std::fs::write(&out, &img)?;
let integrity = sqlite3_clean(&out, "PRAGMA integrity_check;");
let rows = row_count(&out);
println!(
"C2 image {} bytes; header ok {}; NO -wal sidecar written: {}; integrity_check {:?}; rows {:?} vs live source {src}",
img.len(),
img.starts_with(b"SQLite format 3\0"),
!std::path::Path::new(&format!("{out}-wal")).exists(),
integrity,
rows
);
println!(
"C VERDICT: {}",
if rows.as_ref().ok() == Some(&src) && matches!(integrity.as_deref(), Ok("ok")) {
"CORRECT — a foreign-written uncheckpointed WAL hydrates to a vanilla-SQLite image upstream sqlite3 reads"
} else {
"WRONG"
}
);
println!("{}", state_line(&db, "C3 source after the copy (must be unperturbed)")?);
Ok(())
}
fn probe_d() -> Result<()> {
println!("\n=== PROBE D: foreign writer liveness while a CoreWalSeam is open ===");
let p = Probe::new("d")?;
let db = p.path("d.db");
seed_clean(&db, 3)?;
let ctl = append_dirty(&db, 50, 3, 1, 8)?;
println!("D0 control (no seam held): {ctl:?}");
let seam = CoreWalSeam::open(&db)?;
let before = seam.wal_state()?;
println!("D1 seam open; wal_state = (seq {}, max_frame {})", before.checkpoint_seq, before.last_frame);
let w = sqlite3_kill9(
&db,
"PRAGMA busy_timeout=2000;\nBEGIN IMMEDIATE;\nINSERT INTO t(id,v) VALUES(999,'x');\nCOMMIT;",
Duration::from_secs(10),
)?;
println!("D2 foreign WRITE while seam held: {w:?}");
let started = Instant::now();
let ck = Command::new(sqlite3_bin())
.arg(&db)
.arg("PRAGMA busy_timeout=2000; PRAGMA wal_checkpoint(TRUNCATE);")
.output()?;
println!(
"D3 foreign CHECKPOINT while seam held after {:?}: status {} stdout {:?} stderr {:?}",
started.elapsed(),
ck.status,
String::from_utf8_lossy(&ck.stdout).trim(),
String::from_utf8_lossy(&ck.stderr).trim()
);
let after = seam.wal_state()?;
println!(
"D4 the SAME never-reopened seam now reports (seq {}, max_frame {}) — it {} the foreign activity",
after.checkpoint_seq,
after.last_frame,
if after == before { "does NOT see" } else { "sees" }
);
drop(seam);
println!("{}", state_line(&db, "D5 after dropping the seam")?);
Ok(())
}
async fn probe_g() -> Result<()> {
println!("\n=== PROBE G: can turso open a DB a foreign process merely holds open? ===");
let p = Probe::new("g")?;
let db = p.path("g.db");
seed_clean(&db, 3)?;
let mut child = Command::new(sqlite3_bin())
.arg(&db)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()?;
let mut sin = child.stdin.take().expect("piped stdin");
let sout = child.stdout.take().expect("piped stdout");
let mut rd = BufReader::new(sout);
writeln!(sin, "SELECT count(*) FROM t;\nSELECT 'PROBE-SENTINEL';")?;
sin.flush()?;
let mut line = String::new();
loop {
line.clear();
anyhow::ensure!(rd.read_line(&mut line)? > 0, "foreign holder died early");
if line.contains("PROBE-SENTINEL") {
break;
}
}
println!("G1 foreign sqlite3 connection is open and idle on the DB");
match CoreWalSeam::open(&db) {
Ok(s) => println!(
"G2 CoreWalSeam::open SUCCEEDED alongside it -> {:?}",
s.wal_state()?
),
Err(e) => println!("G2 CoreWalSeam::open REFUSED: {e:#}"),
}
match CoreWalSeam::open_reader(&db) {
Ok(s) => println!(
"G2r CoreWalSeam::open_reader (OpenFlags::ReadOnly, no env var) SUCCEEDED alongside it -> {:?}",
s.wal_state()?
),
Err(e) => println!("G2r CoreWalSeam::open_reader REFUSED: {e:#}"),
}
match raw_consistent_copy_live(&db, PAGE_SIZE).await {
Ok(img) => println!("G3 raw_consistent_copy_live SUCCEEDED ({} bytes)", img.len()),
Err(e) => println!("G3 raw_consistent_copy_live REFUSED: {e:#}"),
}
writeln!(sin, "INSERT INTO t(id,v) VALUES(4242,'after-hold');\nSELECT 'PROBE-SENTINEL';")?;
sin.flush()?;
loop {
line.clear();
anyhow::ensure!(rd.read_line(&mut line)? > 0, "foreign holder died mid-write");
if line.contains("PROBE-SENTINEL") {
break;
}
}
let g3c = raw_consistent_copy_live(&db, PAGE_SIZE).await;
match &g3c {
Ok(img) => {
let out = p.path("g-readonly.db");
std::fs::write(&out, img)?;
println!(
"G3c with OpenFlags::ReadOnly (no env var), copy SUCCEEDED ({} bytes); integrity_check {:?}; rows {:?} (source has 4, incl. the row written while we read)",
img.len(),
sqlite3_clean(&out, "PRAGMA integrity_check;"),
row_count(&out)
);
}
Err(e) => println!("G3c with OpenFlags::ReadOnly (no env var), copy REFUSED: {e:#}"),
}
std::env::set_var("LIMBO_DISABLE_FILE_LOCK", "1");
match raw_consistent_copy_live(&db, PAGE_SIZE).await {
Ok(img) => {
let out = p.path("g-nolock.db");
std::fs::write(&out, &img)?;
println!(
"G3b with LIMBO_DISABLE_FILE_LOCK=1, copy SUCCEEDED ({} bytes); integrity_check {:?}; rows {:?} (source has 4, incl. the row written while we read)",
img.len(),
sqlite3_clean(&out, "PRAGMA integrity_check;"),
row_count(&out)
);
}
Err(e) => println!("G3b with LIMBO_DISABLE_FILE_LOCK=1, copy still REFUSED: {e:#}"),
}
std::env::remove_var("LIMBO_DISABLE_FILE_LOCK");
if let Ok(img) = &g3c {
let tgt = target();
let base = upload_base_snapshot(&tgt, img).await?;
match CoreWalSeam::open_reader(&db) {
Ok(seam) => match tail_frames(&seam, &tgt, &cfg(&base)).await {
Ok(o) => println!(
"G6 tail_frames through a ReadOnly seam, foreign holder STILL attached: {}",
short(&o)
),
Err(e) => println!("G6 tail_frames through a ReadOnly seam FAILED: {e:#}"),
},
Err(e) => println!("G6 could not open a ReadOnly seam to tail: {e:#}"),
}
}
let _ = Command::new("kill").arg("-9").arg(child.id().to_string()).status();
let _ = child.wait();
drop(sin);
println!("G4 foreign holder killed; retrying with nobody else on the file:");
match CoreWalSeam::open(&db) {
Ok(s) => println!("G5 CoreWalSeam::open SUCCEEDED -> {:?}", s.wal_state()?),
Err(e) => println!("G5 CoreWalSeam::open still REFUSED: {e:#}"),
}
println!(
"G VERDICT: A WRITABLE turso open and upstream C SQLite are mutually exclusive on one file — \
turso_core takes a whole-file exclusive fcntl lock at open (io/unix.rs lock_file(true)), \
so whichever opens first locks the other out, and headscale holds its connection for its \
whole process lifetime (G2). That is a GUARD, not a format incompatibility, and R858-B18 \
resolved it AT THE HANDLE: OpenFlags::ReadOnly skips the lock for that one handle with no \
env var and no process-wide effect (G2r), and the resulting copy is byte-for-byte as good \
as the env-var one — G3c and G3b agree exactly, same size, integrity_check ok, same 4 rows \
including the one the holder wrote after we had already been refused. So LIMBO_DISABLE_FILE_LOCK \
is never needed and should never be used: it removes the guard for EVERY open in the \
process, including the writable ones. THE FLAG IS ONLY HALF: neither engine observes the \
other's locks (turso whole-file fcntl, C SQLite byte-range + the -shm WAL index), so \
getting in without a lock is not coordination — see probe H for the optimistic validation \
that turns 'may read torn state' into 'detects torn state and refuses', and probe S for \
the same two-part fix on the snapshot tier."
);
Ok(())
}
async fn probe_h() -> Result<()> {
println!("\n=== PROBE H: validation under a foreign writer running CONCURRENTLY with the copy ===");
let hot = h_regime(
"HOT",
HRegime { seed_rows: 4000, autockpt: 16, batch: 40, pad: 400, pause_ms: 0, rounds: 12 },
)
.await?;
let calm = h_regime(
"CALM",
HRegime { seed_rows: 200, autockpt: 1000, batch: 1, pad: 60, pause_ms: 15, rounds: 12 },
)
.await?;
println!(
"H VERDICT: {}",
if hot.corrupt > 0 || calm.corrupt > 0 {
"WRONG — an image was accepted that should not have been; the validation is not sufficient."
} else if hot.refused == 0 {
"UNPROVEN (refuse half) — no refusal was observed even in the HOT regime, so this run did not exercise the detector."
} else if calm.accepted == 0 {
"UNPROVEN (accept half) — the detector refuses, but even the CALM regime never produced a validated image, so the protocol may be too strict to back up a live database at all."
} else {
"CORRECT — under concurrent foreign writes the copy either returns a validated point-in-time image or REFUSES, and never returned a torn image. The CALM (headscale-shaped) regime backs up successfully WHILE the foreign process writes; the HOT regime is correctly reported as unbackupable rather than silently mis-copied."
}
);
Ok(())
}
struct HRegime {
seed_rows: i64,
autockpt: i64,
batch: i64,
pad: usize,
pause_ms: u64,
rounds: u32,
}
struct HResult {
accepted: u32,
refused: u32,
corrupt: u32,
}
async fn h_regime(tag: &str, r: HRegime) -> Result<HResult> {
let p = Probe::new(&format!("h-{}", tag.to_lowercase()))?;
let db = p.path("h.db");
seed_clean(&db, r.seed_rows)?;
println!(
"\nH0[{tag}] seeded {} rows: main {}B | {} | writer autocheckpoint {} pages, {}-row commits every {}ms",
r.seed_rows,
file_len(&db),
hdr_str(&db),
r.autockpt,
r.batch,
r.pause_ms
);
let mut child = Command::new(sqlite3_bin())
.arg(&db)
.stdin(Stdio::piped())
.stdout(Stdio::null())
.stderr(Stdio::piped())
.spawn()?;
let mut sin = child.stdin.take().expect("piped stdin");
let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
let feeder_stop = stop.clone();
let (autockpt, batch, pad, pause_ms) = (r.autockpt, r.batch, r.pad, r.pause_ms);
let feeder = std::thread::spawn(move || {
let mut i = 100_000i64;
let _ = writeln!(sin, "PRAGMA wal_autocheckpoint={autockpt};");
while !feeder_stop.load(std::sync::atomic::Ordering::Relaxed) {
let mut sql = String::from("BEGIN;\n");
for _ in 0..batch {
sql.push_str(&format!(
"INSERT INTO t(id,v) VALUES({i},'{}');\n",
"y".repeat(pad)
));
i += 1;
}
sql.push_str("COMMIT;\n");
if writeln!(sin, "{sql}").is_err() {
break;
}
let _ = sin.flush();
if pause_ms > 0 {
std::thread::sleep(Duration::from_millis(pause_ms));
}
}
drop(sin);
});
let mut res = HResult { accepted: 0, refused: 0, corrupt: 0 };
let mut strict_would_refuse = 0u32;
let mut first_refusal: Option<String> = None;
let mut last_rows: i64 = 0;
for round in 0..r.rounds {
let before = SourceFingerprint::read(&db)?;
let started = Instant::now();
let got = raw_consistent_copy_live(&db, PAGE_SIZE).await;
let elapsed = started.elapsed();
let after = SourceFingerprint::read(&db)?;
if before.wal_len != after.wal_len {
strict_would_refuse += 1;
}
match got {
Ok(img) => {
res.accepted += 1;
let out = p.path(&format!("h-img-{round}.db"));
std::fs::write(&out, &img)?;
let integrity = sqlite3_clean(&out, "PRAGMA integrity_check;");
let rows = row_count(&out);
let ok = matches!(integrity.as_deref(), Ok("ok"))
&& rows.as_ref().is_ok_and(|n| *n >= last_rows);
if ok {
last_rows = *rows.as_ref().unwrap();
} else {
res.corrupt += 1;
println!(
"H1[{tag}/{round}] ACCEPTED a BAD image after {elapsed:?}: {} bytes; integrity {:?}; rows {:?} (previous accepted image had {last_rows})",
img.len(),
integrity,
rows
);
}
let _ = std::fs::remove_file(&out);
}
Err(e) => {
res.refused += 1;
if first_refusal.is_none() {
first_refusal = Some(format!("{e:#}"));
}
}
}
}
stop.store(true, std::sync::atomic::Ordering::Relaxed);
let _ = Command::new("kill").arg("-9").arg(child.id().to_string()).status();
let _ = child.wait();
let _ = feeder.join();
println!(
"H2[{tag}] {} copies under a concurrent foreign writer: {} accepted, {} refused-and-validated, {} accepted-but-corrupt",
r.rounds, res.accepted, res.refused, res.corrupt
);
println!(
"H3[{tag}] every accepted image passed PRAGMA integrity_check and was row-count-monotonic: {} (last accepted image had {last_rows} rows)",
res.corrupt == 0
);
match &first_refusal {
Some(m) => println!("H4[{tag}] first refusal message, verbatim: {m}"),
None => println!("H4[{tag}] no refusal was observed in {} rounds", r.rounds),
}
println!(
"H5[{tag}] the STRICTER rule this crate does not use (wal_len counts as movement) would have refused {strict_would_refuse}/{} of the same copies",
r.rounds
);
Ok(res)
}
async fn probe_s() -> Result<()> {
println!("\n=== PROBE S: the snapshot tier against a live foreign holder ===");
let p = Probe::new("s")?;
let db = p.path("s.db");
seed_clean(&db, 50)?;
let mut child = Command::new(sqlite3_bin())
.arg(&db)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()?;
let mut sin = child.stdin.take().expect("piped stdin");
let sout = child.stdout.take().expect("piped stdout");
let mut rd = BufReader::new(sout);
writeln!(sin, "SELECT count(*) FROM t;\nSELECT 'PROBE-SENTINEL';")?;
sin.flush()?;
let mut line = String::new();
loop {
line.clear();
anyhow::ensure!(rd.read_line(&mut line)? > 0, "foreign holder died early");
if line.contains("PROBE-SENTINEL") {
break;
}
}
println!("S0 foreign sqlite3 connection is open and idle on the DB");
let tgt = target();
match snapshot_and_upload(&db, &tgt).await {
Ok(o) => println!("S1 snapshot_and_upload SUCCEEDED -> {o:?}"),
Err(e) => println!("S1 snapshot_and_upload REFUSED: {e:#}"),
}
for (tag, flags) in [
("S1w VACUUM INTO via write flags (what snapshot.rs did before)", turso_core::OpenFlags::default()),
("S2 VACUUM INTO via OpenFlags::ReadOnly", turso_core::OpenFlags::ReadOnly),
] {
let vac = p.path(&format!("s-vacuum-{:?}.db", flags));
let _ = std::fs::remove_file(&vac);
let got = (|| -> anyhow::Result<()> {
let io: Arc<dyn turso_core::IO> = Arc::new(turso_core::PlatformIO::new()?);
let core = turso_core::Database::open_file_with_flags(
io,
&db,
flags,
turso_core::DatabaseOpts::new(),
None,
)?;
let conn = core.connect()?;
conn.execute(format!("VACUUM INTO '{}'", vac.replace('\'', "''")))?;
Ok(())
})();
match got {
Ok(()) => println!(
"{tag} SUCCEEDED -> {}B; integrity_check {:?}; rows {:?}",
file_len(&vac),
sqlite3_clean(&vac, "PRAGMA integrity_check;"),
row_count(&vac)
),
Err(e) => println!("{tag} REFUSED: {e:#}"),
}
}
match raw_consistent_copy_live(&db, PAGE_SIZE).await {
Ok(img) => {
let key = upload_base_snapshot(&tgt, &img).await?;
let out = p.path("s-reader.db");
std::fs::write(&out, &img)?;
println!(
"S3 raw_consistent_copy_live + upload_base_snapshot SUCCEEDED -> {key} ({}B); integrity_check {:?}; rows {:?}",
img.len(),
sqlite3_clean(&out, "PRAGMA integrity_check;"),
row_count(&out)
);
}
Err(e) => println!("S3 reader path REFUSED: {e:#}"),
}
println!(
"S VERDICT: tier 1a was blocked for exactly the same reason as tier 2 and is fixed the same \
way — S1w (write flags) is still refused while S2 (OpenFlags::ReadOnly) vacuums the same \
live source successfully, so VACUUM INTO does run on a read-only connection and tier 1a \
keeps its byte-deterministic gate-2 hash. snapshot.rs now opens through turso_core with \
that flag instead of turso::Builder (which exposes no OpenFlags), wrapped in the same \
optimistic validation, which is why S1 succeeds where it used to report a locking error."
);
let _ = Command::new("kill").arg("-9").arg(child.id().to_string()).status();
let _ = child.wait();
drop(sin);
Ok(())
}
fn probe_e() -> Result<()> {
println!("\n=== PROBE E: default-autocheckpoint fold rate (how often the restart path fires) ===");
let p = Probe::new("e")?;
let db = p.path("e.db");
println!(
"E0 writer default wal_autocheckpoint = {} pages",
sqlite3_clean(&db, "PRAGMA journal_mode=WAL; PRAGMA wal_autocheckpoint;")?
.lines()
.last()
.unwrap_or("?")
);
sqlite3_clean(
&db,
"CREATE TABLE IF NOT EXISTS t(id INTEGER PRIMARY KEY, v TEXT); CREATE INDEX IF NOT EXISTS ix ON t(v);",
)?;
let mut child = Command::new(sqlite3_bin())
.arg(&db)
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()?;
let mut sin = child.stdin.take().expect("piped stdin");
let sout = child.stdout.take().expect("piped stdout");
let mut rd = BufReader::new(sout);
let mut line = String::new();
let total = 3000i64;
let mut hdr = wal_hdr(&db);
let mut folds: Vec<i64> = Vec::new();
for i in 0..total {
writeln!(sin, "INSERT INTO t(id,v) VALUES({i},'{}');\nSELECT 'PROBE-SENTINEL';", "x".repeat(200))?;
sin.flush()?;
loop {
line.clear();
anyhow::ensure!(rd.read_line(&mut line)? > 0, "writer died at row {i}");
if line.contains("PROBE-SENTINEL") {
break;
}
}
let h = wal_hdr(&db);
if h != hdr {
if hdr.is_some() {
folds.push(i);
}
hdr = h;
}
if i % 500 == 0 {
println!(
"E1 after {i} single-row txns: wal {}B main {}B | {}",
file_len(&format!("{db}-wal")),
file_len(&db),
hdr_str(&db),
);
}
}
println!(
"E2 after {total} single-row txns: wal {}B main {}B | {}",
file_len(&format!("{db}-wal")),
file_len(&db),
hdr_str(&db),
);
println!("E3 WAL-header changes (folds/restarts) at txn #: {folds:?}");
let _ = Command::new("kill").arg("-9").arg(child.id().to_string()).status();
let _ = child.wait();
drop(sin);
let w = turso_watermark(&db)?;
println!(
"E4 with the writer gone, turso reports (seq {}, max_frame {})",
w.checkpoint_seq, w.last_frame
);
println!(
"E VERDICT: {} fold(s) in {total} single-row transactions -> one restart per ~{} writes; turso-observed checkpoint_seq after all of them = {}",
folds.len(),
if folds.is_empty() { "never".to_string() } else { (total / folds.len() as i64).to_string() },
w.checkpoint_seq
);
Ok(())
}
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<()> {
let want: Vec<String> = std::env::args().skip(1).map(|s| s.to_lowercase()).collect();
let on = |c: &str| want.is_empty() || want.iter().any(|w| w == c);
println!(
"foreign writer: {} -> sqlite {}",
sqlite3_bin(),
sqlite3_clean(":memory:", "SELECT sqlite_version();")?
);
if on("a") {
probe_a()?;
}
if on("b") {
println!("\n=== PROBE B: fold, then a SHORTER new WAL (post-fold frames < watermark) ===");
tail_across_fold("b", 6, 3, 3, 1, false).await?;
}
if on("f") {
println!("\n=== PROBE F: fold, then a LONGER new WAL (post-fold frames > watermark) ===");
tail_across_fold("f", 40, 20, 60, 30, true).await?;
}
if on("c") {
probe_c().await?;
}
if on("d") {
probe_d()?;
}
if on("g") {
probe_g().await?;
}
if on("h") {
probe_h().await?;
}
if on("s") {
probe_s().await?;
}
if on("e") {
probe_e()?;
}
Ok(())
}