Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion fuzz/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

7 changes: 7 additions & 0 deletions fuzz/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,5 +18,12 @@ test = false
doc = false
bench = false

[[bin]]
name = "recovery_model"
path = "fuzz_targets/recovery_model.rs"
test = false
doc = false
bench = false

[workspace]
members = ["."]
152 changes: 152 additions & 0 deletions fuzz/fuzz_targets/recovery_model.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
#![no_main]

use libfuzzer_sys::fuzz_target;
use lkv::{Database, DatabaseOptions};
use std::collections::HashMap;
use std::fs::{self, OpenOptions};
use std::io::{Read, Seek, SeekFrom, Write};

struct Input<'a> {
bytes: &'a [u8],
offset: usize,
}

impl Input<'_> {
fn byte(&mut self) -> u8 {
let byte = self.bytes[self.offset % self.bytes.len()];
self.offset += 1;
byte
}

fn number(&mut self) -> u64 {
u64::from_le_bytes(std::array::from_fn(|_| self.byte()))
}
}

fn apply_batch(db: &mut Database, model: &mut HashMap<Vec<u8>, Vec<u8>>, input: &mut Input<'_>) {
let mut transaction = db.begin_write().expect("fuzz transaction must start");
for _ in 0..1 + input.byte() as usize % 8 {
let key = vec![input.byte() % 32, input.byte()];
if input.byte() & 3 == 0 {
transaction.delete(&key).expect("fuzz delete must be valid");
model.remove(&key);
} else {
let value = (0..input.byte() as usize % 96)
.map(|_| input.byte())
.collect::<Vec<_>>();
transaction
.put(&key, &value)
.expect("fuzz put must be valid");
model.insert(key, value);
}
}
transaction.commit().expect("fuzz commit must succeed");
}

fn read_model(db: &Database) -> HashMap<Vec<u8>, Vec<u8>> {
db.iter()
.expect("recovered database must iterate")
.map(|entry| {
let (key, value) = entry.expect("recovered entry must verify");
(key.to_vec(), value.to_vec())
})
.collect()
}

fuzz_target!(|bytes: &[u8]| {
if bytes.is_empty() {
return;
}
let scratch =
std::env::temp_dir().join(format!("lkv-fuzz-recovery-model-{}", std::process::id()));
if scratch.exists() {
fs::remove_dir_all(&scratch).expect("old fuzz scratch directory must be removable");
}
fs::create_dir(&scratch).expect("fuzz scratch directory must be creatable");
let path = scratch.join("database.lkv");
let mut input = Input { bytes, offset: 0 };
let mut db = Database::create_with_options(
&path,
DatabaseOptions::default().with_overlay_memory_limit(usize::MAX),
)
.expect("fuzz database must be created");
let mut model = HashMap::new();

apply_batch(&mut db, &mut model, &mut input);
if input.byte() & 1 != 0 {
db.compact().expect("fuzz compaction must succeed");
}
let log_start = db.stats().expect("stats must succeed").storage_bytes;
let prefix_model = model.clone();
let mut commits = Vec::new();
for _ in 0..1 + input.byte() as usize % 16 {
apply_batch(&mut db, &mut model, &mut input);
commits.push((
db.stats().expect("stats must succeed").storage_bytes,
model.clone(),
));
}
drop(db);

let full_len = fs::metadata(&path)
.expect("fuzz database metadata must exist")
.len();
match input.byte() % 3 {
0 => {
let cut = log_start + input.number() % (full_len - log_start + 1);
let file = OpenOptions::new()
.write(true)
.open(&path)
.expect("fuzz database must open for truncation");
file.set_len(cut).expect("fuzz truncation must succeed");
file.sync_all().expect("fuzz truncation must sync");
drop(file);

let expected = commits
.iter()
.rev()
.find(|(end, _)| *end <= cut)
.map_or(&prefix_model, |(_, model)| model);
let db = Database::open(&path).expect("durable prefix must recover");
assert_eq!(&read_model(&db), expected);
let expected_len = commits
.iter()
.rev()
.find(|(end, _)| *end <= cut)
.map_or(log_start, |(end, _)| *end);
assert_eq!(
db.stats().expect("stats must succeed").storage_bytes,
expected_len
);
}
1 => {
let offset = log_start + input.number() % (full_len - log_start);
let mut file = OpenOptions::new()
.read(true)
.write(true)
.open(&path)
.expect("fuzz database must open for corruption");
file.seek(SeekFrom::Start(offset))
.expect("corruption seek must succeed");
let mut byte = [0];
file.read_exact(&mut byte)
.expect("corruption read must succeed");
byte[0] ^= 1 << (input.byte() & 7);
file.seek(SeekFrom::Start(offset))
.expect("corruption seek must succeed");
file.write_all(&byte)
.expect("corruption write must succeed");
file.sync_all().expect("corruption must sync");
drop(file);
assert!(
Database::open(&path).is_err(),
"committed single-bit corruption must be rejected"
);
}
_ => {
let db = Database::open(&path).expect("clean fuzz database must reopen");
assert_eq!(read_model(&db), model);
}
}
fs::remove_dir_all(&scratch).expect("fuzz scratch directory must be removable");
});
37 changes: 25 additions & 12 deletions lkv/src/database/maintenance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ impl Database {
if let Err(error) = write_compact_marker(&mut self.storage) {
return self.rollback_or_poison(rollback_offset, error);
}
failpoints::crash_process_if_requested("after_compact_marker_write");
if let Err(error) = self.storage.seek(SeekFrom::Start(source_offset)) {
return self.rollback_or_poison(rollback_offset, error.into());
}
Expand All @@ -83,6 +84,7 @@ impl Database {
Error::other("base size changed while compacting"),
);
}
failpoints::crash_process_if_requested("after_compact_base_write");
if let Err(error) = self.storage.sync_data() {
return self.rollback_or_poison(rollback_offset, error);
}
Expand All @@ -97,9 +99,12 @@ impl Database {
source_end,
written.metadata_checksum,
);
if let Err(error) = superblock::write(&mut self.storage, source_superblock)
.and_then(|()| self.storage.sync_all())
{
if let Err(error) = superblock::write(&mut self.storage, source_superblock) {
self.state = HandleState::WritePoisoned;
return Err(error);
}
failpoints::crash_process_if_requested("after_compact_superblock_write");
if let Err(error) = self.storage.sync_all() {
self.state = HandleState::WritePoisoned;
return Err(error);
}
Expand All @@ -121,8 +126,10 @@ impl Database {
"base size changed while relocating compaction output",
));
}
failpoints::crash_process_if_requested("after_compact_destination_base_write");
self.storage.seek(SeekFrom::Start(destination_end))?;
write_compact_marker(&mut self.storage)?;
failpoints::crash_process_if_requested("after_compact_relocation_marker_write");
self.storage.sync_data()?;
failpoints::crash_process_if_requested("after_compact_relocation_sync");

Expand All @@ -135,16 +142,22 @@ impl Database {
destination_end,
destination_written.metadata_checksum,
);
if let Err(error) = superblock::write(&mut self.storage, destination_superblock)
.and_then(|()| self.storage.sync_all())
{
if let Err(error) = superblock::write(&mut self.storage, destination_superblock) {
self.state = HandleState::WritePoisoned;
return Err(error);
}
failpoints::crash_process_if_requested("after_compact_destination_superblock_write");
if let Err(error) = self.storage.sync_all() {
self.state = HandleState::WritePoisoned;
return Err(error);
}
failpoints::crash_process_if_requested("after_compact_destination_superblock_sync");
if let Err(error) = superblock::write_redundant(&mut self.storage, destination_superblock)
.and_then(|()| self.storage.sync_all())
{
if let Err(error) = superblock::write_redundant(&mut self.storage, destination_superblock) {
self.state = HandleState::WritePoisoned;
return Err(error);
}
failpoints::crash_process_if_requested("after_compact_redundant_superblock_write");
if let Err(error) = self.storage.sync_all() {
self.state = HandleState::WritePoisoned;
return Err(error);
}
Expand All @@ -155,9 +168,9 @@ impl Database {
// so replacing our final Base mapping is sufficient on every supported OS.
self.state = HandleState::Unavailable;
self.base = detached;
self.storage
.set_len(destination_end)
.and_then(|()| self.storage.sync_all())?;
self.storage.set_len(destination_end)?;
failpoints::crash_process_if_requested("after_compact_truncate");
self.storage.sync_all()?;
failpoints::crash_process_if_requested("after_compact_truncate_sync");
if let Err(error) = self.install_superblock(destination_superblock) {
self.state = HandleState::Unavailable;
Expand Down
5 changes: 3 additions & 2 deletions lkv/src/database/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -357,6 +357,7 @@ impl Database {
} else {
self.overlay.install_tail(staged, mutations);
}
failpoints::crash_process_if_requested("after_batch_publish");
Ok(())
}

Expand Down Expand Up @@ -697,8 +698,8 @@ mod failpoints {
pub fn crash_process_if_requested(point: &str) {
if std::env::var_os("LKV_TEST_CRASH_POINT").as_deref() == Some(std::ffi::OsStr::new(point))
{
// Deliberately skip destructors to model a process disappearing between
// commit phases. Exit code 86 distinguishes the injected crash.
// Deliberately skip destructors to model a process disappearing between commit phases.
// Exit code 86 distinguishes the injected crash.
std::process::exit(86);
}
}
Expand Down
6 changes: 6 additions & 0 deletions lkv/src/database/tests/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,12 @@ fn temp_dir() -> PathBuf {
path
}

fn copy_test_database(source: &Path) -> Result<PathBuf> {
let target = temp_path();
fs::copy(source, &target)?;
Ok(target)
}

fn remove_test_database(path: &Path) -> Result<()> {
let parent = path
.parent()
Expand Down
Loading