use std::future::Future;
use std::io::{self, Read, Seek, SeekFrom, Write};
use std::ops::Range;
use std::path::{Path, PathBuf};
use crate::wal_fec::replay::recover_wal_fec_image_with_certificates;
use crate::wal_index::{
WAL_INDEX_VERSION, WAL_SHM_SEGMENT_BYTES, WalIndexFrameLocation, WalIndexHdr,
append_native_wal_index_entry, invalidate_shared_wal_index_header,
publish_shared_wal_index_header, read_shared_wal_index_header, replace_shared_wal_index_region,
reset_shared_wal_index_recovery_marks, validate_shared_wal_index_wal_binding,
};
use crate::{WAL_FRAME_HEADER_SIZE, WAL_HEADER_SIZE, WalFrameHeader};
use asupersync::runtime::spawn_blocking;
use fsqlite_error::{FrankenError, Result};
use fsqlite_types::cx::Cx;
use fsqlite_types::flags::VfsOpenFlags;
use fsqlite_vfs::{FileIdentity, ShmRegion, Vfs, VfsFile, host_fs};
use super::{
CapturedSource, ExportReport, IO_CHUNK, NativeFile, NativeVfs, Options, Snapshot, SourceFile,
capture_held, checkpoint, companion, refuse_destination_artifacts, verify_image,
};
struct RepairHandoff {
source: PathBuf,
identity: FileIdentity,
main: SourceFile<NativeFile>,
report: ExportReport,
}
impl RepairHandoff {
#[allow(clippy::future_not_send)]
async fn open<T, F, Fut>(self, cx: &Cx, opener: F) -> (Result<T>, ExportReport)
where
F: FnOnce(PathBuf, FileIdentity) -> Fut,
Fut: Future<Output = Result<T>>,
{
let Self {
source,
identity,
main,
report,
} = self;
let opened = match checkpoint(cx) {
Ok(()) => opener(source, identity).await,
Err(error) => Err(error),
};
drop(main);
(opened, report)
}
}
impl Options {
#[allow(clippy::future_not_send)]
pub async fn repair_and_open<T, F, Fut>(
&self,
cx: &Cx,
opener: F,
) -> Result<(Result<T>, ExportReport)>
where
F: FnOnce(PathBuf, FileIdentity) -> Fut,
Fut: Future<Output = Result<T>>,
{
super::preflight(cx, self)?;
let handoff = run_for_open(&NativeVfs::new(), cx, self).await?;
Ok(handoff.open(cx, opener).await)
}
}
struct RepairPlan {
target: Vec<u8>,
changed: Vec<Range<usize>>,
index_regions: Vec<Vec<u8>>,
index_header: WalIndexHdr,
pages: usize,
certificate_anchors: usize,
}
fn corruption(detail: impl Into<String>) -> FrankenError {
FrankenError::WalCorrupt {
detail: detail.into(),
}
}
impl RepairPlan {
fn build(snapshot: &Snapshot, options: &Options) -> Result<Self> {
let mut database_file_id = [0; 16];
database_file_id.copy_from_slice(
snapshot
.database
.get(76..92)
.ok_or_else(|| corruption("missing captured database identity"))?,
);
let replay = recover_wal_fec_image_with_certificates(
&snapshot.wal,
&snapshot.sidecar,
&snapshot.certificates,
database_file_id,
options.replay,
)?;
let database = replay.database_image(&snapshot.database, options.max_database_bytes)?;
let page_size =
usize::try_from(replay.header().page_size).map_err(|_| FrankenError::TooBig)?;
let pages = database.len() / page_size;
drop(database);
let frame_size = page_size
.checked_add(WAL_FRAME_HEADER_SIZE)
.ok_or(FrankenError::TooBig)?;
let prefix = replay.replayable_prefix();
if prefix.get(..WAL_HEADER_SIZE) != snapshot.wal.get(..WAL_HEADER_SIZE) {
return Err(corruption("repair cannot replace a WAL generation"));
}
let target_len = snapshot.wal.len().max(prefix.len());
if target_len - snapshot.wal.len() != replay.restored_tail_bytes()
|| replay.restored_tail_bytes() > page_size
{
return Err(corruption(
"repair growth lacks verified terminal payload provenance",
));
}
let mut target = Vec::new();
target
.try_reserve_exact(target_len)
.map_err(|_| FrankenError::OutOfMemory)?;
target.extend_from_slice(&snapshot.wal);
target.resize(target_len, 0);
target[..prefix.len()].copy_from_slice(prefix);
let frame_count = replay.committed_frames();
let mut changed = Vec::new();
changed
.try_reserve(usize::try_from(frame_count).map_err(|_| FrankenError::TooBig)?)
.map_err(|_| FrankenError::OutOfMemory)?;
for offset in (WAL_HEADER_SIZE..prefix.len()).step_by(frame_size) {
let range = offset..offset + frame_size;
if snapshot.wal.get(range.clone()) != Some(&target[range.clone()]) {
changed.push(range);
}
}
let last_region = if frame_count == 0 {
0
} else {
WalIndexFrameLocation::new(frame_count)?.region
};
let mut index_regions = Vec::new();
for _ in 0..=last_region {
let mut region = Vec::new();
region
.try_reserve_exact(WAL_SHM_SEGMENT_BYTES)
.map_err(|_| FrankenError::OutOfMemory)?;
region.resize(WAL_SHM_SEGMENT_BYTES, 0);
index_regions.push(region);
}
let mut terminal = None;
for (index, frame) in prefix[WAL_HEADER_SIZE..]
.chunks_exact(frame_size)
.enumerate()
{
let number = u32::try_from(index)
.ok()
.and_then(|index| index.checked_add(1))
.ok_or(FrankenError::TooBig)?;
let marker = WalFrameHeader::from_bytes(frame)?;
let region = usize::try_from(WalIndexFrameLocation::new(number)?.region)
.map_err(|_| FrankenError::TooBig)?;
append_native_wal_index_entry(&mut index_regions[region], number, marker.page_number)?;
terminal = Some((number, marker));
}
let wal_header = replay.header();
let mut index_header = WalIndexHdr {
i_version: WAL_INDEX_VERSION,
unused: 0,
i_change: 0,
is_init: 1,
big_end_cksum: u8::from(wal_header.big_endian_checksum()),
sz_page: if page_size == 65_536 {
1
} else {
u16::try_from(page_size).map_err(|_| FrankenError::TooBig)?
},
mx_frame: frame_count,
n_page: terminal.map_or(0, |(_, marker)| marker.db_size),
a_frame_cksum: terminal.map_or([0, 0], |(_, marker)| {
[marker.checksum.s1, marker.checksum.s2]
}),
a_salt: [wal_header.salts.salt1, wal_header.salts.salt2],
a_cksum: [0, 0],
};
index_header.update_checksum()?;
validate_shared_wal_index_wal_binding(&index_header, wal_header, terminal)?;
Ok(Self {
target,
changed,
index_regions,
index_header,
pages,
certificate_anchors: replay.certificate_anchors().len(),
})
}
}
struct IndexPublication {
zero: ShmRegion,
complete: bool,
}
impl Drop for IndexPublication {
fn drop(&mut self) {
if !self.complete
&& let Err(error) = invalidate_shared_wal_index_header(&self.zero)
{
eprintln!("failed to invalidate incomplete recovery index: {error}");
}
}
}
trait RepairFile: Read + Write + Seek {
fn truncate(&mut self, len: u64) -> io::Result<()>;
}
impl RepairFile for std::fs::File {
fn truncate(&mut self, len: u64) -> io::Result<()> {
self.set_len(len)
}
}
fn write_ranges(
file: &mut (impl Write + Seek),
image: &[u8],
ranges: impl IntoIterator<Item = Range<usize>>,
) -> io::Result<()> {
for range in ranges {
let bytes = image
.get(range.clone())
.ok_or_else(|| io::Error::other("WAL repair range exceeds its source image"))?;
file.seek(SeekFrom::Start(
u64::try_from(range.start)
.map_err(|_| io::Error::other("WAL repair offset overflow"))?,
))?;
file.write_all(bytes)?;
}
Ok(())
}
fn settle_writes<W: RepairFile>(
file: &mut W,
cx: &Cx,
original: &[u8],
plan: &RepairPlan,
mut sync: impl FnMut(&mut W) -> io::Result<()>,
) -> Result<blake3::Hash> {
let write = (|| {
write_ranges(file, &plan.target, plan.changed.iter().cloned())?;
sync(file)?;
verify_image(file, cx, &plan.target)
})();
match write {
Ok(digest) => Ok(digest),
Err(error) => {
let restore = (|| {
let ranges = plan.changed.iter().filter_map(|range| {
let end = range.end.min(original.len());
(range.start < end).then_some(range.start..end)
});
write_ranges(file, original, ranges)?;
if plan.target.len() > original.len() {
file.truncate(
u64::try_from(original.len()).map_err(|_| FrankenError::TooBig)?,
)?;
}
sync(file)?;
verify_image(file, cx, original)
})();
match restore {
Ok(_) => Err(corruption(format!(
"WAL repair failed; exact original WAL restored and synced: {error}"
))),
Err(restore_error) => Err(corruption(format!(
"WAL repair outcome is indeterminate; retain the original backup and do not use the database: repair={error}; restore={restore_error}"
))),
}
}
}
}
fn recheck_names(vfs: &NativeVfs, cx: &Cx, source: &Path, captured: &CapturedSource) -> Result<()> {
let wal_path = companion(source, "-wal");
for (path, identity, flags) in [
(
source,
captured.main_identity,
VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB,
),
(
wal_path.as_path(),
captured.wal_identity,
VfsOpenFlags::READONLY | VfsOpenFlags::WAL,
),
] {
let (probe, _) = vfs.open_with_expected_identity(cx, path, flags, identity)?;
SourceFile::new(probe, cx).finish()?;
}
Ok(())
}
fn backup_original(
vfs: &NativeVfs,
cx: &Cx,
path: &Path,
original: &[u8],
mut sync: impl FnMut(&mut std::fs::File) -> io::Result<()>,
) -> Result<(std::fs::File, FileIdentity)> {
checkpoint(cx)?;
refuse_destination_artifacts(vfs, cx, path)?;
let mut backup = host_fs::reserve_new_file(path)?;
let result: Result<FileIdentity> = (|| {
let identity = FileIdentity::from_file(&backup)?.ok_or(FrankenError::Unsupported)?;
let validate_destination = || {
host_fs::validate_reserved_file_identity(path, identity)?;
refuse_destination_artifacts(vfs, cx, path)
};
validate_destination()?;
for chunk in original.chunks(IO_CHUNK) {
checkpoint(cx)?;
backup.write_all(chunk)?;
}
sync(&mut backup)?;
verify_image(&mut backup, cx, original)?;
validate_destination()?;
vfs.sync_parent_directory(cx, path)?;
validate_destination()?;
Ok(identity)
})();
match result {
Ok(identity) => Ok((backup, identity)),
Err(error) => {
eprintln!(
"Backup is NOT certified for {}; source WAL has not been modified and no output or replacement was removed",
path.display()
);
Err(error)
}
}
}
fn repair_captured(
cx: &Cx,
options: &Options,
mut captured: CapturedSource,
) -> Result<RepairHandoff> {
checkpoint(cx)?;
let vfs = NativeVfs::new();
let mut plan = RepairPlan::build(&captured.snapshot, options)?;
let wal_path = companion(&options.source, "-wal");
recheck_names(&vfs, cx, &options.source, &captured)?;
let mut wal = host_fs::open_wal_for_guarded_repair(&wal_path, captured.wal_identity)?;
verify_image(&mut wal, cx, &captured.snapshot.wal)?;
let region_size = u32::try_from(WAL_SHM_SEGMENT_BYTES).map_err(|_| FrankenError::TooBig)?;
let mut regions = Vec::new();
for number in 0..plan.index_regions.len() {
regions.push(captured.main.file.shm_map(
cx,
u32::try_from(number).map_err(|_| FrankenError::TooBig)?,
region_size,
true,
)?);
}
if let Ok(Some(previous)) = read_shared_wal_index_header(®ions[0]) {
plan.index_header.i_change = previous.i_change.wrapping_add(1);
plan.index_header.update_checksum()?;
}
let (backup, backup_identity) = backup_original(
&vfs,
cx,
&options.destination,
&captured.snapshot.wal,
|file| file.sync_all(),
)?;
recheck_names(&vfs, cx, &options.source, &captured)?;
verify_image(&mut wal, cx, &captured.snapshot.wal)?;
host_fs::validate_reserved_file_identity(&options.destination, backup_identity)?;
refuse_destination_artifacts(&vfs, cx, &options.destination)?;
checkpoint(cx)?;
let _mask = cx.masked();
invalidate_shared_wal_index_header(®ions[0])?;
let mut publication = IndexPublication {
zero: regions[0].share(),
complete: false,
};
let result = (|| {
let digest = settle_writes(&mut wal, cx, &captured.snapshot.wal, &plan, |file| {
file.sync_all()
})?;
for (number, (region, bytes)) in regions.iter().zip(&plan.index_regions).enumerate() {
replace_shared_wal_index_region(
region,
u32::try_from(number).map_err(|_| FrankenError::TooBig)?,
bytes,
)?;
}
reset_shared_wal_index_recovery_marks(®ions[0], plan.index_header.mx_frame)?;
captured.main.file.shm_barrier();
publish_shared_wal_index_header(®ions[0], &plan.index_header)?;
if read_shared_wal_index_header(®ions[0])? != Some(plan.index_header) {
return Err(corruption("repaired WAL index failed publication readback"));
}
host_fs::validate_reserved_file_identity(&options.destination, backup_identity)?;
refuse_destination_artifacts(&vfs, cx, &options.destination)?;
publication.complete = true;
Ok::<_, FrankenError>(digest)
})();
drop(publication);
drop(regions);
drop(wal);
if result.is_err() {
eprintln!(
"WAL repair is NOT certified; verify the original backup at {}",
options.destination.display()
);
}
let digest = result?;
captured.wal.finish()?;
{
let _cleanup_mask = captured.main.cleanup_cx.masked();
captured
.main
.file
.restore_external_maintenance_attempt(&captured.main.cleanup_cx)?;
captured.main.maintenance = false;
}
drop(backup);
let report = ExportReport {
destination: options.destination.clone(),
pages: plan.pages,
wal_frames: plan.index_header.mx_frame,
repaired_frames: plan.changed.len(),
certificate_anchors: plan.certificate_anchors,
digest,
repaired_in_place: true,
};
Ok(RepairHandoff {
source: options.source.clone(),
identity: captured.main_identity,
main: captured.main,
report,
})
}
pub(super) async fn run(vfs: &NativeVfs, cx: &Cx, options: &Options) -> Result<ExportReport> {
let mut handoff = run_for_open(vfs, cx, options).await?;
handoff.main.finish()?;
Ok(handoff.report)
}
async fn run_for_open(vfs: &NativeVfs, cx: &Cx, options: &Options) -> Result<RepairHandoff> {
let source = vfs.full_pathname(cx, &options.source)?;
let destination = vfs.full_pathname(cx, &options.destination)?;
if [
"",
"-wal",
"-shm",
"-journal",
"-wal-fec",
"-wal-fec.lock",
"-wal-cert",
".fsqlite-shm",
]
.iter()
.any(|suffix| destination == companion(&source, suffix))
|| vfs.path_entry_exists(cx, &destination)?
{
return Err(FrankenError::CannotOpen { path: destination });
}
refuse_destination_artifacts(vfs, cx, &destination)?;
let options = Options {
source,
destination,
..*options
};
let captured = capture_held(vfs, cx, &options).await?;
let worker_cx = cx.create_child_for_spawn();
spawn_blocking(move || repair_captured(&worker_cx, &options, captured)).await
}
#[cfg(test)]
mod tests {
use super::*;
use std::cell::Cell;
use std::io::Cursor;
use std::rc::Rc;
use std::task::{Context, Poll, Waker};
impl RepairFile for Cursor<Vec<u8>> {
fn truncate(&mut self, len: u64) -> io::Result<()> {
let len =
usize::try_from(len).map_err(|_| io::Error::other("test file size overflow"))?;
self.get_mut().truncate(len);
Ok(())
}
}
fn with_runtime<F: Future>(future: F) -> F::Output {
asupersync::runtime::RuntimeBuilder::current_thread()
.blocking_threads(1, 2)
.build()
.unwrap()
.block_on(future)
}
fn attached_context() -> Cx {
let cx = Cx::new();
cx.set_native_cx(asupersync::Cx::current().unwrap());
cx
}
#[test]
fn backup_retains_the_verified_descriptor_identity() {
with_runtime(async {
let directory = tempfile::tempdir().unwrap().keep();
let path = directory.join("original.wal");
let original = vec![0x37; IO_CHUNK + 513];
let cx = attached_context();
let (mut backup, identity) =
backup_original(&NativeVfs::new(), &cx, &path, &original, |file| {
file.sync_all()
})
.unwrap();
assert_eq!(FileIdentity::from_file(&backup).unwrap(), Some(identity));
host_fs::validate_reserved_file_identity(&path, identity).unwrap();
assert_eq!(
verify_image(&mut backup, &cx, &original).unwrap(),
blake3::hash(&original)
);
assert_eq!(host_fs::read(&path).unwrap(), original);
});
}
#[test]
fn backup_refuses_missing_or_replaced_path_during_sync() {
with_runtime(async {
for install_replacement in [false, true] {
let directory = tempfile::tempdir().unwrap().keep();
let path = directory.join("original.wal");
let retained = directory.join("retained.wal");
let original = vec![0x37; IO_CHUNK + 513];
let cx = attached_context();
let result = backup_original(&NativeVfs::new(), &cx, &path, &original, |file| {
file.sync_all()?;
host_fs::rename(&path, &retained).unwrap();
if install_replacement {
host_fs::write(&path, &original).unwrap();
}
Ok(())
});
if install_replacement {
assert!(matches!(result, Err(FrankenError::BusyRecovery)));
assert_eq!(host_fs::read(&path).unwrap(), original);
} else {
assert!(matches!(result, Err(FrankenError::Io(error))
if error.kind() == io::ErrorKind::NotFound));
assert!(!path.exists());
}
assert_eq!(host_fs::read(&retained).unwrap(), original);
}
});
}
#[test]
fn backup_refuses_recovery_companions_created_during_sync() {
with_runtime(async {
for suffix in crate::native_recovery::RECOVERY_COMPANION_SUFFIXES {
let directory = tempfile::tempdir().unwrap().keep();
let path = directory.join("original.wal");
let artifact = companion(&path, suffix);
let original = vec![0x37; IO_CHUNK + 513];
let cx = attached_context();
let result = backup_original(&NativeVfs::new(), &cx, &path, &original, |file| {
file.sync_all()?;
host_fs::write(&artifact, b"unowned companion").unwrap();
Ok(())
});
assert!(matches!(result, Err(FrankenError::CannotOpen { path })
if path == artifact));
assert_eq!(host_fs::read(&path).unwrap(), original);
assert_eq!(host_fs::read(&artifact).unwrap(), b"unowned companion");
}
});
}
#[test]
fn backup_sync_and_readback_failures_never_return_a_verified_owner() {
with_runtime(async {
for corrupt_readback in [false, true] {
let directory = tempfile::tempdir().unwrap().keep();
let path = directory.join("original.wal");
let original = vec![0x37; IO_CHUNK + 513];
let cx = attached_context();
let result = backup_original(&NativeVfs::new(), &cx, &path, &original, |file| {
if !corrupt_readback {
return Err(io::Error::other("injected backup sync failure"));
}
file.seek(SeekFrom::Start(0))?;
file.write_all(&[0xff])?;
file.sync_all()
});
if corrupt_readback {
assert!(matches!(result, Err(FrankenError::DatabaseCorrupt { .. })));
assert_eq!(host_fs::read(&path).unwrap()[0], 0xff);
} else {
assert!(matches!(result, Err(FrankenError::Io(_))));
assert_eq!(host_fs::read(&path).unwrap(), original);
}
}
});
}
fn handoff_fixture() -> (Options, Vec<u8>, Vec<u8>) {
use crate::checksum::{SqliteWalChecksum, WalChecksumTransform, WalHeader, WalSalts};
use crate::wal_fec::{
WalFecGroupMeta, WalFecGroupMetaInit, WalFecGroupRecord, append_wal_fec_group,
build_source_page_hashes, generate_wal_fec_repair_symbols,
};
use fsqlite_types::{ObjectId, Oti};
let directory = tempfile::tempdir().unwrap().keep();
let options = Options::new(directory.join("source.db"), directory.join("original.wal"));
let mut page = vec![0_u8; 512];
page[..16].copy_from_slice(b"SQLite format 3\0");
page[16..18].copy_from_slice(&512_u16.to_be_bytes());
page[18..20].copy_from_slice(&[2, 2]);
page[21..24].copy_from_slice(&[64, 32, 32]);
page[24..28].copy_from_slice(&7_u32.to_be_bytes());
page[28..32].copy_from_slice(&1_u32.to_be_bytes());
page[44..48].copy_from_slice(&4_u32.to_be_bytes());
page[56..60].copy_from_slice(&1_u32.to_be_bytes());
page[92..96].copy_from_slice(&7_u32.to_be_bytes());
page[100] = 13;
page[105..107].copy_from_slice(&512_u16.to_be_bytes());
host_fs::write(&options.source, &page).unwrap();
let pages: Vec<_> = (1_u32..=3)
.map(|version| {
let mut current = page.clone();
current[60..64].copy_from_slice(&version.to_be_bytes());
current
})
.collect();
let header = WalHeader {
magic: crate::WAL_MAGIC_LE,
format_version: crate::WAL_FORMAT_VERSION,
page_size: 512,
checkpoint_seq: 1,
salts: WalSalts {
salt1: 123,
salt2: 456,
},
checksum: SqliteWalChecksum::default(),
};
let mut wal = header.to_bytes().unwrap().to_vec();
let mut running = WalHeader::from_bytes(&wal).unwrap().checksum;
for (index, page) in pages.iter().enumerate() {
let start = wal.len();
wal.extend_from_slice(
&WalFrameHeader {
page_number: 1,
db_size: u32::from(index == 2),
salts: header.salts,
checksum: SqliteWalChecksum::default(),
}
.to_bytes(),
);
wal.extend_from_slice(page);
running = WalChecksumTransform::for_wal_frame(&wal[start..], 512, false)
.unwrap()
.apply(running);
wal[start + 16..start + 20].copy_from_slice(&running.s1.to_be_bytes());
wal[start + 20..start + 24].copy_from_slice(&running.s2.to_be_bytes());
}
let meta = WalFecGroupMeta::from_init(WalFecGroupMetaInit {
wal_salt1: 123,
wal_salt2: 456,
start_frame_no: 1,
end_frame_no: 3,
db_size_pages: 1,
page_size: 512,
k_source: 3,
r_repair: 8,
oti: Oti {
f: 1536,
al: 1,
t: 512,
z: 1,
n: 1,
},
object_id: ObjectId::derive_from_canonical_bytes(b"repair-open-handoff"),
page_numbers: vec![1; 3],
source_page_xxh3_128: build_source_page_hashes(&pages),
})
.unwrap();
let symbols = generate_wal_fec_repair_symbols(&meta, &pages).unwrap();
append_wal_fec_group(
&companion(&options.source, "-wal-fec"),
&WalFecGroupRecord::new(meta, symbols).unwrap(),
)
.unwrap();
let repaired = wal.clone();
wal[WAL_HEADER_SIZE + WAL_FRAME_HEADER_SIZE + 60] ^= 0xff;
host_fs::write(&companion(&options.source, "-wal"), &wal).unwrap();
(options, wal, repaired)
}
fn assert_recovery_available(source: &Path, cx: &Cx) {
let (file, _) = NativeVfs::new()
.open(
cx,
Some(source),
VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB,
)
.unwrap();
let mut owner = SourceFile::new(file, cx);
owner.acquire_recovery(cx).unwrap();
owner.finish().unwrap();
}
#[test]
fn handoff_opens_once_with_live_identity_after_repair_fences_are_restored() {
with_runtime(async {
let (options, original, repaired) = handoff_fixture();
let cx = attached_context();
let calls = Rc::new(Cell::new(0));
let calls_in_opener = Rc::clone(&calls);
let opener_cx = &cx;
let (opened, report) = options
.repair_and_open(&cx, |path, identity| async move {
calls_in_opener.set(calls_in_opener.get() + 1);
let (file, _) = NativeVfs::new().open_with_expected_identity(
opener_cx,
&path,
VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB,
identity,
)?;
let mut owner = SourceFile::new(file, opener_cx);
assert_eq!(owner.file.file_identity()?, Some(identity));
owner.acquire_recovery(opener_cx)?;
owner.finish()?;
Ok(calls_in_opener) })
.await
.unwrap();
assert!(Rc::ptr_eq(&opened.unwrap(), &calls));
assert_eq!(calls.get(), 1);
assert_eq!(report.wal_frames, 3);
assert_eq!(report.repaired_frames, 1);
assert!(report.repaired_in_place);
assert_eq!(report.digest, blake3::hash(&repaired));
assert_eq!(host_fs::read(&options.destination).unwrap(), original);
assert_eq!(
host_fs::read(&companion(&options.source, "-wal")).unwrap(),
repaired
);
});
}
#[test]
fn handoff_preserves_repair_receipt_when_opener_fails() {
with_runtime(async {
let (options, original, repaired) = handoff_fixture();
let cx = attached_context();
let (opened, report) = options
.repair_and_open(&cx, |_, _| async {
Err::<(), _>(FrankenError::NoSuchTable {
name: "opener sentinel".to_owned(),
})
})
.await
.unwrap();
assert!(
matches!(opened, Err(FrankenError::NoSuchTable { name }) if name == "opener sentinel")
);
assert_eq!(report.digest, blake3::hash(&repaired));
assert_eq!(host_fs::read(&options.destination).unwrap(), original);
assert_eq!(
host_fs::read(&companion(&options.source, "-wal")).unwrap(),
repaired
);
assert_recovery_available(&options.source, &cx);
});
}
#[test]
fn handoff_does_not_unlock_the_openers_native_claims() {
const CHILD_PATH: &str = "FSQLITE_REPAIR_OPEN_LOCK_CHILD";
const TEST: &str =
"native_recovery::repair::tests::handoff_does_not_unlock_the_openers_native_claims";
if let Some(path) = std::env::var_os(CHILD_PATH) {
with_runtime(async {
let cx = attached_context();
let (file, _) = NativeVfs::new()
.open(
&cx,
Some(Path::new(&path)),
VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB,
)
.unwrap();
let mut peer = SourceFile::new(file, &cx);
assert!(matches!(
peer.acquire_recovery(&cx),
Err(FrankenError::Busy)
));
peer.finish().unwrap();
});
return;
}
with_runtime(async {
let (options, _, _) = handoff_fixture();
let cx = attached_context();
let opener_cx = &cx;
let (opened, _) = options
.repair_and_open(&cx, |path, identity| async move {
let (file, _) = NativeVfs::new().open_with_expected_identity(
opener_cx,
&path,
VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB,
identity,
)?;
let mut new_owner = SourceFile::new(file, opener_cx);
new_owner.acquire_recovery(opener_cx)?;
Ok(new_owner)
})
.await
.unwrap();
let mut owner = opened.unwrap();
let mut child = std::process::Command::new(std::env::current_exe().unwrap())
.args([TEST, "--exact", "--nocapture"])
.env(CHILD_PATH, &options.source)
.spawn()
.unwrap();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15);
loop {
if let Some(status) = child.try_wait().unwrap() {
assert!(
status.success(),
"foreign process must observe the opener's locks"
);
break;
}
if std::time::Instant::now() >= deadline {
let _ = child.kill();
let _ = child.wait();
panic!("foreign lock witness did not terminate");
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
owner.finish().unwrap();
assert_recovery_available(&options.source, &cx);
});
}
#[test]
fn handoff_refusal_never_invokes_opener_or_overwrites_a_backup() {
with_runtime(async {
for existing_backup in [false, true] {
let (options, original, _) = handoff_fixture();
let cx = attached_context();
if existing_backup {
host_fs::write(&options.destination, b"existing backup").unwrap();
} else {
host_fs::write(&companion(&options.source, "-wal"), b"invalid WAL").unwrap();
}
let calls = Cell::new(0);
let result = options
.repair_and_open(&cx, |_, _| {
calls.set(calls.get() + 1);
std::future::ready(Ok(()))
})
.await;
assert!(result.is_err());
assert_eq!(calls.get(), 0);
if existing_backup {
assert_eq!(
host_fs::read(&options.destination).unwrap(),
b"existing backup"
);
assert_eq!(
host_fs::read(&companion(&options.source, "-wal")).unwrap(),
original
);
} else {
assert!(
!NativeVfs::new()
.path_entry_exists(&cx, &options.destination)
.unwrap()
);
}
}
});
}
#[test]
fn handoff_contention_never_invokes_opener_and_can_be_retried() {
with_runtime(async {
let (options, original, _) = handoff_fixture();
let cx = attached_context();
let vfs = NativeVfs::new();
let (file, _) = vfs
.open(
&cx,
Some(&options.source),
VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB,
)
.unwrap();
let mut owner = SourceFile::new(file, &cx);
owner.acquire_recovery(&cx).unwrap();
let calls = Cell::new(0);
let result = options
.repair_and_open(&cx, |_, _| {
calls.set(calls.get() + 1);
std::future::ready(Ok(()))
})
.await;
assert!(result.is_err());
assert_eq!(calls.get(), 0);
assert!(!vfs.path_entry_exists(&cx, &options.destination).unwrap());
assert_eq!(
host_fs::read(&companion(&options.source, "-wal")).unwrap(),
original
);
owner.finish().unwrap();
options
.repair_and_open(&cx, |_, _| std::future::ready(Ok(())))
.await
.unwrap()
.0
.unwrap();
});
}
#[test]
fn handoff_cancellation_after_repair_retains_receipt_without_starting_open() {
with_runtime(async {
let (options, original, repaired) = handoff_fixture();
let cx = attached_context();
let handoff = run_for_open(&NativeVfs::new(), &cx, &options)
.await
.unwrap();
cx.cancel();
let calls = Cell::new(0);
let (opened, report) = handoff
.open(&cx, |_, _| {
calls.set(calls.get() + 1);
std::future::ready(Ok(()))
})
.await;
assert!(matches!(opened, Err(FrankenError::Interrupt)));
assert_eq!(calls.get(), 0);
assert_eq!(report.digest, blake3::hash(&repaired));
assert_eq!(host_fs::read(&options.destination).unwrap(), original);
});
}
#[test]
fn handoff_drop_during_open_releases_managed_identity_guard() {
with_runtime(async {
let (options, original, repaired) = handoff_fixture();
let cx = attached_context();
let handoff = run_for_open(&NativeVfs::new(), &cx, &options)
.await
.unwrap();
let calls = Cell::new(0);
let mut future = Box::pin(handoff.open(&cx, |_, _| {
calls.set(calls.get() + 1);
std::future::pending::<Result<()>>()
}));
assert!(matches!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Pending
));
drop(future);
assert_eq!(calls.get(), 1);
assert_recovery_available(&options.source, &cx);
assert_eq!(host_fs::read(&options.destination).unwrap(), original);
assert_eq!(
host_fs::read(&companion(&options.source, "-wal")).unwrap(),
repaired
);
});
}
#[test]
fn handoff_unpolled_future_and_preflight_failure_have_no_effects() {
with_runtime(async {
let (options, original, _) = handoff_fixture();
let cx = attached_context();
let calls = Cell::new(0);
drop(options.repair_and_open(&cx, |_, _| {
calls.set(1);
std::future::ready(Ok(()))
}));
let detached = Cx::new();
assert!(
options
.repair_and_open(&detached, |_, _| {
calls.set(1);
std::future::ready(Ok(()))
})
.await
.is_err()
);
assert_eq!(calls.get(), 0);
assert!(
!NativeVfs::new()
.path_entry_exists(&cx, &options.destination)
.unwrap()
);
assert_eq!(
host_fs::read(&companion(&options.source, "-wal")).unwrap(),
original
);
});
}
#[test]
fn torn_tail_handoff_preserves_inode_and_backs_up_exact_original_eof() {
use std::os::unix::fs::MetadataExt;
with_runtime(async {
for missing in [1, 256, 512] {
let (options, mut original, repaired) = handoff_fixture();
original.truncate(original.len() - missing);
let wal_path = companion(&options.source, "-wal");
host_fs::write(&wal_path, &original).unwrap();
let before = host_fs::metadata(&wal_path).unwrap();
let main_before = host_fs::read(&options.source).unwrap();
let cx = attached_context();
let calls = Cell::new(0);
let (opened, report) = options
.repair_and_open(&cx, |_, _| {
calls.set(calls.get() + 1);
std::future::ready(Ok(()))
})
.await
.unwrap();
opened.unwrap();
let after = host_fs::metadata(&wal_path).unwrap();
assert_eq!((before.dev(), before.ino()), (after.dev(), after.ino()));
assert_eq!(calls.get(), 1);
assert_eq!(report.wal_frames, 3);
assert_eq!(report.repaired_frames, 2);
assert!(report.repaired_in_place);
assert_eq!(report.digest, blake3::hash(&repaired));
assert_eq!(host_fs::read(&wal_path).unwrap(), repaired);
assert_eq!(host_fs::read(&options.destination).unwrap(), original);
assert_eq!(host_fs::read(&options.source).unwrap(), main_before);
assert_recovery_available(&options.source, &cx);
}
});
}
#[test]
fn torn_tail_export_materializes_latest_page_without_modifying_sources() {
with_runtime(async {
for missing in [1, 512] {
let (options, mut original, repaired) = handoff_fixture();
original.truncate(original.len() - missing);
let wal_path = companion(&options.source, "-wal");
host_fs::write(&wal_path, &original).unwrap();
let main_before = host_fs::read(&options.source).unwrap();
let sidecar_path = companion(&options.source, "-wal-fec");
let sidecar_before = host_fs::read(&sidecar_path).unwrap();
let cx = attached_context();
let report = crate::native_recovery::export_database(&cx, &options)
.await
.unwrap();
let expected = &repaired[repaired.len() - 512..];
assert_eq!(host_fs::read(&report.destination).unwrap(), expected);
assert_eq!(report.digest, blake3::hash(expected));
assert_eq!(report.wal_frames, 3);
assert_eq!(report.repaired_frames, 2);
assert!(!report.repaired_in_place);
assert_eq!(host_fs::read(&options.source).unwrap(), main_before);
assert_eq!(host_fs::read(&wal_path).unwrap(), original);
assert_eq!(host_fs::read(&sidecar_path).unwrap(), sidecar_before);
assert_recovery_available(&options.source, &cx);
}
});
}
#[test]
fn torn_tail_refusals_leave_source_and_backup_namespace_untouched() {
with_runtime(async {
for exceed_budget in [false, true] {
let (mut options, mut original, _) = handoff_fixture();
if !exceed_budget {
let terminal = original.len() - 512 - WAL_FRAME_HEADER_SIZE;
original[terminal + 16] ^= 1;
}
original.truncate(original.len() - 1);
if exceed_budget {
options.replay.max_wal_bytes = original.len();
}
let wal_path = companion(&options.source, "-wal");
host_fs::write(&wal_path, &original).unwrap();
let cx = attached_context();
let calls = Cell::new(0);
assert!(
options
.repair_and_open(&cx, |_, _| {
calls.set(calls.get() + 1);
std::future::ready(Ok(()))
})
.await
.is_err()
);
assert_eq!(calls.get(), 0);
assert_eq!(host_fs::read(&wal_path).unwrap(), original);
assert!(
!NativeVfs::new()
.path_entry_exists(&cx, &options.destination)
.unwrap()
);
assert_recovery_available(&options.source, &cx);
}
});
}
fn byte_plan(original: &[u8]) -> RepairPlan {
let mut target = original.to_vec();
target[40..60].fill(0x77);
target[100..130].fill(0x66);
RepairPlan {
target,
changed: vec![40..60, 100..130],
index_regions: Vec::new(),
index_header: WalIndexHdr {
i_version: WAL_INDEX_VERSION,
unused: 0,
i_change: 0,
is_init: 1,
big_end_cksum: 0,
sz_page: 512,
mx_frame: 0,
n_page: 0,
a_frame_cksum: [0; 2],
a_salt: [0; 2],
a_cksum: [0; 2],
},
pages: 1,
certificate_anchors: 0,
}
}
#[test]
fn physical_settlement_preserves_header_length_and_untouched_bytes() {
let original = vec![0x11; 256];
let plan = byte_plan(&original);
let mut file = Cursor::new(original.clone());
let digest = settle_writes(&mut file, &Cx::new(), &original, &plan, |_| Ok(())).unwrap();
assert_eq!(file.get_ref(), &plan.target);
assert_eq!(digest, blake3::hash(&plan.target));
assert_eq!(
&file.get_ref()[..WAL_HEADER_SIZE],
&original[..WAL_HEADER_SIZE]
);
assert_eq!(file.get_ref().len(), original.len());
}
#[test]
fn failed_sync_restores_the_original_and_never_returns_a_receipt() {
let original = vec![0x11; 256];
let plan = byte_plan(&original);
let mut file = Cursor::new(original.clone());
let mut calls = 0;
let error = settle_writes(&mut file, &Cx::new(), &original, &plan, |_| {
calls += 1;
if calls == 1 {
Err(io::Error::other("injected sync failure"))
} else {
Ok(())
}
})
.unwrap_err();
assert_eq!(calls, 2);
assert_eq!(file.into_inner(), original);
assert!(
error
.to_string()
.contains("original WAL restored and synced")
);
}
#[test]
fn failed_rollback_sync_is_reported_as_indeterminate() {
let original = vec![0x11; 256];
let plan = byte_plan(&original);
let mut file = Cursor::new(original.clone());
let error = settle_writes(&mut file, &Cx::new(), &original, &plan, |_| {
Err(io::Error::other("persistent sync failure"))
})
.unwrap_err();
assert!(error.to_string().contains("indeterminate"));
}
#[test]
fn bad_readback_restores_the_original_before_returning_error() {
let original = vec![0x11; 256];
let plan = byte_plan(&original);
let mut file = Cursor::new(original.clone());
let mut calls = 0;
let error = settle_writes(&mut file, &Cx::new(), &original, &plan, |file| {
calls += 1;
if calls == 1 {
file.get_mut()[45] ^= 1;
}
Ok(())
})
.unwrap_err();
assert_eq!(file.into_inner(), original);
assert!(error.to_string().contains("original WAL restored"));
}
struct TornWrite {
file: Cursor<Vec<u8>>,
bytes_until_failure: usize,
failed: bool,
fail_truncate: bool,
}
impl Read for TornWrite {
fn read(&mut self, bytes: &mut [u8]) -> io::Result<usize> {
self.file.read(bytes)
}
}
impl Seek for TornWrite {
fn seek(&mut self, position: SeekFrom) -> io::Result<u64> {
self.file.seek(position)
}
}
impl Write for TornWrite {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
if !self.failed {
if self.bytes_until_failure == 0 {
self.failed = true;
return Err(io::Error::other("injected torn physical write"));
}
let len = bytes.len().min(self.bytes_until_failure);
self.bytes_until_failure -= len;
return self.file.write(&bytes[..len]);
}
self.file.write(bytes)
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl RepairFile for TornWrite {
fn truncate(&mut self, len: u64) -> io::Result<()> {
if self.fail_truncate {
return Err(io::Error::other("injected rollback truncation failure"));
}
RepairFile::truncate(&mut self.file, len)
}
}
#[test]
fn every_partial_write_boundary_restores_the_exact_original() {
let original = vec![0x11; 256];
let plan = byte_plan(&original);
for prefix_bytes in 0..50 {
let mut file = TornWrite {
file: Cursor::new(original.clone()),
bytes_until_failure: prefix_bytes,
failed: false,
fail_truncate: false,
};
let error =
settle_writes(&mut file, &Cx::new(), &original, &plan, |_| Ok(())).unwrap_err();
assert!(file.failed);
assert_eq!(file.file.into_inner(), original);
assert!(
error
.to_string()
.contains("original WAL restored and synced")
);
}
}
#[test]
fn cancellation_after_mutation_does_not_interrupt_masked_settlement() {
let original = vec![0x11; 256];
let plan = byte_plan(&original);
let mut file = Cursor::new(original.clone());
let cx = Cx::new();
let _mask = cx.masked();
settle_writes(&mut file, &cx, &original, &plan, |_| {
cx.cancel();
Ok(())
})
.unwrap();
assert_eq!(file.into_inner(), plan.target);
}
fn growth_plan(original: &[u8]) -> RepairPlan {
let mut plan = byte_plan(original);
plan.target.resize(320, 0x55);
plan.target[240..].fill(0x55);
plan.changed.push(240..320);
plan
}
#[test]
fn growth_settlement_syncs_exact_target_without_replacing_file_identity() {
use std::os::unix::fs::MetadataExt;
let original = vec![0x11; 256];
let plan = growth_plan(&original);
let mut file = tempfile::tempfile().unwrap();
file.write_all(&original).unwrap();
let before = file.metadata().unwrap();
let digest = settle_writes(&mut file, &Cx::new(), &original, &plan, |file| {
file.sync_all()
})
.unwrap();
let after = file.metadata().unwrap();
assert_eq!((before.dev(), before.ino()), (after.dev(), after.ino()));
assert_eq!(after.len(), 320);
file.seek(SeekFrom::Start(0)).unwrap();
let mut bytes = Vec::new();
file.read_to_end(&mut bytes).unwrap();
assert_eq!(bytes, plan.target);
assert_eq!(digest, blake3::hash(&plan.target));
}
#[test]
fn every_partial_growth_write_boundary_restores_original_bytes_and_eof() {
let original = vec![0x11; 256];
let plan = growth_plan(&original);
for prefix_bytes in 0..130 {
let mut file = TornWrite {
file: Cursor::new(original.clone()),
bytes_until_failure: prefix_bytes,
failed: false,
fail_truncate: false,
};
let error =
settle_writes(&mut file, &Cx::new(), &original, &plan, |_| Ok(())).unwrap_err();
assert!(file.failed);
assert_eq!(file.file.into_inner(), original);
assert!(
error
.to_string()
.contains("original WAL restored and synced")
);
}
}
#[test]
fn failed_growth_sync_truncates_before_syncing_and_verifying_rollback() {
let original = vec![0x11; 256];
let plan = growth_plan(&original);
let mut file = Cursor::new(original.clone());
let mut calls = 0;
let error = settle_writes(&mut file, &Cx::new(), &original, &plan, |file| {
calls += 1;
if calls == 1 {
assert_eq!(file.get_ref().len(), 320);
Err(io::Error::other("injected growth sync failure"))
} else {
assert_eq!(file.get_ref(), &original);
Ok(())
}
})
.unwrap_err();
assert_eq!(calls, 2);
assert_eq!(file.into_inner(), original);
assert!(
error
.to_string()
.contains("original WAL restored and synced")
);
}
#[test]
fn failed_growth_readback_removes_the_corrupted_append() {
let original = vec![0x11; 256];
let plan = growth_plan(&original);
let mut file = Cursor::new(original.clone());
let mut calls = 0;
let error = settle_writes(&mut file, &Cx::new(), &original, &plan, |file| {
calls += 1;
if calls == 1 {
file.get_mut()[319] ^= 1;
}
Ok(())
})
.unwrap_err();
assert_eq!(calls, 2);
assert_eq!(file.into_inner(), original);
assert!(
error
.to_string()
.contains("original WAL restored and synced")
);
}
#[test]
fn failed_growth_rollback_truncation_is_indeterminate_not_restored() {
let original = vec![0x11; 256];
let plan = growth_plan(&original);
let mut file = TornWrite {
file: Cursor::new(original.clone()),
bytes_until_failure: usize::MAX,
failed: false,
fail_truncate: true,
};
let error = settle_writes(&mut file, &Cx::new(), &original, &plan, |_| {
Err(io::Error::other("injected growth sync failure"))
})
.unwrap_err();
assert_eq!(file.file.get_ref().len(), 320);
assert!(error.to_string().contains("indeterminate"));
assert!(error.to_string().contains("rollback truncation failure"));
assert!(
!error
.to_string()
.contains("original WAL restored and synced")
);
}
}