use crate::bulk::{pack_tree, ExternalSort};
use crate::budget::MemoryBudget;
use crate::io::{open_file, FileIo};
use crate::meta::Meta;
use crate::page::{PageKind, PageRef, PAGE_SIZE};
use crate::pool::BufferPool;
use crate::store::Config;
use crate::Result;
use std::collections::{HashMap, HashSet};
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
mod safe;
pub use safe::{recover_to, RecoveryClass, SafeRecoveryReport};
mod reader;
pub use reader::{CandidateReader, LeafCandidate, LeafEvent, LeafScanReport};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LostLeaf {
pub page_no: u32,
pub after_key: Option<Vec<u8>>,
}
#[derive(Debug, Default)]
pub struct RecoveryReport {
pub pages_scanned: u64,
pub leaves_kept: u64,
pub leaves_lost: u64,
pub entries_recovered: u64,
pub lost_ranges: Vec<LostLeaf>,
pub truncated_tail_bytes: u64,
pub wal_quarantined: Option<std::path::PathBuf>,
pub wal_bytes_kept: u64,
pub wal_bytes_set_aside: u64,
}
enum ReadOutcome<'b> {
Verified(PageRef<'b>),
ReadOk,
ReadFailed,
}
fn read_page<'b>(file: &dyn FileIo, buf: &'b mut [u8], no: u64) -> ReadOutcome<'b> {
if file.read_at(buf, no * PAGE_SIZE as u64).is_err() { return ReadOutcome::ReadFailed; }
match PageRef::open(buf, no as u32) {
Ok(p) => ReadOutcome::Verified(p),
Err(_) => ReadOutcome::ReadOk,
}
}
fn read_verified<'b>(file: &dyn FileIo, buf: &'b mut [u8], no: u64) -> Option<PageRef<'b>> {
match read_page(file, buf, no) {
ReadOutcome::Verified(p) => Some(p),
ReadOutcome::ReadOk | ReadOutcome::ReadFailed => None,
}
}
fn decode_leaf_record(rec: &[u8], page_no: u32) -> Result<(&[u8], &[u8], bool)> {
match crate::verify::decode_record(rec, page_no, PageKind::Leaf)? {
crate::verify::DecodedRecord::Leaf { key, value, overflow } => {
Ok((key, value, overflow))
}
crate::verify::DecodedRecord::Interior { .. } => unreachable!(),
}
}
fn copy_overflow(file: &dyn FileIo, pool: &BufferPool, marker: &[u8]) -> Result<Vec<u8>> {
use crate::btree::{OV_CAP, OV_DATA, OV_NEXT, OV_USED};
use crate::page::PageMut;
let bad = |page_no, why| crate::Error::Corrupt { page_no, why };
if marker.len() != 12 { return Err(bad(0, "recovery overflow marker wrong size")); }
let total = u32::from_le_bytes(marker[0..4].try_into().unwrap()) as usize;
let mut source = u32::from_le_bytes(marker[4..8].try_into().unwrap());
let want_crc = u32::from_le_bytes(marker[8..12].try_into().unwrap());
let expected_pages = total.div_ceil(OV_CAP).max(1);
let file_pages = file.len()? / PAGE_SIZE as u64;
let (mut seen, mut bytes, mut crc) = (0usize, 0usize, 0u32);
let (mut head, mut previous) = (0u32, 0u32);
let mut buf = [0u8; PAGE_SIZE];
while source != 0 {
seen += 1;
if seen > expected_pages || source as u64 >= file_pages {
return Err(bad(source, "recovery overflow chain cycles or exceeds file"));
}
file.read_at(&mut buf, source as u64 * PAGE_SIZE as u64)?;
let page = PageRef::open(&buf, source)?;
if page.kind() != PageKind::Overflow || page.tree_id() != 0 {
return Err(bad(source, "recovery overflow chain reaches wrong page kind"));
}
let used = u16::from_le_bytes(buf[OV_USED..OV_USED + 2].try_into().unwrap()) as usize;
let next = u32::from_le_bytes(buf[OV_NEXT..OV_NEXT + 4].try_into().unwrap());
if used > OV_CAP || bytes.checked_add(used).is_none_or(|n| n > total) {
return Err(bad(source, "recovery overflow length out of bounds"));
}
let chunk = &buf[OV_DATA..OV_DATA + used];
crc = crc32c::crc32c_append(crc, chunk);
bytes += used;
let mut w = pool.allocate()?;
let dest = w.page_no();
let b = w.bytes_mut();
PageMut::init(b, PageKind::Overflow, 0, dest).finalise(0);
b[OV_USED..OV_USED + 2].copy_from_slice(&(used as u16).to_le_bytes());
b[OV_DATA..OV_DATA + used].copy_from_slice(chunk);
drop(w);
if previous == 0 { head = dest; } else {
let mut w = pool.get_mut(previous)?;
w.bytes_mut()[OV_NEXT..OV_NEXT + 4].copy_from_slice(&dest.to_le_bytes());
}
previous = dest;
source = next;
}
if seen != expected_pages || bytes != total || crc != want_crc {
return Err(bad(0, "recovery overflow value fails manifest checksum"));
}
let mut relocated = marker.to_vec();
relocated[4..8].copy_from_slice(&head.to_le_bytes());
Ok(relocated)
}
fn tag(gen: u64, page_no: u32, val: &[u8]) -> Vec<u8> {
let mut t = Vec::with_capacity(12 + val.len());
t.extend_from_slice(&gen.to_be_bytes());
t.extend_from_slice(&page_no.to_be_bytes());
t.extend_from_slice(val);
t
}
fn untag(t: &[u8]) -> (&[u8], &[u8]) {
(&t[0..12], &t[12..])
}
static SEQ: AtomicU64 = AtomicU64::new(0);
pub fn recover(dir: &Path, cfg: Config) -> Result<RecoveryReport> {
if matches!(crate::limits::read(dir), Ok(Some(_))) {
return Err(crate::Error::ResourceLimit("in-place repair disabled for constrained stores; recover_to a separately provisioned destination"));
}
let data = dir.join("data");
let (file, _) = open_file(&data, cfg.io)?;
recover_impl(&*file, dir, cfg)
}
fn recover_impl(file: &dyn FileIo, dir: &Path, cfg: Config) -> Result<RecoveryReport> {
recover_impl_with_before_publish(file, dir, cfg, |_| Ok(()))
}
fn recover_impl_with_before_publish<F>(
file: &dyn FileIo,
dir: &Path,
cfg: Config,
before_publish: F,
) -> Result<RecoveryReport>
where
F: FnOnce(&Path) -> Result<()>,
{
let data = dir.join("data");
let len = file.len()?;
let total = len / PAGE_SIZE as u64;
let seq = SEQ.fetch_add(1, Ordering::Relaxed);
let scratch = std::env::temp_dir().join(format!("kernel-recover-{}-{}", std::process::id(), seq));
let arena = (cfg.budget_bytes / 3).max(4 << 20);
let mut sorter = ExternalSort::new(&scratch, arena)?;
let mut rep = RecoveryReport::default();
rep.truncated_tail_bytes = len % PAGE_SIZE as u64;
let mut max_lsn: u64 = 0;
let mut lost_pages: Vec<u32> = Vec::new();
let mut buf = vec![0u8; PAGE_SIZE];
for no in 1..total {
rep.pages_scanned += 1;
let page = match read_page(file, &mut buf, no) {
ReadOutcome::Verified(p) => p,
ReadOutcome::ReadOk => {
let kind = u16::from_le_bytes([buf[6], buf[7]]);
if kind == PageKind::Leaf as u16 {
rep.leaves_lost += 1;
lost_pages.push(no as u32);
}
continue;
}
ReadOutcome::ReadFailed => {
rep.leaves_lost += 1;
lost_pages.push(no as u32);
continue;
}
};
max_lsn = max_lsn.max(page.lsn());
if page.kind() != PageKind::Leaf { continue; }
rep.leaves_kept += 1;
let n = page.nentries();
for i in 0..n {
let rec = page.slot(i);
let (k, v, is_marker) = decode_leaf_record(rec, page.page_no())?;
sorter.push_flagged(k.to_vec(), tag(page.lsn(), page.page_no(), v), is_marker)?;
}
}
let lost_set: HashSet<u32> = lost_pages.iter().copied().collect();
let mut after_by_page: HashMap<u32, Vec<u8>> = HashMap::with_capacity(lost_pages.len());
if !lost_set.is_empty() {
for no in 1..total {
rep.pages_scanned += 1;
let Some(page) = read_verified(file, &mut buf, no) else { continue };
if page.kind() != PageKind::Leaf { continue; }
if !lost_set.contains(&page.next_leaf()) { continue; }
let n = page.nentries();
if n == 0 { continue; }
let (last_k, _, _) = decode_leaf_record(page.slot(n - 1), page.page_no())?;
after_by_page.insert(page.next_leaf(), last_k.to_vec());
}
}
for page_no in lost_pages {
rep.lost_ranges.push(LostLeaf { page_no, after_key: after_by_page.get(&page_no).cloned() });
}
let tmp = dir.join("data.rebuild");
let _ = std::fs::remove_file(&tmp);
let mut entries_recovered = 0u64;
let next_lsn = max_lsn.checked_add(1).ok_or(crate::Error::Corrupt {
page_no: 0,
why: "recovered LSN high-water mark is exhausted",
})?;
let mut rebuilt_roots = [0u32; crate::meta::MAX_TREES];
let build: Result<()> = (|| {
let (nf, _) = open_file(&tmp, cfg.io)?;
let budget = Arc::new(MemoryBudget::new(cfg.budget_bytes));
let frames = ((cfg.budget_bytes * 2 / 3) / PAGE_SIZE).max(16);
let pool = BufferPool::new(nf.into(), budget, frames)?;
let _meta_page = pool.allocate()?; drop(_meta_page); let _slot_b = pool.allocate()?; drop(_slot_b);
Meta::init_slot_b(&pool)?;
let mut runs = sorter.finish()?;
let merged = runs.iter()?;
struct Dedup<I> { inner: I, winner: Option<(Vec<u8>, Vec<u8>, bool)>, done: bool }
impl<I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>> Iterator for Dedup<I> {
type Item = Result<(Vec<u8>, Vec<u8>, bool)>;
fn next(&mut self) -> Option<Self::Item> {
if self.done { return None; }
loop {
match self.inner.next() {
Some(Ok((k, tagged_v, m))) => {
match &mut self.winner {
None => self.winner = Some((k, tagged_v, m)),
Some((wk, wv, wm)) if *wk == k => {
let (worder, _) = untag(wv);
let (order, _) = untag(&tagged_v);
if order > worder { *wv = tagged_v; *wm = m; }
}
Some(_) => {
let (out_k, out_v, out_m) = self.winner.take().unwrap();
self.winner = Some((k, tagged_v, m));
let (_, real_v) = untag(&out_v);
return Some(Ok((out_k, real_v.to_vec(), out_m)));
}
}
}
Some(Err(e)) => { self.done = true; return Some(Err(e)); }
None => {
self.done = true;
return self.winner.take().map(|(k, v, m)| {
let (_, real_v) = untag(&v);
Ok((k, real_v.to_vec(), m))
});
}
}
}
}
}
let recovered = Dedup { inner: merged, winner: None, done: false }
.filter(|item| !matches!(item, Ok((key, _, _))
if crate::keys::is_field_aggregate_key(key)));
let relocated = recovered.map(|item| {
let (key, value, overflow) = item?;
let value = if overflow { copy_overflow(file, &pool, &value)? } else { value };
Ok((key, value, overflow))
});
let counting = CountingIter { inner: relocated, count: &mut entries_recovered };
let root = pack_tree(&pool, 1, counting, 0.9, &scratch)?;
rebuilt_roots[0] = root;
Meta { format_version: crate::meta::FORMAT_VERSION, roots: rebuilt_roots, next_lsn, generation: 0 }.write(&pool)?;
Meta::mark_salvaged(&pool)?;
pool.flush_all(crate::io::Barrier::Full)?;
Ok(())
})();
if let Err(e) = build {
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
if let Err(e) = before_publish(&tmp).and_then(|_| {
let expected = Meta {
format_version: crate::meta::FORMAT_VERSION,
roots: rebuilt_roots,
next_lsn,
generation: 0,
};
let verified = crate::verify::verify_rebuild(&tmp, cfg.io, &expected, entries_recovered)?;
debug_assert_eq!(verified.rows, entries_recovered);
debug_assert!(verified.pages > 0);
Ok(())
}) {
let _ = std::fs::remove_file(&tmp);
return Err(e);
}
let rebuilt_pages = std::fs::metadata(&tmp)?.len() / PAGE_SIZE as u64;
let rebuilt_pages = u32::try_from(rebuilt_pages).map_err(|_| crate::Error::Corrupt {
page_no: 0,
why: "rebuilt file has too many pages for its freelist bound",
})?;
let empty_free = BufferPool::empty_free(0);
crate::verify::publish_freelist(dir, empty_free, 0, rebuilt_pages, file, true, true)?;
std::fs::rename(&tmp, &data)?;
let (f2, _) = open_file(&data, cfg.io)?;
f2.sync_dir()?;
quarantine_damaged_wal(dir, cfg, &mut rep)?;
rep.entries_recovered = entries_recovered;
Ok(rep)
}
fn quarantine_damaged_wal(dir: &Path, cfg: Config, rep: &mut RecoveryReport) -> Result<()> {
let wal = dir.join("wal");
if !wal.exists() { return Ok(()); }
let scan = crate::wal::Wal::inspect(&wal, cfg.io)?;
match scan.stop {
crate::wal::Stop::End(_) => return Ok(()),
crate::wal::Stop::Damaged { .. } => {}
}
let mut n = 0u32;
let aside = loop {
let c = dir.join(format!("wal.corrupt.{n}"));
if !c.exists() { break c; }
n = n.checked_add(1).ok_or(crate::Error::TooLarge)?;
};
let total = std::fs::metadata(&wal)?.len();
let source_hash = crate::wal::hash_prefix(&wal, total)?;
std::fs::copy(&wal, &aside)?;
std::fs::File::open(&aside)?.sync_all()?;
if std::fs::metadata(&aside)?.len() != total
|| crate::wal::hash_prefix(&aside, total)? != source_hash
{
let _ = std::fs::remove_file(&aside);
return Err(crate::Error::CorruptWal {
offset: scan.end,
why: "quarantined WAL copy failed independent read-back hashing",
});
}
let tmp = dir.join("wal.rebuild");
let _ = std::fs::remove_file(&tmp);
let salvaged = crate::wal::Wal::salvage_committed(&wal, &tmp)?;
let verified = crate::wal::Wal::inspect(&tmp, cfg.io)?;
if verified.end != salvaged.bytes
|| !matches!(verified.stop, crate::wal::Stop::End(_))
|| std::fs::metadata(&tmp)?.len() != salvaged.bytes
|| crate::wal::hash_prefix(&tmp, salvaged.bytes)? != salvaged.hash
{
let _ = std::fs::remove_file(&tmp);
return Err(crate::Error::CorruptWal {
offset: scan.end,
why: "reconstructed WAL failed independent parse and hash verification",
});
}
std::fs::rename(&tmp, &wal)?;
let (f, _) = open_file(&wal, crate::io::IoMode::Buffered)?;
f.sync_dir()?;
rep.wal_quarantined = Some(aside);
rep.wal_bytes_kept = salvaged.bytes;
rep.wal_bytes_set_aside = total;
Ok(())
}
struct CountingIter<'a, I> { inner: I, count: &'a mut u64 }
impl<'a, I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>> Iterator for CountingIter<'a, I> {
type Item = Result<(Vec<u8>, Vec<u8>, bool)>;
fn next(&mut self) -> Option<Self::Item> {
let item = self.inner.next();
if let Some(Ok(_)) = &item { *self.count += 1; }
item
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::io::IoMode;
use crate::store::{Config, Store, SyncMode};
fn cfg() -> Config { Config { budget_bytes: 16 << 20, io: IoMode::Buffered, sync: SyncMode::Off } }
#[test]
fn a_damaged_interior_page_costs_nothing() {
let d = tempfile::tempdir().unwrap();
let n = 5_000u64;
{ let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"v".to_vec()))).unwrap();
s.commit().unwrap(); s.checkpoint().unwrap(); }
let path = d.path().join("data");
let mut bytes = std::fs::read(&path).unwrap();
let ps = PAGE_SIZE;
let mut wrecked_interior = false;
for p in 1..bytes.len() / ps {
let kind = u16::from_le_bytes([bytes[p * ps + 6], bytes[p * ps + 7]]);
if kind == PageKind::Interior as u16 && !wrecked_interior {
bytes[p * ps + 100] ^= 0xff;
wrecked_interior = true;
}
}
assert!(wrecked_interior, "the fixture needs at least one interior page");
std::fs::write(&path, &bytes).unwrap();
let report = recover(d.path(), cfg()).unwrap();
assert_eq!(report.leaves_lost, 0, "a damaged interior page must cost no leaf loss");
assert_eq!(report.entries_recovered, n, "and no data may be missing either");
}
#[test]
fn a_duplicate_key_across_two_surviving_generations_keeps_the_newer_one() {
let d = tempfile::tempdir().unwrap();
let n = 500u64;
let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"stale".to_vec()))).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"current".to_vec()))).unwrap();
s.commit().unwrap(); s.checkpoint().unwrap();
drop(s);
let bytes = std::fs::read(d.path().join("data")).unwrap();
let mut stale_leaves = 0;
let mut current_leaves = 0;
for p in 1..bytes.len() / PAGE_SIZE {
let b = &bytes[p * PAGE_SIZE..(p + 1) * PAGE_SIZE];
if let Ok(pr) = PageRef::open(b, p as u32) {
if pr.kind() == PageKind::Leaf && pr.nentries() > 0 {
let (_, v, _) = decode_leaf_record(pr.slot(0), pr.page_no()).unwrap();
if v == b"stale" { stale_leaves += 1; }
if v == b"current" { current_leaves += 1; }
}
}
}
assert!(stale_leaves > 0 && current_leaves > 0,
"fixture must retain both generations' leaves unreclaimed");
let report = recover(d.path(), cfg()).unwrap();
assert_eq!(report.entries_recovered, n,
"duplicates must collapse to one entry per key, not {}",
stale_leaves + current_leaves);
let reopened = Store::open(d.path(), cfg()).unwrap();
for i in (0..n).step_by(37) {
assert_eq!(
reopened.get(&i.to_be_bytes()).unwrap().as_deref(),
Some(&b"current"[..]),
"key {i} must resolve to the newer generation, not the stale one"
);
}
}
fn leaves_with_last_keys(bytes: &[u8]) -> Vec<(u32, Vec<u8>)> {
let mut out = Vec::new();
for p in 1..bytes.len() / PAGE_SIZE {
let b = &bytes[p * PAGE_SIZE..(p + 1) * PAGE_SIZE];
if let Ok(pr) = PageRef::open(b, p as u32) {
if pr.kind() == PageKind::Leaf && pr.nentries() > 0 {
let (last_k, _, _) =
decode_leaf_record(pr.slot(pr.nentries() - 1), pr.page_no()).unwrap();
out.push((p as u32, last_k.to_vec()));
}
}
}
out
}
fn wreck(bytes: &mut [u8], page_no: u32) {
let base = page_no as usize * PAGE_SIZE;
bytes[base + 50] ^= 0xff;
}
fn forge_next_leaf(bytes: &mut [u8], page_no: u32, target: u32) {
let base = page_no as usize * PAGE_SIZE;
let page = &mut bytes[base..base + PAGE_SIZE];
let gen = u64::from_le_bytes(page[24..32].try_into().unwrap());
let mut p = crate::page::PageMut::reopen(page);
p.set_next_leaf(target);
crate::page::seal(page, gen);
}
#[test]
fn a_lost_leafs_after_key_comes_from_the_sibling_chain_not_page_order() {
let d = tempfile::tempdir().unwrap();
let n = 2_000u64;
{ let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"v".to_vec()))).unwrap();
s.commit().unwrap(); s.checkpoint().unwrap(); }
let path = d.path().join("data");
let mut bytes = std::fs::read(&path).unwrap();
let leaves = leaves_with_last_keys(&bytes);
assert!(leaves.len() >= 2, "fixture needs at least two non-empty leaves");
let (pointer_no, pointer_last_key) = leaves[0].clone();
let (victim_no, _) = leaves[1].clone();
assert_ne!(pointer_no, victim_no);
forge_next_leaf(&mut bytes, pointer_no, victim_no);
wreck(&mut bytes, victim_no);
let bytes_for_half_b = bytes.clone();
std::fs::write(&path, &bytes).unwrap();
let report = recover(d.path(), cfg()).unwrap();
let victim_entry = report.lost_ranges.iter().find(|l| l.page_no == victim_no)
.expect("the victim page must be reported as a loss");
assert_eq!(
victim_entry.after_key,
Some(pointer_last_key.clone()),
"after_key must come from the leaf whose next_leaf names the lost page"
);
let d2 = tempfile::tempdir().unwrap();
let mut bytes2 = bytes_for_half_b;
wreck(&mut bytes2, pointer_no);
std::fs::write(d2.path().join("data"), &bytes2).unwrap();
let report2 = recover(d2.path(), cfg()).unwrap();
let victim_entry2 = report2.lost_ranges.iter().find(|l| l.page_no == victim_no)
.expect("the victim page must still be reported as a loss");
assert_eq!(
victim_entry2.after_key, None,
"after_key must be None once the only pointer to the lost page is itself unreadable"
);
}
#[test]
fn a_truncated_tail_is_named_not_silently_dropped() {
let d = tempfile::tempdir().unwrap();
let n = 2_000u64;
{ let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"v".to_vec()))).unwrap();
s.commit().unwrap(); s.checkpoint().unwrap(); }
let path = d.path().join("data");
let full_len = std::fs::metadata(&path).unwrap().len();
assert_eq!(full_len % PAGE_SIZE as u64, 0, "fixture must start page-aligned");
let short_by = 137u64;
let truncated_len = full_len - short_by;
let f = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
f.set_len(truncated_len).unwrap();
drop(f);
let report = recover(d.path(), cfg()).unwrap();
assert_eq!(
report.truncated_tail_bytes,
PAGE_SIZE as u64 - short_by,
"the dangling bytes at the end of the file must be named exactly"
);
assert_eq!(report.leaves_lost, 0, "a truncated tail must not be counted as a lost leaf");
assert_eq!(report.entries_recovered, n, "every complete page's data must still be recovered");
}
struct FailAt { inner: Box<dyn FileIo>, fail_offset: u64 }
impl FileIo for FailAt {
fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> {
if off == self.fail_offset {
return Err(std::io::Error::new(std::io::ErrorKind::Other, "injected read failure").into());
}
self.inner.read_at(buf, off)
}
fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
fn len(&self) -> Result<u64> { self.inner.len() }
fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
}
#[test]
fn an_unreadable_page_is_counted_regardless_of_the_stale_buffer_it_follows() {
let d = tempfile::tempdir().unwrap();
let n = 3_000u64;
let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"gen1".to_vec()))).unwrap();
s.bulk_load((n..2 * n).map(|i| (i.to_be_bytes().to_vec(), b"gen2".to_vec()))).unwrap();
s.commit().unwrap(); s.checkpoint().unwrap();
drop(s);
let path = d.path().join("data");
let bytes = std::fs::read(&path).unwrap();
let total = bytes.len() / PAGE_SIZE;
let kind_of = |p: usize| -> Option<PageKind> {
PageRef::open(&bytes[p * PAGE_SIZE..(p + 1) * PAGE_SIZE], p as u32).ok().map(|pr| pr.kind())
};
let victim = (2..total as u32)
.find(|&p| kind_of(p as usize - 1) == Some(PageKind::Interior)
&& kind_of(p as usize) == Some(PageKind::Leaf))
.expect("two generations must produce an interior-then-leaf boundary");
let (real_file, _) = crate::io::open_file(&path, IoMode::Buffered).unwrap();
let failing = FailAt { inner: real_file, fail_offset: victim as u64 * PAGE_SIZE as u64 };
let report = recover_impl(&failing, d.path(), cfg()).unwrap();
assert!(
report.lost_ranges.iter().any(|l| l.page_no == victim),
"page {victim} (predecessor is Interior) failed to read and must be counted as lost, \
not silently skipped because the stale buffer said Interior"
);
let buf = &bytes[(victim as usize - 1) * PAGE_SIZE..victim as usize * PAGE_SIZE];
let stale_kind = u16::from_le_bytes([buf[6], buf[7]]);
assert_eq!(
stale_kind,
PageKind::Interior as u16,
"sanity: the stale buffer a failed read would have left behind really does say Interior"
);
let old_behaviour_would_count_it = stale_kind == PageKind::Leaf as u16;
assert!(
!old_behaviour_would_count_it,
"the old kind-byte-fallback logic would NOT have counted page {victim} -- \
confirming this is a genuine fix, not a no-op"
);
}
#[test]
fn recover_does_not_panic_on_an_empty_but_crc_valid_meta_page() {
let d = tempfile::tempdir().unwrap();
let n = 500u64;
{ let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..n).map(|i| (i.to_be_bytes().to_vec(), b"v".to_vec()))).unwrap();
s.commit().unwrap(); s.checkpoint().unwrap(); }
let path = d.path().join("data");
let mut bytes = std::fs::read(&path).unwrap();
{
let page0 = &mut bytes[0..PAGE_SIZE];
let mut p = crate::page::PageMut::init(page0, PageKind::Meta, 0, 0);
p.finalise(0);
}
std::fs::write(&path, &bytes).unwrap();
let report = recover(d.path(), cfg())
.expect("recover() must degrade on a malformed superblock, not panic");
assert!(report.leaves_kept > 0, "the surviving leaves must still be found and repacked");
assert_eq!(report.entries_recovered, n, "every leaf's entries survive an empty page 0");
}
#[test]
fn a_corrupted_rebuild_never_replaces_the_original_data_file() {
let d = tempfile::tempdir().unwrap();
{
let mut s = Store::create(d.path(), cfg()).unwrap();
s.bulk_load((0..2_000u64).map(|i| {
(i.to_be_bytes().to_vec(), i.to_le_bytes().to_vec())
})).unwrap();
}
let data = d.path().join("data");
let original = std::fs::read(&data).unwrap();
let (file, _) = open_file(&data, cfg().io).unwrap();
let result = recover_impl_with_before_publish(&*file, d.path(), cfg(), |fresh| {
let mut bytes = std::fs::read(fresh)?;
bytes[PAGE_SIZE * 2 + 100] ^= 0x80;
std::fs::write(fresh, bytes)?;
Ok(())
});
assert!(result.is_err(), "a rebuild corrupted before publication must be refused");
assert_eq!(std::fs::read(&data).unwrap(), original,
"the original data file must remain byte-for-byte authoritative");
let s = Store::open(d.path(), cfg()).unwrap();
assert_eq!(s.get(&1999u64.to_be_bytes()).unwrap().as_deref(),
Some(&1999u64.to_le_bytes()[..]));
}
#[test]
fn an_unreadable_wal_is_an_error_while_no_wal_is_success() {
let damaged = tempfile::tempdir().unwrap();
{
let mut s = Store::create(damaged.path(), cfg()).unwrap();
s.put(b"published", b"yes").unwrap();
s.commit().unwrap();
s.checkpoint().unwrap();
s.put(b"committed-tail", b"must not be stranded").unwrap();
s.commit().unwrap();
}
let wal = damaged.path().join("wal");
std::fs::rename(&wal, damaged.path().join("wal.committed-tail")).unwrap();
std::fs::create_dir(&wal).unwrap();
let result = recover(damaged.path(), cfg());
assert!(result.is_err(), "an unreadable committed log must not be reported as recovered");
let healthy = tempfile::tempdir().unwrap();
{
let mut s = Store::create(healthy.path(), cfg()).unwrap();
s.put(b"published", b"yes").unwrap();
s.commit().unwrap();
s.checkpoint().unwrap();
}
std::fs::remove_file(healthy.path().join("wal")).unwrap();
assert!(recover(healthy.path(), cfg()).is_ok(), "there is nothing to inspect when no log exists");
}
}