use std::collections::BTreeMap;
use std::fs;
use std::time::{Duration, Instant};
use regolith::{Db, Error, Options, WriteBatch};
use tempfile::TempDir;
mod common;
#[test]
fn crash_child() {
common::fault::child_entrypoint(common::fault::builtin_workload);
}
fn peak_rss_kib() -> Option<u64> {
let status = fs::read_to_string("/proc/self/status").ok()?;
status
.lines()
.find_map(|line| line.strip_prefix("VmHWM:"))
.and_then(|rest| rest.split_whitespace().next())
.and_then(|n| n.parse().ok())
}
fn reset_peak_rss() -> bool {
fs::write("/proc/self/clear_refs", b"5\n").is_ok()
}
fn measured<T>(label: &str, f: impl FnOnce() -> T) -> T {
let reset = reset_peak_rss();
let started = Instant::now();
let out = f();
let elapsed = started.elapsed();
match (peak_rss_kib(), reset) {
(Some(kib), true) => println!(
"[resource_limits] {label}: peak RSS {kib} KiB, {elapsed:.2?} \
(counter reset before the workload; only isolated under --test-threads=1)"
),
(Some(kib), false) => println!(
"[resource_limits] {label}: peak RSS {kib} KiB, {elapsed:.2?} \
(process-wide since start - the peak counter could not be reset)"
),
(None, _) => println!(
"[resource_limits] {label}: peak RSS not measured \
(this kernel does not expose /proc/self/status VmHWM), {elapsed:.2?}"
),
}
out
}
fn opts_with(write_buffer_size: usize) -> Options {
Options {
write_buffer_size,
..Options::default()
}
}
fn expect_invalid_argument(result: regolith::Result<()>, what: &str) {
match result {
Err(Error::InvalidArgument(message)) => {
assert!(
message.contains(what),
"the rejection must name the limit it enforced, got {message:?}"
);
}
Err(other) => panic!("expected InvalidArgument for {what}, got {other:?}"),
Ok(()) => panic!("{what} must be rejected, but the write was accepted"),
}
}
fn seeded_bytes(seed: u64, len: usize) -> Vec<u8> {
let mut state = seed | 1;
let mut out = Vec::with_capacity(len);
while out.len() < len {
state ^= state << 13;
state ^= state >> 7;
state ^= state << 17;
out.extend_from_slice(&state.to_le_bytes());
}
out.truncate(len);
out
}
fn scan_in_order(db: &Db) -> BTreeMap<Vec<u8>, Vec<u8>> {
let mut out = BTreeMap::new();
let mut previous: Option<Vec<u8>> = None;
let mut iter = db.iter();
iter.seek_to_first();
while iter.valid() {
let key = iter.key().expect("a valid iterator has a key").to_vec();
let value = iter.value().expect("a valid iterator has a value").to_vec();
if let Some(prev) = &previous {
assert!(
prev.as_slice() < key.as_slice(),
"scan produced keys out of order: {:?} then {:?}",
&prev[..prev.len().min(16)],
&key[..key.len().min(16)]
);
}
previous = Some(key.clone());
out.insert(key, value);
iter.next();
}
iter.status().expect("the scan must not end in an error");
out
}
fn wait_until(deadline: Duration, mut check: impl FnMut() -> bool) -> bool {
let started = Instant::now();
let mut backoff = Duration::from_millis(1);
loop {
if check() {
return true;
}
if started.elapsed() >= deadline {
return false;
}
std::thread::sleep(backoff);
backoff = (backoff * 2).min(Duration::from_millis(20));
}
}
fn deepest_populated_level(db: &Db) -> Option<usize> {
(0..7).rev().find(|level| {
db.get_int_property(&format!("regolith.num-files-at-level{level}"))
.unwrap_or(0)
> 0
})
}
#[test]
fn a_zero_length_key_and_a_zero_length_value_survive_flush_and_compaction() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts_with(4 * 1024)).unwrap();
db.put(b"", b"").unwrap();
db.put(b"a", b"").unwrap();
db.put(b"b", b"nonempty").unwrap();
assert_eq!(db.get(b"").unwrap(), Some(Vec::new()));
assert_eq!(db.get(b"absent").unwrap(), None);
db.compact_range(None, None).unwrap();
assert_eq!(db.get(b"").unwrap(), Some(Vec::new()));
assert_eq!(db.get(b"a").unwrap(), Some(Vec::new()));
let scanned = scan_in_order(&db);
assert_eq!(
scanned.keys().next().map(|k| k.as_slice()),
Some(&b""[..]),
"the empty key must sort first"
);
assert_eq!(scanned.len(), 3);
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts_with(4 * 1024)).unwrap();
assert_eq!(db.get(b"").unwrap(), Some(Vec::new()));
assert_eq!(db.get(b"b").unwrap(), Some(b"nonempty".to_vec()));
db.delete(b"").unwrap();
db.compact_range(None, None).unwrap();
assert_eq!(
db.get(b"").unwrap(),
None,
"deleting the empty key must actually hide it"
);
assert_eq!(scan_in_order(&db).len(), 2);
}
#[test]
fn a_write_batch_with_zero_operations_is_a_successful_no_op() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts_with(4 * 1024)).unwrap();
db.put(b"k", b"v").unwrap();
let before = scan_in_order(&db);
let memtable_before = db.get_int_property("regolith.cur-size-all-mem-tables");
let snapshot = db.snapshot();
let empty = WriteBatch::new();
assert_eq!(empty.len(), 0);
assert!(empty.is_empty());
db.write(empty).unwrap();
assert_eq!(scan_in_order(&db), before, "an empty batch changed state");
assert_eq!(
db.get_int_property("regolith.cur-size-all-mem-tables"),
memtable_before,
"an empty batch put bytes into the memtable"
);
assert_eq!(snapshot.get(b"k").unwrap(), Some(b"v".to_vec()));
db.close().unwrap();
drop(snapshot);
drop(db);
let db = Db::open(dir.path(), opts_with(4 * 1024)).unwrap();
assert_eq!(
scan_in_order(&db),
before,
"an empty batch left something behind in the WAL"
);
}
#[test]
fn delete_range_covers_exactly_zero_one_and_all_keys() {
let keys: Vec<&[u8]> = vec![b"a", b"b", b"c", b"d", b"e"];
let fresh = |dir: &TempDir| {
let db = Db::open(dir.path(), opts_with(4 * 1024)).unwrap();
for key in &keys {
db.put(key, b"v").unwrap();
}
db
};
let dir = TempDir::new().unwrap();
let db = fresh(&dir);
db.delete_range(b"m", b"z").unwrap();
assert_eq!(scan_in_order(&db).len(), 5);
db.compact_range(None, None).unwrap();
assert_eq!(scan_in_order(&db).len(), 5, "an empty range deleted data");
let dir = TempDir::new().unwrap();
let db = fresh(&dir);
db.delete_range(b"c", b"d").unwrap();
let after = scan_in_order(&db);
assert_eq!(after.len(), 4);
assert!(!after.contains_key(b"c".as_slice()));
assert!(
after.contains_key(b"d".as_slice()),
"the exclusive upper bound must survive"
);
db.compact_range(None, None).unwrap();
assert_eq!(scan_in_order(&db), after);
let dir = TempDir::new().unwrap();
let db = fresh(&dir);
db.delete_range(b"", b"\xff").unwrap();
assert!(scan_in_order(&db).is_empty());
db.compact_range(None, None).unwrap();
assert!(scan_in_order(&db).is_empty());
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts_with(4 * 1024)).unwrap();
assert!(
scan_in_order(&db).is_empty(),
"a full-range delete did not survive reopen"
);
}
#[test]
fn the_value_size_limit_is_exact_on_the_point_and_batch_paths() {
const LIMIT: usize = 4096;
let dir = TempDir::new().unwrap();
let db = Db::open(
dir.path(),
Options {
max_value_size: LIMIT,
write_buffer_size: 64 * 1024,
..Options::default()
},
)
.unwrap();
let at_limit = seeded_bytes(0x5123_A11E, LIMIT);
let over_limit = seeded_bytes(0x5123_A11E, LIMIT + 1);
db.put(b"exact", &at_limit).unwrap();
assert_eq!(db.get(b"exact").unwrap(), Some(at_limit.clone()));
expect_invalid_argument(db.put(b"over", &over_limit), "max_value_size");
expect_invalid_argument(
db.put_opt(®olith::WriteOptions::sync(), b"over", &over_limit),
"max_value_size",
);
let mut batch = WriteBatch::new();
batch.put(b"batch_exact", &at_limit);
db.write(batch).unwrap();
assert_eq!(db.get(b"batch_exact").unwrap(), Some(at_limit));
let mut batch = WriteBatch::new();
batch.put(b"batch_over", &over_limit);
expect_invalid_argument(db.write(batch), "max_value_size");
assert_eq!(
db.get(b"batch_over").unwrap(),
None,
"a rejected batch must not apply any of its operations"
);
}
#[test]
fn the_key_size_limit_is_exact_on_the_point_and_batch_paths() {
const LIMIT: usize = 64;
let dir = TempDir::new().unwrap();
let db = Db::open(
dir.path(),
Options {
max_key_size: LIMIT,
write_buffer_size: 64 * 1024,
..Options::default()
},
)
.unwrap();
let at_limit = vec![b'k'; LIMIT];
let over_limit = vec![b'k'; LIMIT + 1];
db.put(&at_limit, b"v").unwrap();
assert_eq!(db.get(&at_limit).unwrap(), Some(b"v".to_vec()));
db.delete(&at_limit).unwrap();
db.put(&at_limit, b"v").unwrap();
db.compact_range(Some(&at_limit), None).unwrap();
expect_invalid_argument(db.put(&over_limit, b"v"), "max_key_size");
expect_invalid_argument(db.delete(&over_limit), "max_key_size");
expect_invalid_argument(db.delete_range(&over_limit, b"\xff"), "max_key_size");
expect_invalid_argument(db.delete_range(b"", &over_limit), "max_key_size");
expect_invalid_argument(db.compact_range(Some(&over_limit), None), "max_key_size");
let mut batch = WriteBatch::new();
batch.put(&at_limit, b"batched");
batch.delete_range(&at_limit, b"\xff");
db.write(batch).unwrap();
let mut batch = WriteBatch::new();
batch.put(&over_limit, b"v");
expect_invalid_argument(db.write(batch), "max_key_size");
let mut batch = WriteBatch::new();
batch.delete(&over_limit);
expect_invalid_argument(db.write(batch), "max_key_size");
let mut batch = WriteBatch::new();
batch.delete_range(&over_limit, b"\xff");
expect_invalid_argument(db.write(batch), "max_key_size");
assert_eq!(
db.get(&at_limit).unwrap(),
None,
"the accepted batch's own range delete should have removed the key"
);
}
#[test]
fn one_byte_over_the_default_limits_is_rejected() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), Options::default()).unwrap();
let over_value = vec![0u8; regolith::DEFAULT_MAX_VALUE_SIZE + 1];
expect_invalid_argument(db.put(b"k", &over_value), "max_value_size");
drop(over_value);
let over_key = vec![b'k'; regolith::DEFAULT_MAX_KEY_SIZE + 1];
expect_invalid_argument(db.put(&over_key, b"v"), "max_key_size");
}
#[test]
fn the_default_size_maxima_round_trip_at_exactly_the_limit() {
measured("64 MiB value + 8 MiB key at the limit", || {
let dir = TempDir::new().unwrap();
let value = seeded_bytes(0x0BAD_F00D, regolith::DEFAULT_MAX_VALUE_SIZE);
let key = {
let mut k = seeded_bytes(0x00C0_FFEE, regolith::DEFAULT_MAX_KEY_SIZE);
k[0] = b'z';
k
};
let db = Db::open(dir.path(), Options::default()).unwrap();
db.put(b"biggest", &value).unwrap();
db.put(&key, b"long key").unwrap();
assert_eq!(db.get(b"biggest").unwrap().as_ref(), Some(&value));
assert_eq!(db.get(&key).unwrap(), Some(b"long key".to_vec()));
db.compact_range(None, None).unwrap();
assert_eq!(db.get(b"biggest").unwrap().as_ref(), Some(&value));
assert_eq!(db.get(&key).unwrap(), Some(b"long key".to_vec()));
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), Options::default()).unwrap();
assert_eq!(
db.get(b"biggest").unwrap().as_ref(),
Some(&value),
"a value at exactly max_value_size did not survive reopen"
);
assert_eq!(db.get(&key).unwrap(), Some(b"long key".to_vec()));
});
}
#[test]
fn one_key_overwritten_100_000_times_collapses_to_a_single_version() {
measured("100 000 overwrites of one key", || {
const WRITES: usize = 100_000;
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts_with(64 * 1024)).unwrap();
for i in 0..WRITES {
db.put(b"hot", format!("v{i:07}").as_bytes()).unwrap();
}
let last = format!("v{:07}", WRITES - 1).into_bytes();
assert_eq!(db.get(b"hot").unwrap(), Some(last.clone()));
db.compact_range(None, None).unwrap();
let live = scan_in_order(&db);
assert_eq!(live.len(), 1, "compaction kept more than the live version");
assert_eq!(live.get(b"hot".as_slice()), Some(&last));
let on_disk = db
.get_int_property("regolith.total-sst-files-size")
.expect("regolith.total-sst-files-size is a supported property");
assert!(
on_disk < 64 * 1024,
"one live 11-byte value occupies {on_disk} bytes on disk; \
shadowed versions were not collapsed"
);
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts_with(64 * 1024)).unwrap();
assert_eq!(db.get(b"hot").unwrap(), Some(last));
assert_eq!(scan_in_order(&db).len(), 1);
});
}
#[test]
fn one_hundred_thousand_distinct_keys_scan_in_order() {
measured("100 000 distinct keys, full ordered scan", || {
const KEYS: usize = 100_000;
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts_with(256 * 1024)).unwrap();
for i in 0..KEYS {
db.put(
format!("key_{i:07}").as_bytes(),
format!("value_{i:07}").as_bytes(),
)
.unwrap();
}
let mut seen = 0usize;
let mut previous: Option<Vec<u8>> = None;
let mut iter = db.iter();
iter.seek_to_first();
while iter.valid() {
let key = iter.key().unwrap().to_vec();
let value = iter.value().unwrap().to_vec();
assert_eq!(
value,
format!("value_{seen:07}").into_bytes(),
"entry {seen} carries the wrong value"
);
assert_eq!(key, format!("key_{seen:07}").into_bytes());
if let Some(prev) = &previous {
assert!(prev < &key, "scan went backwards at entry {seen}");
}
previous = Some(key);
seen += 1;
iter.next();
}
iter.status().unwrap();
assert_eq!(seen, KEYS, "the scan lost or invented keys");
db.compact_range(None, None).unwrap();
assert_eq!(scan_in_order(&db).len(), KEYS);
});
}
#[test]
fn a_write_batch_of_100_000_operations_applies_atomically() {
measured("100 000-operation WriteBatch", || {
const OPS: usize = 100_000;
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts_with(256 * 1024)).unwrap();
db.put(b"sentinel", b"before").unwrap();
let before = db.snapshot();
let mut batch = WriteBatch::new();
for i in 0..OPS {
batch.put(format!("b_{i:07}").as_bytes(), b"batched");
}
for i in (0..OPS).step_by(10) {
batch.delete(format!("b_{i:07}").as_bytes());
}
assert_eq!(batch.len(), OPS + OPS.div_ceil(10));
db.write(batch).unwrap();
assert_eq!(
before.get(b"b_0000001").unwrap(),
None,
"a snapshot taken before the batch observed part of it"
);
let expected_live = OPS - OPS.div_ceil(10);
let live = scan_in_order(&db);
assert_eq!(
live.len(),
expected_live + 1,
"sentinel plus surviving keys"
);
assert_eq!(live.get(b"b_0000000".as_slice()), None);
assert_eq!(
live.get(b"b_0000001".as_slice()),
Some(&b"batched".to_vec())
);
db.close().unwrap();
drop(before);
drop(db);
let db = Db::open(dir.path(), opts_with(256 * 1024)).unwrap();
assert_eq!(
scan_in_order(&db).len(),
expected_live + 1,
"the batch did not survive reopen intact"
);
});
}
#[test]
fn megabyte_keys_mixed_with_short_keys_survive_prefix_compression() {
measured("1 MiB keys mixed with short keys", || {
const MIB: usize = 1024 * 1024;
let filler = vec![b'x'; MIB];
let mut expected: BTreeMap<Vec<u8>, Vec<u8>> = BTreeMap::new();
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts_with(4 * MIB)).unwrap();
for i in 0..6u32 {
let mut long = format!("p{i:02}_").into_bytes();
long.extend_from_slice(&filler);
let short = format!("p{i:02}_short").into_bytes();
expected.insert(long.clone(), format!("long{i}").into_bytes());
expected.insert(short.clone(), format!("short{i}").into_bytes());
db.put(&long, format!("long{i}").as_bytes()).unwrap();
db.put(&short, format!("short{i}").as_bytes()).unwrap();
}
for j in 0..4u32 {
let mut key = b"q_".to_vec();
key.extend_from_slice(&filler);
key.extend_from_slice(format!("{j:04}").as_bytes());
expected.insert(key.clone(), format!("shared{j}").into_bytes());
db.put(&key, format!("shared{j}").as_bytes()).unwrap();
}
let verify = |db: &Db, stage: &str| {
for (key, value) in &expected {
assert_eq!(
db.get(key).unwrap().as_ref(),
Some(value),
"{stage}: a {}-byte key lost its value",
key.len()
);
}
assert_eq!(scan_in_order(db), expected, "{stage}: scan mismatch");
let target = expected.keys().last().unwrap();
let mut iter = db.iter();
iter.seek(target);
assert_eq!(iter.key(), Some(target.as_slice()), "{stage}: seek missed");
iter.seek_for_prev(target);
assert_eq!(
iter.key(),
Some(target.as_slice()),
"{stage}: seek_for_prev missed"
);
};
verify(&db, "in memory");
db.compact_range(None, None).unwrap();
verify(&db, "after compaction");
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts_with(4 * MIB)).unwrap();
verify(&db, "after reopen");
});
}
#[test]
fn a_deep_level_structure_still_answers_every_read() {
measured("cascade to L6", || {
const KEYS: usize = 6_000;
let dir = TempDir::new().unwrap();
let opts = Options {
write_buffer_size: 8 * 1024,
block_size: 1024,
target_file_size: 8 * 1024,
level_base_bytes: 1024,
level_size_multiplier: 2,
l0_compaction_trigger: 2,
..Options::default()
};
let db = Db::open(dir.path(), opts.clone()).unwrap();
let mut expected = BTreeMap::new();
for i in 0..KEYS {
let key = format!("k{i:06}").into_bytes();
let value = format!("value_for_{i:06}").into_bytes();
db.put(&key, &value).unwrap();
expected.insert(key, value);
}
assert_eq!(
scan_in_order(&db),
expected,
"a scan during background compaction lost data"
);
let deep = wait_until(Duration::from_secs(120), || {
deepest_populated_level(&db).is_some_and(|level| level >= 6)
});
assert!(
deep,
"the tree never reached L6 within 120 s; levels now:\n{}",
db.get_property("regolith.levelstats").unwrap_or_default()
);
let levels: Vec<(usize, u64)> = (0..7)
.map(|l| {
(
l,
db.get_int_property(&format!("regolith.num-files-at-level{l}"))
.unwrap_or(0),
)
})
.filter(|(_, n)| *n > 0)
.collect();
println!("[resource_limits] populated levels (level, files): {levels:?}");
assert!(
deepest_populated_level(&db).is_some_and(|level| level >= 3),
"the deepest populated level regressed above L3"
);
assert!(
levels.iter().filter(|(l, _)| *l >= 1).count() >= 3,
"the reads never had to cross a nested tree: only {levels:?} are populated"
);
for (key, value) in &expected {
assert_eq!(db.get(key).unwrap().as_ref(), Some(value));
}
assert_eq!(scan_in_order(&db), expected);
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts).unwrap();
assert_eq!(
scan_in_order(&db),
expected,
"the deep tree did not survive reopen"
);
});
}
#[path = "resource_limits/enospc.rs"]
mod enospc;