use kernel::graph::Graph;
use kernel::store::{Config, Store};
const PAGE_SIZE: usize = 4096;
fn stamps(dir: &std::path::Path) -> std::collections::BTreeMap<u32, u64> {
let data = std::fs::read(dir.join("data")).unwrap();
let mut out = std::collections::BTreeMap::new();
for (i, page) in data.chunks_exact(PAGE_SIZE).enumerate() {
let magic = u32::from_le_bytes(page[0..4].try_into().unwrap());
let kind = u16::from_le_bytes(page[6..8].try_into().unwrap());
if magic != 0x53454B32 || kind == 0 { continue; } let no = u32::from_le_bytes(page[12..16].try_into().unwrap());
if no as usize != i { continue; }
let gen = u64::from_le_bytes(page[24..32].try_into().unwrap());
out.insert(no, gen);
}
out
}
#[test]
fn every_flushed_page_carries_its_publishing_generation() {
let d = tempfile::TempDir::new().unwrap();
let mut g = Graph::new(Store::create(d.path(), Config::default()).unwrap()).unwrap();
for i in 1..=500u64 {
g.add_node(None, 7, format!("railway ledger row {i}").as_bytes()).unwrap();
}
g.commit().unwrap();
g.checkpoint().unwrap();
let first = stamps(d.path());
assert!(!first.is_empty());
let gens: std::collections::HashSet<u64> = first.values().copied().collect();
assert_eq!(gens.len(), 1, "one epoch published: one generation everywhere, got {gens:?}");
let g1 = *gens.iter().next().unwrap();
assert!(g1 > 0, "stamp must be a real generation, not the reserved 0");
for i in 501..=550u64 {
g.add_node(None, 7, format!("second epoch row {i}").as_bytes()).unwrap();
}
g.commit().unwrap();
g.checkpoint().unwrap();
let second = stamps(d.path());
let g2max = *second.values().max().unwrap();
assert!(g2max > g1, "the second epoch must stamp a higher generation");
assert!(second.values().any(|&v| v == g1),
"untouched pages keep the first epoch's stamp");
}
#[test]
fn recovery_prefers_higher_generation_over_higher_page_number() {
use kernel::store::{Config, Store};
const PS: usize = 4096;
let d = tempfile::TempDir::new().unwrap();
{
let mut s = Store::create(d.path(), Config::default()).unwrap();
for i in 0..2_000u64 {
s.put(&i.to_be_bytes(), b"current-generation-v").unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
let path = d.path().join("data");
let mut bytes = std::fs::read(&path).unwrap();
let mut found: Option<(usize, Vec<u8>, usize, usize)> = None; for (i, page) in bytes.chunks_exact(PS).enumerate() {
let magic = u32::from_le_bytes(page[0..4].try_into().unwrap());
let kind = u16::from_le_bytes(page[6..8].try_into().unwrap());
let no = u32::from_le_bytes(page[12..16].try_into().unwrap());
let nentries = u16::from_le_bytes(page[10..12].try_into().unwrap());
let tid = u16::from_le_bytes(page[8..10].try_into().unwrap());
if magic != 0x53454B32 || kind != 2 || tid != 1 || no as usize != i || nentries < 10 { continue; }
let off = u16::from_le_bytes(page[40..42].try_into().unwrap()) as usize;
let raw = u16::from_le_bytes(page[off..off + 2].try_into().unwrap()) as usize;
let compact = cfg!(feature = "compact-cells") && raw & 0xf000 == 0x4000;
let klen = if compact { raw & 0x0fff } else { raw };
let key = page[off + 2..off + 2 + klen].to_vec();
let (value_at, vlen) = if compact {
let slot_len = u16::from_le_bytes(page[42..44].try_into().unwrap()) as usize;
(off + 2 + klen, slot_len - 2 - klen)
} else {
(off + 4 + klen, u16::from_le_bytes(page[off + 2 + klen..off + 4 + klen].try_into().unwrap()) as usize)
};
if vlen != 20 { continue; } found = Some((i, key, value_at, vlen));
break;
}
let (low, key, val_off, val_len) = found.expect("a decodable leaf");
let end_no = bytes.len() / PS;
let mut dup = bytes[low * PS..(low + 1) * PS].to_vec();
dup[12..16].copy_from_slice(&(end_no as u32).to_le_bytes());
dup[val_off..val_off + val_len].copy_from_slice(b"stale-generation-vXX");
let orig_gen = u64::from_le_bytes(dup[24..32].try_into().unwrap());
assert!(orig_gen >= 1, "pages must carry a stamp by now");
kernel::page::seal(&mut dup, orig_gen.saturating_sub(1));
bytes.extend_from_slice(&dup);
std::fs::write(&path, &bytes).unwrap();
kernel::recover::recover(d.path(), Config::default()).unwrap();
let s = Store::open(d.path(), Config::default()).unwrap();
let v = s.get(&key).unwrap().expect("the key must survive recovery");
assert_eq!(&v, b"current-generation-v",
"the higher-generation copy must win over a higher page number");
}
#[test]
fn fold_insert_fold_cycles_lose_nothing() {
use kernel::graph::Graph;
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
let mut g = Graph::new(Store::create(d.path(), Config::default()).unwrap()).unwrap();
let count = |g: &Graph| -> usize {
let mut n = 0;
let it = g.store_ref().scan(&[0x0C]).unwrap();
it.for_each_ref(|k, _| {
if k.first() == Some(&0x0C) { n += 1; true } else { false }
}).unwrap();
n
};
let mut expected_seg_rows = 0usize;
for round in 1..=3u64 {
for i in 0..100u64 {
let id = round * 1000 + i;
g.index_text(1, id, &format!("railway survey number {id}")).unwrap();
}
g.commit().unwrap();
g.checkpoint().unwrap();
g.fold_text(1).unwrap();
g.checkpoint().unwrap();
expected_seg_rows += 103;
assert_eq!(count(&g), expected_seg_rows,
"round {round}: all folded segments must survive");
let hits = g.text_search(1, "1000", 5).unwrap();
assert!(!hits.is_empty(), "round {round}: round-1 data must stay searchable");
}
}
#[test]
fn overwrite_churn_plateaus_the_file() {
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
let mut s = Store::create(d.path(), Config::default()).unwrap();
let mut size_at = Vec::new();
for round in 0..6u64 {
for i in 0..2_000u64 {
s.put(format!("churn-{i:05}").as_bytes(),
format!("round {round} value of {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
size_at.push(std::fs::metadata(d.path().join("data")).unwrap().len());
}
assert_eq!(size_at[2], size_at[5],
"constant live data must mean a constant file: {size_at:?}");
}
#[test]
fn a_corrupt_registered_reader_marker_halts_recycling() {
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
let mut s = Store::create(d.path(), Config::default()).unwrap();
for i in 0..2_000u64 {
s.put(format!("base-{i:05}").as_bytes(), format!("stable value {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
let tiny = Config { budget_bytes: 1 << 16, ..Config::default() };
let reader = Store::open_snapshot(d.path(), tiny).unwrap();
let marker = std::fs::read_dir(d.path().join("readers"))
.unwrap().next().unwrap().unwrap().path();
let mut marker_bytes = std::fs::read(&marker).unwrap();
let last = marker_bytes.len() - 1;
marker_bytes[last] ^= 0x80;
std::fs::write(&marker, marker_bytes).unwrap();
let before: Vec<_> = (0..2_000u64)
.map(|i| reader.get(format!("base-{i:05}").as_bytes()).unwrap()).collect();
for round in 0..8u64 {
for i in 0..2_000u64 {
s.put(format!("base-{i:05}").as_bytes(),
format!("round {round} rewrite {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
let after: Vec<_> = (0..2_000u64)
.map(|i| reader.get(format!("base-{i:05}").as_bytes()).unwrap()).collect();
assert_eq!(before, after, "a pinned reader's world moved");
drop(reader);
}
#[test]
fn the_freelist_survives_reopen() {
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
{
let mut s = Store::create(d.path(), Config::default()).unwrap();
for round in 0..3u64 {
for i in 0..2_000u64 {
s.put(format!("churn-{i:05}").as_bytes(),
format!("round {round} value {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
assert!(s.pool_ref().free_pages_pending() > 10, "churn must have freed pages");
}
let s = Store::open(d.path(), Config::default()).unwrap();
assert!(s.pool_ref().free_pages_pending() > 10,
"the persisted freelist must survive a reopen");
}
#[test]
fn a_corrupt_freelist_sidecar_only_leaks() {
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
{
let mut s = Store::create(d.path(), Config::default()).unwrap();
for i in 0..1_000u64 {
s.put(format!("row-{i:05}").as_bytes(), b"value bytes here").unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
std::fs::write(d.path().join("free"), b"not a freelist at all").unwrap();
let s = Store::open(d.path(), Config::default()).unwrap();
assert_eq!(s.pool_ref().free_pages_pending(), 0, "corrupt sidecar: empty list");
assert_eq!(s.get(b"row-00500").unwrap().as_deref(), Some(&b"value bytes here"[..]),
"data untouched by sidecar corruption");
}
#[test]
fn a_crash_after_data_publish_cannot_reimport_the_previous_freelist() {
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
let cfg = Config::default();
let stale = {
let mut s = Store::create(d.path(), cfg).unwrap();
for round in 0..2u64 {
for i in 0..2_000u64 {
s.put(format!("row-{i:05}").as_bytes(),
format!("round {round} value of row {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
assert!(s.pool_ref().free_pages_pending() > 0,
"generation 2 must persist pages generation 3 can recycle");
let stale = std::fs::read(d.path().join("free")).unwrap();
for i in 0..2_000u64 {
s.put(format!("row-{i:05}").as_bytes(),
format!("published generation 3 row {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
stale
};
std::fs::write(d.path().join("free"), stale).unwrap();
let mut reopened = Store::open(d.path(), cfg).unwrap();
assert_eq!(reopened.pool_ref().free_pages_pending(), 0,
"a sidecar from another published generation must be discarded");
for i in 0..2_000u64 {
reopened.put(format!("after-{i:05}").as_bytes(), b"allocate after crash").unwrap();
}
reopened.commit().unwrap();
reopened.checkpoint().unwrap();
for i in 0..2_000u64 {
assert_eq!(reopened.get(format!("row-{i:05}").as_bytes()).unwrap().as_deref(),
Some(format!("published generation 3 row {i}").as_bytes()),
"post-crash allocation overwrote committed row {i}");
}
}
#[test]
fn recovery_does_not_leave_a_freelist_that_claims_live_pages() {
use kernel::store::{Config, Store};
let d = tempfile::TempDir::new().unwrap();
{
let mut s = Store::create(d.path(), Config::default()).unwrap();
for round in 0..4u64 {
for i in 0..2_000u64 {
s.put(format!("row-{i:05}").as_bytes(),
format!("round {round} value of row {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
assert!(s.pool_ref().free_pages_pending() > 0, "the churn must have freed pages");
}
assert!(d.path().join("free").exists(), "and persisted them");
let old_freelist = std::fs::read(d.path().join("free")).unwrap();
kernel::recover::recover(d.path(), Config::default()).unwrap();
let rebuilt_freelist = std::fs::read(d.path().join("free")).unwrap();
assert_ne!(rebuilt_freelist, old_freelist,
"recovery must not leave the old file's page-number list beside the rebuild");
let after_recovery = Store::open(d.path(), Config::default()).unwrap();
assert_eq!(after_recovery.pool_ref().free_pages_pending(), 0,
"a rebuilt file starts with no reusable pages");
drop(after_recovery);
std::fs::write(d.path().join("free"), old_freelist).unwrap();
let mut s = Store::open(d.path(), Config::default()).unwrap();
assert_eq!(s.pool_ref().free_pages_pending(), 0,
"a pre-recovery sidecar must never survive renumbering logically or physically");
for round in 0..4u64 {
for i in 0..2_000u64 {
s.put(format!("after-{i:05}").as_bytes(),
format!("post-recovery round {round} row {i}").as_bytes()).unwrap();
}
s.commit().unwrap();
s.checkpoint().unwrap();
}
for i in 0..2_000u64 {
let got = s.get(format!("row-{i:05}").as_bytes()).unwrap();
assert_eq!(got.as_deref(), Some(format!("round 3 value of row {i}").as_bytes()),
"row {i} was lost or overwritten after recovery recycled pages");
}
}