mod common;
use std::fs;
use std::path::Path;
use std::sync::Arc;
use std::thread;
use regolith::{
Db, DurabilityMode, Error, IngestOptions, Options, Range, SstFileWriter, WriteBatch,
WriteOptions,
};
use tempfile::TempDir;
use common::fault::{file_len, find_ssts, first_sst, overwrite_range};
#[path = "lifecycle/options.rs"]
mod options;
#[test]
fn crash_child() {
common::fault::child_entrypoint(common::fault::builtin_workload);
}
fn opts() -> Options {
Options {
write_buffer_size: 4 * 1024,
..Options::default()
}
}
fn key(i: usize) -> Vec<u8> {
format!("key_{i:06}").into_bytes()
}
fn value(i: usize) -> Vec<u8> {
format!("val_{i:06}").into_bytes()
}
fn write_range(db: &Db, from: usize, to: usize) {
for i in from..to {
db.put(&key(i), &value(i)).unwrap();
}
}
fn assert_range_present(db: &Db, from: usize, to: usize, ctx: &str) {
for i in from..to {
assert_eq!(
db.get(&key(i)).unwrap(),
Some(value(i)),
"{ctx}: key {i} is missing or wrong"
);
}
assert_eq!(
db.scan(None, None).unwrap().len(),
to - from,
"{ctx}: the database holds keys it was never given, or lost some"
);
}
fn assert_closed<T: std::fmt::Debug>(what: &str, result: regolith::Result<T>) {
match result {
Err(Error::Closed) => {}
other => panic!("{what}: expected Error::Closed after close, got {other:?}"),
}
}
fn assert_read_only<T: std::fmt::Debug>(what: &str, result: regolith::Result<T>) {
match result {
Err(Error::ReadOnly) => {}
other => panic!("{what}: expected Error::ReadOnly, got {other:?}"),
}
}
fn sweep_every_mutation(db: &Db, scratch: &Path, expect: fn(&str, regolith::Result<()>)) {
let cf = db.default_cf();
let wo = WriteOptions::new();
let batch = || {
let mut b = WriteBatch::new();
b.put(b"k", b"v");
b
};
expect("put", db.put(b"k", b"v"));
expect("put_opt", db.put_opt(&WriteOptions::sync(), b"k", b"v"));
expect("delete", db.delete(&key(0)));
expect("delete_opt", db.delete_opt(&wo, &key(0)));
expect("merge", db.merge(b"k", b"op"));
expect("merge_opt", db.merge_opt(&wo, b"k", b"op"));
expect("delete_range", db.delete_range(&key(0), &key(9)));
expect("write", db.write(batch()));
expect("write_opt", db.write_opt(&wo, batch()));
let sync = DurabilityMode::Immediate;
expect(
"write_with_durability",
db.write_with_durability(batch(), sync),
);
expect("compact_range", db.compact_range(None, None));
expect("drop_all", db.drop_all());
expect("checkpoint", db.checkpoint(scratch.join("cp")));
expect(
"create_column_family",
db.create_column_family("late").map(drop),
);
expect("drop_column_family", db.drop_column_family(cf.clone()));
expect("put_cf", db.put_cf(&cf, b"k", b"v"));
expect("delete_cf", db.delete_cf(&cf, b"k"));
expect("merge_cf", db.merge_cf(&cf, b"k", b"op"));
expect("delete_range_cf", db.delete_range_cf(&cf, &key(0), &key(9)));
let ingest = db.ingest_external_files(&[], IngestOptions::default());
expect("ingest_external_files", ingest);
}
#[test]
fn a_hundred_open_close_reopen_cycles_accumulate_every_write() {
const CYCLES: usize = 100;
const PER_CYCLE: usize = 10;
let dir = TempDir::new().unwrap();
for cycle in 0..CYCLES {
let db = Db::open(dir.path(), opts()).unwrap();
let already = cycle * PER_CYCLE;
assert_eq!(
db.scan(None, None).unwrap().len(),
already,
"cycle {cycle}: reopen did not recover exactly the {already} keys written so far"
);
write_range(&db, already, already + PER_CYCLE);
if cycle % 2 == 0 {
db.close().unwrap();
}
drop(db);
}
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, CYCLES * PER_CYCLE, "final reopen");
}
#[test]
fn close_is_idempotent_and_a_dropped_handle_persists_the_same_data() {
let closed = TempDir::new().unwrap();
let db = Db::open(closed.path(), opts()).unwrap();
write_range(&db, 0, 300);
db.close().unwrap();
db.close().unwrap();
db.close().unwrap();
drop(db);
let dropped = TempDir::new().unwrap();
let db = Db::open(dropped.path(), opts()).unwrap();
write_range(&db, 0, 300);
drop(db);
for (label, dir) in [("closed", &closed), ("dropped", &dropped)] {
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 300, label);
db.close().unwrap();
}
}
#[test]
fn a_second_read_write_open_of_the_same_directory_is_refused() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 50);
let err = Db::open(dir.path(), opts()).unwrap_err();
let msg = err.to_string();
assert!(
msg.contains("locked"),
"a second read-write open must name the lock conflict, got: {msg}"
);
write_range(&db, 50, 100);
assert_range_present(&db, 0, 100, "holder after a refused second open");
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 100, "after the holder released the lock");
}
#[test]
fn eight_concurrent_opens_of_a_held_directory_are_all_refused() {
let dir = TempDir::new().unwrap();
let holder = Db::open(dir.path(), opts()).unwrap();
write_range(&holder, 0, 100);
let path: Arc<Path> = Arc::from(dir.path());
let racers: Vec<_> = (0..8)
.map(|_| {
let path = Arc::clone(&path);
thread::spawn(move || Db::open(&*path, opts()).map(|_| ()))
})
.collect();
for (i, racer) in racers.into_iter().enumerate() {
let result = racer.join().expect("racing opener panicked");
assert!(
result.is_err(),
"racer {i} opened a directory that was already held read-write"
);
}
assert_range_present(&holder, 0, 100, "holder after eight refused opens");
holder.close().unwrap();
drop(holder);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 100, "after the race");
}
#[test]
fn read_only_and_read_write_handles_exclude_each_other_but_readers_share() {
let dir = TempDir::new().unwrap();
let writer = Db::open(dir.path(), opts()).unwrap();
write_range(&writer, 0, 200);
writer.close().unwrap();
drop(writer);
let writer = Db::open(dir.path(), opts()).unwrap();
let err = Db::open_read_only(dir.path(), opts()).unwrap_err();
assert!(
err.to_string().contains("lock"),
"a read-only open under a live writer must name the lock conflict, got: {err}"
);
writer.close().unwrap();
drop(writer);
let reader_a = Db::open_read_only(dir.path(), opts()).unwrap();
let reader_b = Db::open_read_only(dir.path(), opts()).unwrap();
assert_range_present(&reader_a, 0, 200, "first read-only handle");
assert_range_present(&reader_b, 0, 200, "second read-only handle");
let err = Db::open(dir.path(), opts()).unwrap_err();
assert!(
err.to_string().contains("lock"),
"a read-write open under a live reader must name the lock conflict, got: {err}"
);
drop(reader_a);
drop(reader_b);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 200, "after every reader went away");
}
#[test]
fn a_read_only_handle_refuses_every_mutation_and_still_reads() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 100);
db.close().unwrap();
drop(db);
let ro = Db::open_read_only(dir.path(), opts()).unwrap();
sweep_every_mutation(&ro, dir.path(), assert_read_only);
assert_range_present(&ro, 0, 100, "read-only handle after refusing every write");
}
#[test]
fn a_read_only_handle_refuses_a_write_that_carries_no_work() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 10);
db.close().unwrap();
drop(db);
let ro = Db::open_read_only(dir.path(), opts()).unwrap();
let cf = ro.default_cf();
assert_read_only("write(empty batch)", ro.write(WriteBatch::new()));
assert_read_only(
"write_opt(empty batch)",
ro.write_opt(&WriteOptions::new(), WriteBatch::new()),
);
assert_read_only(
"write_with_durability(empty batch)",
ro.write_with_durability(WriteBatch::new(), DurabilityMode::Eventual),
);
assert_read_only("delete_range(empty range)", ro.delete_range(b"z", b"a"));
assert_read_only(
"delete_range_cf(empty range)",
ro.delete_range_cf(&cf, b"z", b"a"),
);
}
#[test]
fn every_fallible_method_returns_closed_after_close() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 50);
let cf = db.default_cf();
let snap = db.snapshot();
db.close().unwrap();
sweep_every_mutation(&db, dir.path(), assert_closed);
assert_closed("get", db.get(&key(0)));
assert_closed("multi_get", db.multi_get(&[&key(0), &key(1)]));
assert_closed("scan", db.scan(None, None));
assert_closed("scan_page", db.scan_page(None, None, 4));
assert_closed("get_cf", db.get_cf(&cf, &key(0)));
assert_closed("multi_get_cf", db.multi_get_cf(&cf, &[&key(0)]));
assert_closed("scan_cf", db.scan_cf(&cf, None, None));
assert_closed("scan_page_cf", db.scan_page_cf(&cf, None, None, 4));
let mut it = db.iter();
it.seek_to_first();
assert_closed("iter().status()", it.status());
assert!(!it.valid(), "a closed iterator must not claim a position");
let mut it = db.iter_cf(&cf);
it.seek_to_first();
assert_closed("iter_cf().status()", it.status());
let mut tail = db.iter_tailing();
tail.seek_to_first();
assert_closed("iter_tailing().status()", tail.status());
assert_closed("snapshot.get", snap.get(&key(0)));
assert_closed("snapshot.multi_get", snap.multi_get(&[&key(0)]));
assert_closed("snapshot.scan", snap.scan(None, None));
assert_closed("snapshot.scan_page", snap.scan_page(None, None, 4));
assert_closed("snapshot.get_cf", snap.get_cf(&cf, &key(0)));
let mut snap_it = snap.iter();
snap_it.seek_to_first();
assert_closed("snapshot.iter().status()", snap_it.status());
let _ = db.get_property("regolith.stats");
let _ = db.get_property("regolith.sstables");
let _ = db.get_property("regolith.levelstats");
let _ = db.get_property("regolith.options");
let _ = db.get_int_property("regolith.estimate-num-keys");
let _ = db.get_approximate_sizes(&[Range::new(&key(0), &key(50))]);
let _ = db.get_approximate_memtable_stats(Range::new(&key(0), &key(50)));
let _ = db.list_column_families();
let _ = db.column_family("default");
let _ = db.snapshot();
assert!(!format!("{db:?}").is_empty());
}
#[test]
fn a_snapshot_and_an_owned_iterator_outlive_the_handle_that_made_them() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 300);
let snap = db.snapshot();
let mut owned = db.snapshot().into_owned_iter();
let mut tail = db.iter_tailing();
drop(db);
assert_eq!(snap.get(&key(7)).unwrap(), Some(value(7)));
assert_eq!(snap.scan(None, None).unwrap().len(), 300);
owned.seek_to_first();
let mut walked = 0usize;
while owned.valid() {
assert_eq!(owned.key(), Some(key(walked).as_slice()));
walked += 1;
owned.next();
}
owned.status().unwrap();
assert_eq!(
walked, 300,
"the owned iterator lost entries after Db::drop"
);
tail.seek_to_first();
let mut tailed = 0usize;
while tail.valid() {
tailed += 1;
tail.next();
}
tail.status().unwrap();
assert_eq!(
tailed, 300,
"the tailing iterator lost entries after Db::drop"
);
assert!(
Db::open(dir.path(), opts()).is_err(),
"the lock must survive as long as a snapshot pins the engine"
);
drop(snap);
drop(owned);
drop(tail);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 300, "after every reader went away");
}
#[test]
fn drop_all_empties_the_database_and_the_emptiness_survives_a_reopen() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 500);
db.compact_range(None, None).unwrap();
db.drop_all().unwrap();
assert!(
db.scan(None, None).unwrap().is_empty(),
"drop_all left data behind"
);
assert_eq!(db.get(&key(0)).unwrap(), None);
write_range(&db, 1000, 1100);
for i in 1000..1100 {
assert_eq!(db.get(&key(i)).unwrap(), Some(value(i)));
}
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 1000, 1100, "reopen after drop_all");
for i in 0..500 {
assert_eq!(
db.get(&key(i)).unwrap(),
None,
"key {i} came back from the dead after drop_all + reopen"
);
}
}
#[test]
fn repeated_drop_all_cycles_leave_no_sstables_and_exactly_one_wal() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
for round in 0..10 {
write_range(&db, 0, 500);
db.compact_range(None, None).unwrap();
assert!(
common::count_sst_files(dir.path()) > 0,
"round {round}: nothing was written, so drop_all would have nothing to reclaim"
);
db.drop_all().unwrap();
assert_eq!(
common::count_sst_files(dir.path()),
0,
"round {round}: drop_all left SSTables on disk"
);
assert_eq!(
common::count_wal_files(dir.path()),
1,
"round {round}: drop_all left more than the one live WAL on disk"
);
}
}
#[test]
fn open_creates_a_missing_parent_chain() {
let dir = TempDir::new().unwrap();
let deep = dir.path().join("a").join("b").join("c").join("db");
let db = Db::open(&deep, opts()).unwrap();
write_range(&db, 0, 20);
db.close().unwrap();
drop(db);
assert!(deep.is_dir(), "the leaf directory was not created");
let db = Db::open(&deep, opts()).unwrap();
assert_range_present(&db, 0, 20, "reopen of a deep path");
}
#[cfg(unix)]
#[test]
fn a_symlink_to_a_directory_is_an_alias_for_the_database_behind_it() {
use std::os::unix::fs::symlink;
let dir = TempDir::new().unwrap();
let target = dir.path().join("real");
fs::create_dir(&target).unwrap();
let link = dir.path().join("link");
symlink(&target, &link).unwrap();
let db = Db::open(&link, opts()).unwrap();
write_range(&db, 0, 100);
let err = Db::open(&target, opts()).unwrap_err();
assert!(
err.to_string().contains("locked"),
"opening the target while the symlink is held must hit the same lock, got: {err}"
);
db.close().unwrap();
drop(db);
let db = Db::open(&target, opts()).unwrap();
assert_range_present(&db, 0, 100, "opened through the target after the link");
db.close().unwrap();
drop(db);
let db = Db::open(&link, opts()).unwrap();
assert_range_present(&db, 0, 100, "opened through the link again");
}
#[cfg(unix)]
#[test]
fn a_symlink_that_is_not_a_directory_is_refused_and_creates_nothing() {
use std::os::unix::fs::symlink;
let dir = TempDir::new().unwrap();
let missing = dir.path().join("does-not-exist");
let dangling = dir.path().join("dangling");
symlink(&missing, &dangling).unwrap();
let err = Db::open(&dangling, opts()).unwrap_err();
assert!(
matches!(err, Error::Io(_)),
"a dangling symlink must fail as an I/O error, got {err:?}"
);
assert!(
!missing.exists(),
"the refused open materialized the symlink's target at {}",
missing.display()
);
let a_file = dir.path().join("a_file");
fs::write(&a_file, b"not a database").unwrap();
let to_file = dir.path().join("to_file");
symlink(&a_file, &to_file).unwrap();
let err = Db::open(&to_file, opts()).unwrap_err();
assert!(
matches!(err, Error::Io(_)),
"a symlink to a regular file must fail as an I/O error, got {err:?}"
);
assert_eq!(
fs::read(&a_file).unwrap(),
b"not a database".to_vec(),
"the refused open rewrote the file the symlink pointed at"
);
assert!(
!dir.path().join("a_file").join("sst").exists(),
"the refused open created database subdirectories"
);
}
#[test]
fn opening_a_directory_with_unrelated_files_leaves_them_untouched() {
let dir = TempDir::new().unwrap();
fs::create_dir_all(dir.path().join("sst")).unwrap();
fs::create_dir_all(dir.path().join("wal")).unwrap();
let strangers = [
(dir.path().join("README.txt"), "top level"),
(dir.path().join("sst").join("notes.txt"), "inside sst/"),
(dir.path().join("wal").join("notes.txt"), "inside wal/"),
];
for (path, body) in &strangers {
fs::write(path, body.as_bytes()).unwrap();
}
fs::create_dir(dir.path().join("subdir")).unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 500);
db.compact_range(None, None).unwrap();
db.drop_all().unwrap();
write_range(&db, 0, 100);
db.close().unwrap();
drop(db);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 100, "reopen alongside unrelated files");
db.close().unwrap();
for (path, body) in &strangers {
assert_eq!(
fs::read_to_string(path).unwrap_or_default(),
*body,
"regolith disturbed an unrelated file at {}",
path.display()
);
}
assert!(
dir.path().join("subdir").is_dir(),
"regolith removed an unrelated subdirectory"
);
}
fn sst_version_byte_offset(path: &Path) -> u64 {
file_len(path) - 8
}
fn read_byte(path: &Path, offset: u64) -> u8 {
let bytes = fs::read(path).unwrap();
bytes[offset as usize]
}
fn sst_format_versions(db_dir: &Path) -> Vec<u8> {
let mut versions: Vec<u8> = find_ssts(db_dir)
.iter()
.map(|sst| read_byte(sst, sst_version_byte_offset(sst)))
.collect();
versions.sort_unstable();
versions.dedup();
versions
}
#[test]
fn an_sstable_from_an_unknown_format_version_is_rejected_with_a_clear_error() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 500);
db.compact_range(None, None).unwrap();
db.close().unwrap();
drop(db);
let sst = first_sst(dir.path());
let offset = sst_version_byte_offset(&sst);
let known = read_byte(&sst, offset);
assert!(
known == 5 || known == 6,
"expected a version regolith writes today (5 = flat index, 6 = partitioned, \
both under the REGOSST magic), found {known} at offset {offset} of {} - \
this test targets the wrong byte",
sst.display()
);
for future_version in [0x07u8, 0x7F, 0xFF] {
overwrite_range(&sst, offset, &[future_version]);
match Db::open(dir.path(), opts()) {
Ok(_) => panic!(
"regolith opened a database whose SSTable claims format version \
{future_version}, which it does not implement"
),
Err(Error::Corruption(source)) => {
let msg = source.to_string();
assert!(
msg.contains("magic"),
"the rejection must name the magic/version it did not recognize, got: {msg}"
);
}
Err(other) => panic!(
"expected a corruption error for format version {future_version}, got {other:?}"
),
}
}
overwrite_range(&sst, offset, &[known]);
let db = Db::open(dir.path(), opts()).unwrap();
assert_range_present(&db, 0, 500, "after restoring the format version byte");
}
#[test]
fn ingesting_an_sstable_from_an_unknown_format_version_changes_nothing() {
let dir = TempDir::new().unwrap();
let db = Db::open(dir.path(), opts()).unwrap();
write_range(&db, 0, 200);
db.compact_range(None, None).unwrap();
let build = |path: &Path| {
let mut writer = SstFileWriter::create(path, &Options::default()).unwrap();
for i in 5000..5100 {
writer.put(&key(i), &value(i)).unwrap();
}
writer.finish().unwrap();
};
let clean = dir.path().join("clean.sst");
build(&clean);
db.ingest_external_files(&[clean], IngestOptions::default())
.unwrap();
for i in 5000..5100 {
assert_eq!(db.get(&key(i)).unwrap(), Some(value(i)), "ingested key {i}");
}
let future = dir.path().join("future.sst");
build(&future);
overwrite_range(&future, sst_version_byte_offset(&future), &[0x42]);
let before_files = find_ssts(dir.path()).len();
let before_state = db.scan(None, None).unwrap();
let err = db
.ingest_external_files(&[future], IngestOptions::default())
.unwrap_err();
assert!(
err.to_string().contains("magic"),
"a refused ingest must name the format it did not recognize, got: {err}"
);
assert_eq!(
find_ssts(dir.path()).len(),
before_files,
"a refused ingest left an SSTable behind"
);
assert_eq!(
db.scan(None, None).unwrap(),
before_state,
"a refused ingest changed the visible state"
);
}