use rand::RngExt;
use rand::prelude::SliceRandom;
#[cfg(feature = "experimental-api-5")]
use redb::KeyRange;
use redb::backends::{FileBackend, InMemoryBackend};
use redb::{
AccessGuard, Builder, CommitError, CompactionError, Database, Durability, Key, MultimapRange,
MultimapTableDefinition, MultimapValue, Range, ReadableDatabase, ReadableTable,
ReadableTableMetadata, SetDurabilityError, StorageBackend, TableDefinition, TableStats,
TransactionError, Value, WriteTransaction,
};
use redb::{DatabaseError, ReadableMultimapTable, SavepointError, StorageError, TableError};
use std::borrow::Borrow;
use std::fs;
use std::io::{ErrorKind, Write};
use std::marker::PhantomData;
#[cfg(not(feature = "experimental-api-5"))]
use std::ops::RangeBounds;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex, RwLock, mpsc};
use std::thread;
use std::time::Duration;
const ELEMENTS: usize = 100;
const PAGE_REUSE_VALUE_LEN: usize = 512 * 1024;
const MIN_REUSED_VALUE_PAGES: u64 = 64;
const MAX_PAGE_REUSE_METADATA_GROWTH: u64 = 16;
const SLICE_TABLE: TableDefinition<&[u8], &[u8]> = TableDefinition::new("slice");
const SLICE_TABLE2: TableDefinition<&[u8], &[u8]> = TableDefinition::new("slice2");
const STR_TABLE: TableDefinition<&str, &str> = TableDefinition::new("x");
const U64_TABLE: TableDefinition<u64, u64> = TableDefinition::new("u64");
fn create_tempfile() -> tempfile::NamedTempFile {
if cfg!(target_os = "wasi") {
tempfile::NamedTempFile::new_in("/tmp").unwrap()
} else {
tempfile::NamedTempFile::new().unwrap()
}
}
#[derive(Debug)]
struct BlockingSyncState {
block_next: AtomicBool,
blocked: Mutex<bool>,
blocked_cvar: Condvar,
release: Mutex<bool>,
release_cvar: Condvar,
}
impl BlockingSyncState {
fn new() -> Self {
Self {
block_next: AtomicBool::new(false),
blocked: Mutex::new(false),
blocked_cvar: Condvar::new(),
release: Mutex::new(false),
release_cvar: Condvar::new(),
}
}
fn block_next_sync(&self) {
*self.blocked.lock().unwrap() = false;
*self.release.lock().unwrap() = false;
self.block_next.store(true, Ordering::SeqCst);
}
fn wait_until_blocked(&self) {
let mut blocked = self.blocked.lock().unwrap();
while !*blocked {
let (next, timeout) = self
.blocked_cvar
.wait_timeout(blocked, Duration::from_secs(10))
.unwrap();
blocked = next;
assert!(!timeout.timed_out(), "timed out waiting for sync_data()");
}
}
fn release(&self) {
*self.release.lock().unwrap() = true;
self.release_cvar.notify_all();
}
fn maybe_block(&self) {
if !self.block_next.swap(false, Ordering::SeqCst) {
return;
}
*self.blocked.lock().unwrap() = true;
self.blocked_cvar.notify_all();
let mut release = self.release.lock().unwrap();
while !*release {
release = self.release_cvar.wait(release).unwrap();
}
}
}
#[derive(Debug)]
struct BlockingSyncBackend {
inner: FileBackend,
state: Arc<BlockingSyncState>,
}
impl StorageBackend for BlockingSyncBackend {
fn len(&self) -> Result<u64, std::io::Error> {
self.inner.len()
}
fn read(&self, offset: u64, out: &mut [u8]) -> Result<(), std::io::Error> {
self.inner.read(offset, out)
}
fn set_len(&self, len: u64) -> Result<(), std::io::Error> {
self.inner.set_len(len)
}
fn sync_data(&self) -> Result<(), std::io::Error> {
self.state.maybe_block();
self.inner.sync_data()
}
fn write(&self, offset: u64, data: &[u8]) -> Result<(), std::io::Error> {
self.inner.write(offset, data)
}
fn close(&self) -> Result<(), std::io::Error> {
self.inner.close()
}
}
#[derive(Clone, Debug, Default)]
struct SharedInMemoryBackend {
inner: Arc<RwLock<Vec<u8>>>,
}
impl StorageBackend for SharedInMemoryBackend {
fn len(&self) -> Result<u64, std::io::Error> {
Ok(self.inner.read().unwrap().len() as u64)
}
fn read(&self, offset: u64, out: &mut [u8]) -> Result<(), std::io::Error> {
let offset = usize::try_from(offset).unwrap();
let end = offset + out.len();
let guard = self.inner.read().unwrap();
if end > guard.len() {
return Err(std::io::Error::from(ErrorKind::UnexpectedEof));
}
out.copy_from_slice(&guard[offset..end]);
Ok(())
}
fn set_len(&self, len: u64) -> Result<(), std::io::Error> {
self.inner
.write()
.unwrap()
.resize(len.try_into().unwrap(), 0);
Ok(())
}
fn sync_data(&self) -> Result<(), std::io::Error> {
Ok(())
}
fn write(&self, offset: u64, data: &[u8]) -> Result<(), std::io::Error> {
let offset = usize::try_from(offset).unwrap();
let end = offset + data.len();
let mut guard = self.inner.write().unwrap();
if end > guard.len() {
return Err(std::io::Error::from(ErrorKind::UnexpectedEof));
}
guard[offset..end].copy_from_slice(data);
Ok(())
}
}
#[derive(Debug)]
struct WriteCountingBackend {
inner: InMemoryBackend,
writes: Arc<AtomicU64>,
}
impl StorageBackend for WriteCountingBackend {
fn len(&self) -> Result<u64, std::io::Error> {
self.inner.len()
}
fn read(&self, offset: u64, out: &mut [u8]) -> Result<(), std::io::Error> {
self.inner.read(offset, out)
}
fn set_len(&self, len: u64) -> Result<(), std::io::Error> {
self.inner.set_len(len)
}
fn sync_data(&self) -> Result<(), std::io::Error> {
self.inner.sync_data()
}
fn write(&self, offset: u64, data: &[u8]) -> Result<(), std::io::Error> {
self.writes.fetch_add(1, Ordering::SeqCst);
self.inner.write(offset, data)
}
}
fn random_data(count: usize, key_size: usize, value_size: usize) -> Vec<(Vec<u8>, Vec<u8>)> {
let mut pairs = vec![];
for _ in 0..count {
let key: Vec<u8> = (0..key_size).map(|_| rand::rng().random()).collect();
let value: Vec<u8> = (0..value_size).map(|_| rand::rng().random()).collect();
pairs.push((key, value));
}
pairs
}
#[derive(Debug)]
struct FailingBackend {
inner: FileBackend,
fail_reads: Arc<AtomicBool>,
fail_syncs: Arc<AtomicBool>,
}
impl FailingBackend {
fn new(backend: FileBackend) -> Self {
Self {
inner: backend,
fail_reads: Arc::new(AtomicBool::new(false)),
fail_syncs: Arc::new(AtomicBool::new(false)),
}
}
}
impl StorageBackend for FailingBackend {
fn len(&self) -> Result<u64, std::io::Error> {
self.inner.len()
}
fn read(&self, offset: u64, out: &mut [u8]) -> Result<(), std::io::Error> {
if self.fail_reads.load(Ordering::SeqCst) {
return Err(std::io::Error::from(ErrorKind::Other));
}
self.inner.read(offset, out)
}
fn set_len(&self, len: u64) -> Result<(), std::io::Error> {
self.inner.set_len(len)
}
fn sync_data(&self) -> Result<(), std::io::Error> {
if self.fail_syncs.load(Ordering::SeqCst) {
return Err(std::io::Error::from(ErrorKind::Other));
}
self.inner.sync_data()
}
fn write(&self, offset: u64, data: &[u8]) -> Result<(), std::io::Error> {
self.inner.write(offset, data)
}
}
#[test]
fn previous_io_error() {
let tmpfile = create_tempfile();
let backend = FailingBackend::new(FileBackend::new(tmpfile.into_file()).unwrap());
let fail_syncs = backend.fail_syncs.clone();
let db = Database::builder().create_with_backend(backend).unwrap();
fail_syncs.store(true, Ordering::SeqCst);
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&0, &0).unwrap();
}
assert!(txn.commit().is_err());
assert!(matches!(
db.begin_write().err().unwrap(),
TransactionError::Storage(StorageError::PreviousIo)
));
}
#[test]
fn read_past_eof_errors() {
let mut tmpfile = create_tempfile();
tmpfile.write_all(&[1u8; 8]).unwrap();
tmpfile.flush().unwrap();
let backend = FileBackend::new(tmpfile.into_file()).unwrap();
let mut buf = [0u8; 16];
assert_eq!(
backend.read(64, &mut buf).unwrap_err().kind(),
ErrorKind::UnexpectedEof
);
assert_eq!(
backend.read(0, &mut buf).unwrap_err().kind(),
ErrorKind::UnexpectedEof
);
}
#[test]
fn extract_if_error_latches() {
let tmpfile = create_tempfile();
let backend = FailingBackend::new(FileBackend::new(tmpfile.into_file()).unwrap());
let fail_reads = backend.fail_reads.clone();
let db = Builder::new()
.set_cache_size(0)
.create_with_backend(backend)
.unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..10_000u64 {
table.insert(&i, &i).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
let mut iter = table.extract_if(|_, _| true).unwrap();
assert!(iter.next().unwrap().is_ok());
fail_reads.store(true, Ordering::SeqCst);
loop {
match iter.next() {
Some(Ok(_)) => {}
Some(Err(_)) => break,
None => panic!("iterator must not report exhaustion"),
}
}
assert!(matches!(iter.next(), Some(Err(StorageError::PreviousIo))));
assert!(matches!(
iter.next_back(),
Some(Err(StorageError::PreviousIo))
));
fail_reads.store(false, Ordering::SeqCst);
}
drop(txn);
}
#[test]
fn extract_if_error_commit_reports_storage_failure() {
let tmpfile = create_tempfile();
let backend = FailingBackend::new(FileBackend::new(tmpfile.into_file()).unwrap());
let fail_reads = backend.fail_reads.clone();
let db = Builder::new()
.set_cache_size(0)
.create_with_backend(backend)
.unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..1000u64 {
table.insert(&i, &i).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
fail_reads.store(true, Ordering::SeqCst);
let mut iter = table.extract_if(|_, _| false).unwrap();
assert!(iter.next().unwrap().is_err());
fail_reads.store(false, Ordering::SeqCst);
}
let err = txn.commit().unwrap_err();
assert!(matches!(
err,
CommitError::Storage(StorageError::PreviousIo)
));
}
#[test]
fn close_called_exactly_once_on_shutdown_io_error() {
use redb::backends::InMemoryBackend;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::atomic::{AtomicI64, AtomicU64};
#[derive(Debug)]
struct FaultBackend {
inner: InMemoryBackend,
ops: Arc<AtomicU64>,
countdown: Arc<AtomicI64>,
close_calls: Arc<AtomicU64>,
}
impl FaultBackend {
fn new() -> (Self, Arc<AtomicU64>, Arc<AtomicI64>, Arc<AtomicU64>) {
let ops = Arc::new(AtomicU64::new(0));
let countdown = Arc::new(AtomicI64::new(i64::MAX));
let close_calls = Arc::new(AtomicU64::new(0));
(
Self {
inner: InMemoryBackend::new(),
ops: ops.clone(),
countdown: countdown.clone(),
close_calls: close_calls.clone(),
},
ops,
countdown,
close_calls,
)
}
fn should_fail(&self) -> bool {
self.ops.fetch_add(1, Ordering::SeqCst);
self.countdown.fetch_sub(1, Ordering::SeqCst) == 1
}
}
impl StorageBackend for FaultBackend {
fn len(&self) -> Result<u64, std::io::Error> {
self.inner.len()
}
fn read(&self, offset: u64, out: &mut [u8]) -> Result<(), std::io::Error> {
self.inner.read(offset, out)
}
fn set_len(&self, len: u64) -> Result<(), std::io::Error> {
if self.should_fail() {
return Err(std::io::Error::other("injected fault"));
}
self.inner.set_len(len)
}
fn sync_data(&self) -> Result<(), std::io::Error> {
if self.should_fail() {
return Err(std::io::Error::other("injected fault"));
}
self.inner.sync_data()
}
fn write(&self, offset: u64, data: &[u8]) -> Result<(), std::io::Error> {
if self.should_fail() {
return Err(std::io::Error::other("injected fault"));
}
self.inner.write(offset, data)
}
fn close(&self) -> Result<(), std::io::Error> {
self.close_calls.fetch_add(1, Ordering::SeqCst);
Ok(())
}
}
fn workload(db: &Database) {
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&0, &0).unwrap();
}
txn.commit().unwrap();
}
let (backend, ops, _countdown, close_calls) = FaultBackend::new();
let db = Database::builder().create_with_backend(backend).unwrap();
workload(&db);
ops.store(0, Ordering::SeqCst);
drop(db);
let shutdown_ops = ops.load(Ordering::SeqCst);
assert_eq!(close_calls.load(Ordering::SeqCst), 1);
assert!(shutdown_ops > 0);
for fail_at in 0..shutdown_ops {
let (backend, ops, countdown, close_calls) = FaultBackend::new();
let db = Database::builder().create_with_backend(backend).unwrap();
workload(&db);
ops.store(0, Ordering::SeqCst);
countdown.store(i64::try_from(fail_at).unwrap() + 1, Ordering::SeqCst);
catch_unwind(AssertUnwindSafe(|| drop(db)))
.unwrap_or_else(|_| panic!("shutdown fail_at={fail_at}: panic during Database::drop"));
assert_eq!(
close_calls.load(Ordering::SeqCst),
1,
"shutdown fail_at={fail_at}: close() call count"
);
}
}
#[test]
fn close_called_exactly_once_when_open_fails() {
use std::sync::atomic::AtomicU64;
#[derive(Debug)]
struct CountingBackend {
data: Vec<u8>,
len_fails: bool,
close_calls: Arc<AtomicU64>,
}
impl StorageBackend for CountingBackend {
fn len(&self) -> Result<u64, std::io::Error> {
if self.len_fails {
return Err(std::io::Error::other("injected fault"));
}
Ok(self.data.len() as u64)
}
fn read(&self, offset: u64, out: &mut [u8]) -> Result<(), std::io::Error> {
let start = usize::try_from(offset).unwrap();
out.copy_from_slice(&self.data[start..start + out.len()]);
Ok(())
}
fn set_len(&self, _len: u64) -> Result<(), std::io::Error> {
Ok(())
}
fn sync_data(&self) -> Result<(), std::io::Error> {
Ok(())
}
fn write(&self, _offset: u64, _data: &[u8]) -> Result<(), std::io::Error> {
Ok(())
}
fn close(&self) -> Result<(), std::io::Error> {
self.close_calls.fetch_add(1, Ordering::SeqCst);
Ok(())
}
}
let close_calls = Arc::new(AtomicU64::new(0));
let backend = CountingBackend {
data: vec![0xAB; 4096],
len_fails: false,
close_calls: close_calls.clone(),
};
assert!(Database::builder().create_with_backend(backend).is_err());
assert_eq!(close_calls.load(Ordering::SeqCst), 1);
let close_calls = Arc::new(AtomicU64::new(0));
let backend = CountingBackend {
data: Vec::new(),
len_fails: true,
close_calls: close_calls.clone(),
};
assert!(Database::builder().create_with_backend(backend).is_err());
assert_eq!(close_calls.load(Ordering::SeqCst), 1);
}
#[test]
fn mixed_durable_commit() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&0, &0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
}
#[test]
fn non_durable_commit_persistence() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
let pairs = random_data(100, 16, 20);
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..ELEMENTS {
let (key, value) = &pairs[i % pairs.len()];
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
drop(db);
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(SLICE_TABLE).unwrap();
let mut key_order: Vec<usize> = (0..ELEMENTS).collect();
key_order.shuffle(&mut rand::rng());
{
for i in &key_order {
let (key, value) = &pairs[*i % pairs.len()];
assert_eq!(table.get(key.as_slice()).unwrap().unwrap().value(), value);
}
}
}
#[test]
fn non_durable_commit_issues_no_backend_writes() {
let writes = Arc::new(AtomicU64::new(0));
let backend = WriteCountingBackend {
inner: InMemoryBackend::new(),
writes: writes.clone(),
};
let db = Database::builder().create_with_backend(backend).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&0, &0).unwrap();
}
txn.commit().unwrap();
let writes_before = writes.load(Ordering::SeqCst);
for i in 1..=10u64 {
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&i, &(i * 10)).unwrap();
}
txn.commit().unwrap();
}
assert_eq!(writes.load(Ordering::SeqCst), writes_before);
{
let read = db.begin_read().unwrap();
let table = read.open_table(U64_TABLE).unwrap();
for i in 1..=10u64 {
assert_eq!(table.get(&i).unwrap().unwrap().value(), i * 10);
}
}
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&100, &100).unwrap();
}
txn.commit().unwrap();
assert!(writes.load(Ordering::SeqCst) > writes_before);
}
#[test]
fn durable_overwrite_does_not_write_post_commit_metadata() {
const COMMITS: u64 = 20;
let writes = Arc::new(AtomicU64::new(0));
let backend = WriteCountingBackend {
inner: InMemoryBackend::new(),
writes: writes.clone(),
};
let db = Database::builder().create_with_backend(backend).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&0, &0).unwrap();
}
txn.commit().unwrap();
let writes_before = writes.load(Ordering::SeqCst);
for value in 1..=COMMITS {
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&0, &value).unwrap();
}
txn.commit().unwrap();
}
let commit_writes = writes.load(Ordering::SeqCst) - writes_before;
assert!(
commit_writes <= 7 * COMMITS,
"{COMMITS} durable overwrites issued {commit_writes} backend writes"
);
}
fn test_persistence(durability: Durability) {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
let pairs = random_data(100, 16, 20);
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..ELEMENTS {
let (key, value) = &pairs[i % pairs.len()];
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
drop(db);
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(SLICE_TABLE).unwrap();
let mut key_order: Vec<usize> = (0..ELEMENTS).collect();
key_order.shuffle(&mut rand::rng());
{
for i in &key_order {
let (key, value) = &pairs[*i % pairs.len()];
assert_eq!(table.get(key.as_slice()).unwrap().unwrap().value(), value);
}
}
}
#[test]
fn immediate_persistence() {
test_persistence(Durability::Immediate);
}
#[test]
fn immediate_free() {
test_free(Durability::Immediate);
}
#[test]
fn nondurable_free() {
test_free(Durability::None);
}
#[test]
fn page_reuse() {
test_page_reuse(false);
test_page_reuse(true);
}
#[test]
fn page_reuse_after_persistent_savepoint_delete() {
test_page_reuse_after_persistent_savepoint_delete(false);
test_page_reuse_after_persistent_savepoint_delete(true);
}
#[test]
fn page_reuse_with_racing_reader() {
test_page_reuse_with_racing_reader(false);
test_page_reuse_with_racing_reader(true);
}
#[test]
fn page_reuse_after_unclean_reopen() {
test_page_reuse_after_unclean_reopen(false);
test_page_reuse_after_unclean_reopen(true);
}
#[test]
fn stats_with_open_table() {
const BIG_TABLE: TableDefinition<u64, &[u8]> = TableDefinition::new("big");
const BIG_ELEMENTS: u64 = 2000;
const BIG_VALUE_LEN: usize = 16000;
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..5000 {
table.insert(i, i).unwrap();
}
}
let mut table = {
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..5000 {
table.remove(i).unwrap();
}
table
};
{
let mut big = txn.open_table(BIG_TABLE).unwrap();
let value = vec![0xEE; BIG_VALUE_LEN];
for i in 0..BIG_ELEMENTS {
big.insert(i, value.as_slice()).unwrap();
}
}
let big_stored = BIG_ELEMENTS * (u64::fixed_width().unwrap() as u64 + BIG_VALUE_LEN as u64);
let stats = txn.stats().unwrap();
assert_eq!(stats.stored_bytes(), big_stored);
table.insert(0, 0).unwrap();
drop(table);
let stats = txn.stats().unwrap();
assert_eq!(
stats.stored_bytes(),
big_stored + 2 * u64::fixed_width().unwrap() as u64
);
txn.abort().unwrap();
}
#[test]
fn stats_with_renamed_open_table() {
const RENAMED_TABLE: TableDefinition<u64, u64> = TableDefinition::new("renamed");
const BIG_TABLE: TableDefinition<u64, &[u8]> = TableDefinition::new("big");
const BIG_ELEMENTS: u64 = 2000;
const BIG_VALUE_LEN: usize = 16000;
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
const COMMITTED_ELEMENTS: u64 = 5000;
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..COMMITTED_ELEMENTS {
table.insert(i, i).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in COMMITTED_ELEMENTS..2 * COMMITTED_ELEMENTS {
table.insert(i, i).unwrap();
}
}
txn.rename_table(U64_TABLE, RENAMED_TABLE).unwrap();
let table = {
let mut table = txn.open_table(RENAMED_TABLE).unwrap();
for i in 0..2 * COMMITTED_ELEMENTS {
table.remove(i).unwrap();
}
table
};
{
let mut big = txn.open_table(BIG_TABLE).unwrap();
let value = vec![0xEE; BIG_VALUE_LEN];
for i in 0..BIG_ELEMENTS {
big.insert(i, value.as_slice()).unwrap();
}
}
let element_stored = 2 * u64::fixed_width().unwrap() as u64;
let big_stored = BIG_ELEMENTS * (u64::fixed_width().unwrap() as u64 + BIG_VALUE_LEN as u64);
let stats = txn.stats().unwrap();
assert_eq!(
stats.stored_bytes(),
big_stored + COMMITTED_ELEMENTS * element_stored
);
drop(table);
let stats = txn.stats().unwrap();
assert_eq!(stats.stored_bytes(), big_stored);
txn.abort().unwrap();
}
#[test]
fn corrupted_persistent_savepoint_record() {
let tmpfile = create_tempfile();
let path = tmpfile.path();
let savepoint_id = {
let db = Database::create(path).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
id
};
let mut data = fs::read(path).unwrap();
let mut offsets = vec![];
for i in 0..data.len() - 18 {
let transaction_id = u64::from_le_bytes(data[i + 9..i + 17].try_into().unwrap());
if data[i] == 3
&& u64::from_le_bytes(data[i + 1..i + 9].try_into().unwrap()) == savepoint_id
&& transaction_id > 0
&& transaction_id < 100
&& data[i + 17] == 1
{
offsets.push(i);
}
}
assert_eq!(offsets.len(), 1);
data[offsets[0]] = 9;
fs::write(path, &data).unwrap();
let result = Database::open(path);
assert!(matches!(
result,
Err(DatabaseError::Storage(StorageError::Corrupted(_)))
));
}
#[test]
fn failed_reopen_preserves_staged_table() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..100 {
table.insert(i, i).unwrap();
}
}
let wrong_type: TableDefinition<&str, &str> = TableDefinition::new("u64");
assert!(txn.open_table(wrong_type).is_err());
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
assert_eq!(txn.open_table(U64_TABLE).unwrap().len().unwrap(), 100);
}
fn begin_page_reuse_write(db: &Database, quick_repair: bool) -> WriteTransaction {
let mut txn = db.begin_write().unwrap();
txn.set_quick_repair(quick_repair);
txn
}
fn make_page_reuse_value(byte: u8) -> Vec<u8> {
vec![byte; PAGE_REUSE_VALUE_LEN]
}
fn test_page_reuse(quick_repair: bool) {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 1).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let allocated_pages = txn.stats().unwrap().allocated_pages();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 2).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let allocated_after_reuse = txn.stats().unwrap().allocated_pages();
assert!(
allocated_after_reuse <= allocated_pages + MAX_PAGE_REUSE_METADATA_GROWTH,
"allocated_pages={allocated_pages}, allocated_after_reuse={allocated_after_reuse}, quick_repair={quick_repair}"
);
}
fn test_page_reuse_after_persistent_savepoint_delete(quick_repair: bool) {
let tmpfile = create_tempfile();
let key = [0u8; 16];
let db = Database::create(tmpfile.path()).unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
let value = make_page_reuse_value(0);
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let savepoint = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let allocated_with_savepoint = txn.stats().unwrap().allocated_pages();
txn.abort().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
assert!(table.remove(key.as_slice()).unwrap().is_some());
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let allocated_pinned = txn.stats().unwrap().allocated_pages();
assert!(
allocated_pinned >= allocated_with_savepoint,
"allocated_with_savepoint={allocated_with_savepoint}, allocated_pinned={allocated_pinned}, quick_repair={quick_repair}"
);
txn.abort().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
assert!(txn.delete_persistent_savepoint(savepoint).unwrap());
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let allocated_after_delete = txn.stats().unwrap().allocated_pages();
assert!(
allocated_with_savepoint > allocated_after_delete + MIN_REUSED_VALUE_PAGES,
"allocated_with_savepoint={allocated_with_savepoint}, allocated_after_delete={allocated_after_delete}, quick_repair={quick_repair}"
);
txn.abort().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
let value = make_page_reuse_value(1);
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
let allocated_after_reinsert = txn.stats().unwrap().allocated_pages();
assert!(
allocated_after_reinsert <= allocated_with_savepoint + MAX_PAGE_REUSE_METADATA_GROWTH,
"allocated_with_savepoint={allocated_with_savepoint}, allocated_after_reinsert={allocated_after_reinsert}, quick_repair={quick_repair}"
);
txn.abort().unwrap();
drop(db);
let db = Database::open(tmpfile.path()).unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(SLICE_TABLE).unwrap();
let value = table.get(key.as_slice()).unwrap().unwrap();
assert_eq!(1, value.value()[0]);
assert_eq!(PAGE_REUSE_VALUE_LEN, value.value().len());
}
fn test_page_reuse_after_unclean_reopen(quick_repair: bool) {
let backend = SharedInMemoryBackend::default();
let db = Database::builder()
.create_with_backend(backend.clone())
.unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
txn.commit().unwrap();
let txn = begin_page_reuse_write(&db, quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 1).unwrap();
}
txn.commit().unwrap();
std::mem::forget(db);
let mut db = Database::builder().create_with_backend(backend).unwrap();
{
let txn = db.begin_read().unwrap();
let table = txn.open_table(U64_TABLE).unwrap();
assert_eq!(1, table.get(&0).unwrap().unwrap().value());
}
assert!(db.check_integrity().unwrap(), "quick_repair={quick_repair}");
}
fn test_page_reuse_with_racing_reader(quick_repair: bool) {
let tmpfile = create_tempfile();
let sync_state = Arc::new(BlockingSyncState::new());
let backend = BlockingSyncBackend {
inner: FileBackend::new(tmpfile.into_file()).unwrap(),
state: sync_state.clone(),
};
let db = Arc::new(Database::builder().create_with_backend(backend).unwrap());
let txn = begin_page_reuse_write(db.as_ref(), quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = begin_page_reuse_write(db.as_ref(), quick_repair);
txn.commit().unwrap();
let txn = begin_page_reuse_write(db.as_ref(), quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 1).unwrap();
}
sync_state.block_next_sync();
let commit_handle = thread::spawn(move || txn.commit().unwrap());
sync_state.wait_until_blocked();
let db_for_read = db.clone();
let (read_started, read_started_rx) = mpsc::channel();
let (check_read, check_read_rx) = mpsc::channel();
let (read_checked, read_checked_rx) = mpsc::channel();
let (drop_read, drop_read_rx) = mpsc::channel();
let read_handle = thread::spawn(move || {
let read_txn = db_for_read.begin_read().unwrap();
let value = {
let table = read_txn.open_table(U64_TABLE).unwrap();
table.get(&0).unwrap().unwrap().value()
};
read_started.send(value).unwrap();
check_read_rx.recv().unwrap();
let value = {
let table = read_txn.open_table(U64_TABLE).unwrap();
table.get(&0).unwrap().unwrap().value()
};
read_checked.send(value).unwrap();
drop_read_rx.recv().unwrap();
drop(read_txn);
});
let read_value = read_started_rx
.recv_timeout(Duration::from_secs(10))
.expect("begin_read() blocked during sync_data()");
assert_eq!(0, read_value, "quick_repair={quick_repair}");
sync_state.release();
commit_handle.join().unwrap();
let txn = begin_page_reuse_write(db.as_ref(), quick_repair);
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 1..1000 {
table.insert(i, i).unwrap();
}
}
txn.commit().unwrap();
check_read.send(()).unwrap();
let read_value = read_checked_rx
.recv_timeout(Duration::from_secs(10))
.expect("read transaction blocked after racing with commit");
assert_eq!(0, read_value, "quick_repair={quick_repair}");
drop_read.send(()).unwrap();
read_handle.join().unwrap();
for _ in 0..3 {
let txn = begin_page_reuse_write(db.as_ref(), quick_repair);
txn.commit().unwrap();
}
}
fn test_free(durability: Durability) {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
{
let _table = txn.open_table(SLICE_TABLE).unwrap();
let mut table = txn.open_table(SLICE_TABLE2).unwrap();
table.insert([].as_slice(), [].as_slice()).unwrap();
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
{
let mut table = txn.open_table(SLICE_TABLE2).unwrap();
table.remove([].as_slice()).unwrap();
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
let allocated_pages = txn.stats().unwrap().allocated_pages();
let key = vec![0; 100];
let value = vec![0u8; 1024];
let target_db_size = 8 * 1024 * 1024;
let num_writes = target_db_size / 10 / (key.len() + value.len());
assert!(num_writes > 64);
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..num_writes {
let mut mut_key = key.clone();
mut_key.extend_from_slice(&(i as u64).to_le_bytes());
table.insert(mut_key.as_slice(), value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
{
let key_range: Vec<usize> = (0..num_writes).collect();
for chunk in key_range.chunks(10) {
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in chunk {
let mut mut_key = key.clone();
mut_key.extend_from_slice(&(*i as u64).to_le_bytes());
table.remove(mut_key.as_slice()).unwrap();
}
}
txn.commit().unwrap();
}
}
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(durability).unwrap();
assert!(txn.stats().unwrap().allocated_pages() <= allocated_pages);
txn.abort().unwrap();
}
#[test]
fn nondurable_live_and_free() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let allocated_pages = txn.stats().unwrap().allocated_pages();
txn.abort().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 1).unwrap();
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
for i in 0..5 {
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, i).unwrap();
}
txn.commit().unwrap();
}
{
let table = read_txn.open_table(U64_TABLE).unwrap();
assert_eq!(table.get(0).unwrap().unwrap().value(), 1);
}
drop(read_txn);
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(0).unwrap();
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
assert!(txn.stats().unwrap().allocated_pages() <= allocated_pages * 2 + 2);
}
#[test]
fn large_values() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
let mut key = vec![0u8; 1024];
let value = vec![0u8; 2_000_000];
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..5 {
key[0] = i;
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..5 {
key[0] = i;
table.remove(key.as_slice()).unwrap();
}
}
txn.commit().unwrap();
}
#[test]
#[cfg(target_pointer_width = "64")]
fn value_too_large() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
let small_value = vec![0u8; 1024];
let too_big_value = vec![0u8; 3 * 1024 * 1024 * 1024 + 1];
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
assert!(matches!(
table.insert(small_value.as_slice(), too_big_value.as_slice()),
Err(StorageError::ValueTooLarge(_))
));
assert!(matches!(
table.insert(too_big_value.as_slice(), small_value.as_slice()),
Err(StorageError::ValueTooLarge(_))
));
assert!(matches!(
table.insert(too_big_value.as_slice(), too_big_value.as_slice()),
Err(StorageError::ValueTooLarge(_))
));
table
.insert(small_value.as_slice(), small_value.as_slice())
.unwrap();
{
let mut guard = table.get_mut(small_value.as_slice()).unwrap().unwrap();
assert!(matches!(
guard.insert(too_big_value.as_slice()),
Err(StorageError::ValueTooLarge(_))
));
}
assert!(matches!(
table
.entry(small_value.as_slice())
.unwrap()
.and_modify(|g| g.insert(too_big_value.as_slice())),
Err(StorageError::ValueTooLarge(_))
));
assert_eq!(
table.get(small_value.as_slice()).unwrap().unwrap().value(),
small_value.as_slice()
);
table.remove(small_value.as_slice()).unwrap();
drop(too_big_value);
let almost_big_value = vec![0u8; 2 * 1024 * 1024 * 1024];
assert!(matches!(
table.insert(almost_big_value.as_slice(), almost_big_value.as_slice()),
Err(StorageError::ValueTooLarge(_))
));
}
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(SLICE_TABLE).unwrap();
assert!(table.is_empty().unwrap());
}
#[test]
fn small_db_is_small_file() {
let tmpfile = create_tempfile();
const TABLE: TableDefinition<u32, u32> = TableDefinition::new("TABLE");
let mut db = Database::create(tmpfile.path()).unwrap();
let wtx = db.begin_write().unwrap();
let mut table = wtx.open_table(TABLE).unwrap();
table.insert(0, 0).unwrap();
drop(table);
wtx.commit().unwrap();
db.compact().unwrap();
drop(db);
let metadata = tmpfile.as_file().metadata().unwrap();
assert!(
metadata.len() < 40 * 1024,
"File size: {:?}",
metadata.len()
);
}
#[test]
fn many_pairs() {
let tmpfile = create_tempfile();
const TABLE: TableDefinition<u32, u32> = TableDefinition::new("TABLE");
let db = Database::create(tmpfile.path()).unwrap();
let wtx = db.begin_write().unwrap();
let mut table = wtx.open_table(TABLE).unwrap();
for i in 0..200_000 {
table.insert(i, i).unwrap();
if i % 10_000 == 0 {
eprintln!("{i}");
}
}
drop(table);
wtx.commit().unwrap();
}
#[test]
fn explicit_close() {
let tmpfile = create_tempfile();
const TABLE: TableDefinition<u32, u32> = TableDefinition::new("TABLE");
let db = Database::create(tmpfile.path()).unwrap();
let wtx = db.begin_write().unwrap();
wtx.open_table(TABLE).unwrap();
wtx.commit().unwrap();
let tx = db.begin_read().unwrap();
let table = tx.open_table(TABLE).unwrap();
assert!(matches!(
tx.close(),
Err(TransactionError::ReadTransactionStillInUse(_))
));
drop(table);
let tx2 = db.begin_read().unwrap();
tx2.close().unwrap();
}
#[test]
fn read_only_get_guard_keeps_transaction_alive() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..100 {
table.insert(i, i).unwrap();
}
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(U64_TABLE).unwrap();
let guard = table.get_owned(&5).unwrap().unwrap();
drop(table);
assert!(matches!(
read_txn.close(),
Err(TransactionError::ReadTransactionStillInUse(_))
));
for _ in 0..10 {
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..100 {
table.insert(i, i + 1).unwrap();
}
}
write_txn.commit().unwrap();
}
assert_eq!(guard.value(), 5);
}
#[test]
fn read_only_range_guards_keep_transaction_alive() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..100 {
table.insert(i, i).unwrap();
}
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(U64_TABLE).unwrap();
let mut range = table.range_owned(5..).unwrap();
let (key, value) = range.next().unwrap().unwrap();
let (back_key, back_value) = range.next_back().unwrap().unwrap();
drop(range);
drop(table);
assert!(matches!(
read_txn.close(),
Err(TransactionError::ReadTransactionStillInUse(_))
));
for _ in 0..10 {
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..100 {
table.insert(i, i + 1).unwrap();
}
}
write_txn.commit().unwrap();
}
assert_eq!(key.value(), 5);
assert_eq!(value.value(), 5);
assert_eq!(back_key.value(), 99);
assert_eq!(back_value.value(), 99);
}
#[test]
fn large_keys() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
let mut key = vec![0u8; 1024];
let value = vec![0u8; 1];
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..100 {
key[0] = i;
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..100 {
key[0] = i;
table.remove(key.as_slice()).unwrap();
}
}
txn.commit().unwrap();
}
#[test]
fn dynamic_growth() {
let tmpfile = create_tempfile();
let table_definition: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let big_value = vec![0u8; 1024];
let expected_size = 10 * 1024 * 1024;
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_definition).unwrap();
table.insert(&0, big_value.as_slice()).unwrap();
}
txn.commit().unwrap();
let initial_file_size = tmpfile.as_file().metadata().unwrap().len();
assert!(initial_file_size < (expected_size / 2) as u64);
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_definition).unwrap();
for i in 0..2048 {
table.insert(&i, big_value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let file_size = tmpfile.as_file().metadata().unwrap().len();
assert!(file_size > initial_file_size);
}
#[test]
fn multi_page_kv() {
let tmpfile = create_tempfile();
let elements = 4;
let page_size = 4096;
let db = Builder::new().create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
let mut key = vec![0u8; page_size + 1];
let mut value = vec![0; page_size + 1];
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..elements {
key[0] = i;
value[0] = i;
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..elements {
key[0] = i;
value[0] = i;
assert_eq!(&value, table.get(key.as_slice()).unwrap().unwrap().value());
}
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..elements {
key[0] = i;
table.remove(key.as_slice()).unwrap();
}
}
txn.commit().unwrap();
}
#[test]
fn regression() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&1, &1).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&6, &9).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&12, &10).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&18, &27).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&24, &33).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&30, &14).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(&30).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(U64_TABLE).unwrap();
let v = table.get(&6).unwrap().unwrap().value();
assert_eq!(v, 9);
}
#[test]
fn regression2() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let tx = db.begin_write().unwrap();
let a_def: TableDefinition<&str, &str> = TableDefinition::new("a");
let b_def: TableDefinition<&str, &str> = TableDefinition::new("b");
let c_def: TableDefinition<&str, &str> = TableDefinition::new("c");
let _c = tx.open_table(c_def).unwrap();
let b = tx.open_table(b_def).unwrap();
let mut a = tx.open_table(a_def).unwrap();
a.insert("hi", "1").unwrap();
assert!(b.get("hi").unwrap().is_none());
}
#[test]
fn regression3() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(SLICE_TABLE).unwrap();
let big_value = vec![0u8; 1000];
for i in 0..20u8 {
t.insert([i].as_slice(), big_value.as_slice()).unwrap();
}
for i in (10..20u8).rev() {
t.remove([i].as_slice()).unwrap();
for j in 0..i {
assert!(t.get([j].as_slice()).unwrap().is_some());
}
}
}
tx.commit().unwrap();
}
#[test]
fn regression7() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let big_value = vec![0u8; 4063];
t.insert(&35723, big_value.as_slice()).unwrap();
t.remove(&145278).unwrap();
t.remove(&145227).unwrap();
}
tx.commit().unwrap();
let mut tx = db.begin_write().unwrap();
tx.set_durability(Durability::None).unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 47];
t.insert(&66469, v.as_slice()).unwrap();
let v = vec![0u8; 2414];
t.insert(&146255, v.as_slice()).unwrap();
let v = vec![0u8; 159];
t.insert(&153701, v.as_slice()).unwrap();
let v = vec![0u8; 1186];
t.insert(&145227, v.as_slice()).unwrap();
let v = vec![0u8; 223];
t.insert(&118749, v.as_slice()).unwrap();
t.remove(&145227).unwrap();
let mut iter = t.range(138763..(138763 + 232359)).unwrap().rev();
assert_eq!(iter.next().unwrap().unwrap().0.value(), 153701);
assert_eq!(iter.next().unwrap().unwrap().0.value(), 146255);
assert!(iter.next().is_none());
}
tx.commit().unwrap();
}
#[test]
fn regression8() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut tx = db.begin_write().unwrap();
tx.set_durability(Durability::None).unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 1186];
t.insert(&145227, v.as_slice()).unwrap();
let v = vec![0u8; 1585];
t.insert(&565922, v.as_slice()).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 2040];
t.insert(&94937, v.as_slice()).unwrap();
let v = vec![0u8; 2058];
t.insert(&130571, v.as_slice()).unwrap();
t.remove(&145227).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 947];
t.insert(&118749, v.as_slice()).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
{
let t = tx.open_table(table_def).unwrap();
let mut iter = t.range(118749..142650).unwrap();
assert_eq!(iter.next().unwrap().unwrap().0.value(), 118749);
assert_eq!(iter.next().unwrap().unwrap().0.value(), 130571);
assert!(iter.next().is_none());
}
tx.commit().unwrap();
}
#[test]
fn regression9() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 118665];
t.insert(&452, v.as_slice()).unwrap();
t.len().unwrap();
}
tx.commit().unwrap();
}
#[test]
fn regression10() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 1043];
t.insert(&118749, v.as_slice()).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 952];
t.insert(&118757, v.as_slice()).unwrap();
}
tx.abort().unwrap();
let tx = db.begin_write().unwrap();
{
let t = tx.open_table(table_def).unwrap();
t.get(&829513).unwrap();
}
tx.abort().unwrap();
}
#[test]
fn regression11() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 1204];
t.insert(&118749, v.as_slice()).unwrap();
let v = vec![0u8; 2062];
t.insert(&153697, v.as_slice()).unwrap();
let v = vec![0u8; 2980];
t.insert(&110557, v.as_slice()).unwrap();
let v = vec![0u8; 1999];
t.insert(&677853, v.as_slice()).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let v = vec![0u8; 691];
t.insert(&103591, v.as_slice()).unwrap();
let v = vec![0u8; 952];
t.insert(&118757, v.as_slice()).unwrap();
}
tx.abort().unwrap();
let tx = db.begin_write().unwrap();
tx.commit().unwrap();
}
#[test]
fn regression12() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, u64> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
t.insert(&0, &0).unwrap();
assert_eq!(t.get(&0).unwrap().unwrap().value(), 0);
drop(t);
let t2 = tx.open_table(table_def).unwrap();
assert_eq!(t2.get(&0).unwrap().unwrap().value(), 0);
}
tx.commit().unwrap();
}
#[test]
fn regression13() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: MultimapTableDefinition<u64, &[u8]> = MultimapTableDefinition::new("x");
let mut tx = db.begin_write().unwrap();
tx.set_durability(Durability::None).unwrap();
{
let mut t = tx.open_multimap_table(table_def).unwrap();
let value = vec![0; 1026];
t.insert(&539717, value.as_slice()).unwrap();
let value = vec![0; 530];
t.insert(&539717, value.as_slice()).unwrap();
}
tx.abort().unwrap();
}
#[test]
fn regression14() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: MultimapTableDefinition<u64, &[u8]> = MultimapTableDefinition::new("x");
let mut tx = db.begin_write().unwrap();
tx.set_durability(Durability::None).unwrap();
{
let mut t = tx.open_multimap_table(table_def).unwrap();
let value = vec![0; 1424];
t.insert(&539749, value.as_slice()).unwrap();
}
tx.commit().unwrap();
let mut tx = db.begin_write().unwrap();
tx.set_durability(Durability::None).unwrap();
{
let mut t = tx.open_multimap_table(table_def).unwrap();
let value = vec![0; 2230];
t.insert(&776971, value.as_slice()).unwrap();
let mut iter = t.range(514043..(514043 + 514043)).unwrap().rev();
{
let (key, mut value_iter) = iter.next().unwrap().unwrap();
assert_eq!(key.value(), 776971);
assert_eq!(value_iter.next().unwrap().unwrap().value(), &[0; 2230]);
}
{
let (key, mut value_iter) = iter.next().unwrap().unwrap();
assert_eq!(key.value(), 539749);
assert_eq!(value_iter.next().unwrap().unwrap().value(), &[0; 1424]);
}
}
tx.abort().unwrap();
}
#[test]
fn regression17() {
let tmpfile = create_tempfile();
let db = Database::builder().create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut tx = db.begin_write().unwrap();
tx.set_durability(Durability::None).unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let value = vec![0; 4578];
t.insert(&671325, value.as_slice()).unwrap();
let mut value = t.insert_reserve(&723904, 2246).unwrap();
value.as_mut().fill(0xFF);
}
tx.abort().unwrap();
}
#[test]
fn regression18() {
let tmpfile = create_tempfile();
let db = Database::builder().create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
let savepoint0 = tx.ephemeral_savepoint().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let mut value = t.insert_reserve(&118749, 817).unwrap();
value.as_mut().fill(0xFF);
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
let savepoint1 = tx.ephemeral_savepoint().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let mut value = t.insert_reserve(&65373, 1807).unwrap();
value.as_mut().fill(0xFF);
}
tx.commit().unwrap();
let mut tx = db.begin_write().unwrap();
let savepoint2 = tx.ephemeral_savepoint().unwrap();
tx.restore_savepoint(&savepoint2).unwrap();
tx.commit().unwrap();
drop(savepoint0);
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let mut value = t.insert_reserve(&118749, 2494).unwrap();
value.as_mut().fill(0xFF);
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
let savepoint4 = tx.ephemeral_savepoint().unwrap();
tx.abort().unwrap();
drop(savepoint1);
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let mut value = t.insert_reserve(&429469, 667).unwrap();
value.as_mut().fill(0xFF);
drop(value);
let mut value = t.insert_reserve(&266845, 1614).unwrap();
value.as_mut().fill(0xFF);
}
tx.commit().unwrap();
let mut tx = db.begin_write().unwrap();
tx.restore_savepoint(&savepoint4).unwrap();
tx.commit().unwrap();
drop(savepoint2);
drop(savepoint4);
}
#[test]
fn regression19() {
let tmpfile = create_tempfile();
let db = Database::builder().create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let value = vec![0xFF; 100];
t.insert(&1, value.as_slice()).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
let savepoint0 = tx.ephemeral_savepoint().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let value = vec![0xFF; 101];
t.insert(&1, value.as_slice()).unwrap();
}
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
{
let mut t = tx.open_table(table_def).unwrap();
let value = vec![0xFF; 102];
t.insert(&1, value.as_slice()).unwrap();
}
tx.commit().unwrap();
let mut tx = db.begin_write().unwrap();
tx.restore_savepoint(&savepoint0).unwrap();
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
tx.open_table(table_def).unwrap();
}
fn savepoint_integrity_value(i: u64, len: usize) -> Vec<u8> {
vec![(i % 250) as u8 + 1; len]
}
#[test]
fn check_integrity_after_persistent_savepoint_delete() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 0..1000u64 {
table
.insert(&k, savepoint_integrity_value(k, 60).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint_id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 2000..2500u64 {
table
.insert(&k, savepoint_integrity_value(k, 60).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 2000..2500u64 {
table.remove(&k).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.delete_persistent_savepoint(savepoint_id).unwrap();
txn.commit().unwrap();
assert!(
db.check_integrity().unwrap(),
"false corruption report after deleting a persistent savepoint"
);
}
#[test]
fn check_integrity_with_persistent_savepoint_then_restore() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 0..1000u64 {
table
.insert(&k, savepoint_integrity_value(k, 60).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint_id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 0..900u64 {
table.remove(&k).unwrap();
}
}
txn.commit().unwrap();
assert!(
db.check_integrity().unwrap(),
"false corruption report while a persistent savepoint exists"
);
let mut txn = db.begin_write().unwrap();
let savepoint = txn.get_persistent_savepoint(savepoint_id).unwrap();
txn.restore_savepoint(&savepoint).unwrap();
drop(savepoint);
txn.commit().unwrap();
let read = db.begin_read().unwrap();
let table = read.open_table(table_def).unwrap();
assert_eq!(table.len().unwrap(), 1000);
drop(table);
drop(read);
let txn = db.begin_write().unwrap();
txn.delete_persistent_savepoint(savepoint_id).unwrap();
txn.commit().unwrap();
assert!(db.check_integrity().unwrap());
}
#[test]
fn check_integrity_after_persistent_savepoint_restore_and_delete() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for i in 0..100u64 {
table
.insert(&i, savepoint_integrity_value(i, 100).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for i in 0..50u64 {
table.remove(&i).unwrap();
}
for i in 100..200u64 {
table
.insert(&i, savepoint_integrity_value(i, 100).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
let savepoint = txn.get_persistent_savepoint(id).unwrap();
txn.restore_savepoint(&savepoint).unwrap();
drop(savepoint);
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
assert!(txn.delete_persistent_savepoint(id).unwrap());
assert!(!txn.delete_persistent_savepoint(id).unwrap());
txn.commit().unwrap();
assert!(db.check_integrity().unwrap());
let txn = db.begin_read().unwrap();
let table = txn.open_table(table_def).unwrap();
assert_eq!(table.len().unwrap(), 100);
for i in 0..100u64 {
assert_eq!(
table.get(&i).unwrap().unwrap().value(),
savepoint_integrity_value(i, 100).as_slice()
);
}
}
#[test]
fn check_integrity_ephemeral_savepoint_drop_racing_commit() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let sync_state = Arc::new(BlockingSyncState::new());
let backend = BlockingSyncBackend {
inner: FileBackend::new(tmpfile.into_file()).unwrap(),
state: sync_state.clone(),
};
let mut db = Database::builder().create_with_backend(backend).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 0..1000u64 {
table
.insert(&k, savepoint_integrity_value(k, 60).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 2000..2500u64 {
table
.insert(&k, savepoint_integrity_value(k, 60).as_slice())
.unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for k in 2000..2500u64 {
table.remove(&k).unwrap();
}
}
sync_state.block_next_sync();
let commit_handle = thread::spawn(move || txn.commit().unwrap());
sync_state.wait_until_blocked();
drop(savepoint);
sync_state.release();
commit_handle.join().unwrap();
assert!(
db.check_integrity().unwrap(),
"false corruption after an ephemeral savepoint was dropped during commit"
);
}
#[test]
fn ephemeral_savepoint_racing_first_open_no_leak() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
table
.insert(&0, savepoint_integrity_value(0, 100).as_slice())
.unwrap();
}
txn.commit().unwrap();
const CREATORS: usize = 6;
for i in 0..600u64 {
let mut txn = db.begin_write().unwrap();
let barrier = std::sync::Barrier::new(CREATORS + 1);
let savepoints = thread::scope(|s| {
let handles: Vec<_> = (0..CREATORS)
.map(|_| {
s.spawn(|| {
barrier.wait();
txn.ephemeral_savepoint()
})
})
.collect();
barrier.wait();
{
let mut table = txn.open_table(table_def).unwrap();
table
.insert(&(1000 + i), savepoint_integrity_value(i, 200).as_slice())
.unwrap();
}
handles
.into_iter()
.map(|h| h.join().unwrap())
.collect::<Vec<_>>()
});
let created: Vec<_> = savepoints.into_iter().flatten().collect();
if let Some(savepoint) = created.first() {
{
let mut table = txn.open_table(table_def).unwrap();
for k in 3000..3100u64 {
table
.insert(&k, savepoint_integrity_value(k, 200).as_slice())
.unwrap();
}
}
txn.restore_savepoint(savepoint).unwrap();
txn.commit().unwrap();
} else {
txn.abort().unwrap();
}
}
assert!(
db.check_integrity().unwrap(),
"ephemeral savepoint racing the first table-open leaked pages"
);
}
#[test]
fn regression20() {
let tmpfile = create_tempfile();
let table_def: MultimapTableDefinition<'static, u128, u128> =
MultimapTableDefinition::new("some-table");
for _ in 0..3 {
let mut db = Database::builder().create(tmpfile.path()).unwrap();
db.check_integrity().unwrap();
let txn = db.begin_write().unwrap();
let mut table = txn.open_multimap_table(table_def).unwrap();
for i in 0..1024 {
table.insert(0, i).unwrap();
}
drop(table);
txn.commit().unwrap();
}
}
#[test]
fn regression21() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let write_tx = db.begin_write().unwrap();
let read_tx = db.begin_read().unwrap();
let mut write_table = write_tx
.open_table::<&str, &str>(TableDefinition::new("example"))
.unwrap();
write_table.insert("example", "example").unwrap();
drop(write_table);
write_tx.commit().unwrap();
assert!(matches!(
db.compact().unwrap_err(),
CompactionError::TransactionInProgress
));
drop(read_tx);
}
#[test]
fn regression22() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let allocated_pages = txn.stats().unwrap().allocated_pages();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(0).unwrap();
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
db.begin_write().unwrap().commit().unwrap();
drop(read_txn);
let txn = db.begin_write().unwrap();
let after = txn.stats().unwrap().allocated_pages();
assert!(
after <= allocated_pages + MAX_PAGE_REUSE_METADATA_GROWTH,
"allocated_pages={allocated_pages}, after={after}"
);
}
#[test]
fn regression23() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
#[allow(unused_must_use)]
{
txn.list_persistent_savepoints().unwrap();
}
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let allocated_pages = txn.stats().unwrap().allocated_pages();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.remove(0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.restore_savepoint(&savepoint).unwrap();
txn.commit().unwrap();
drop(savepoint);
db.begin_write().unwrap().commit().unwrap();
db.begin_write().unwrap().commit().unwrap();
let txn = db.begin_write().unwrap();
assert_eq!(allocated_pages, txn.stats().unwrap().allocated_pages());
}
#[test]
fn regression24() {
let tmpfile = create_tempfile();
let table_def: MultimapTableDefinition<u64, u64> = MultimapTableDefinition::new("x");
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let id = txn.persistent_savepoint().unwrap();
txn.delete_persistent_savepoint(id).unwrap();
#[allow(unused_must_use)]
{
txn.list_persistent_savepoints().unwrap();
}
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
txn.delete_table(U64_TABLE).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let allocated_pages = txn.stats().unwrap().allocated_pages();
{
let mut table = txn.open_multimap_table(table_def).unwrap();
table.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
txn.delete_multimap_table(table_def).unwrap();
}
txn.commit().unwrap();
db.begin_write().unwrap().commit().unwrap();
let txn = db.begin_write().unwrap();
assert_eq!(allocated_pages, txn.stats().unwrap().allocated_pages());
}
#[test]
fn regression25() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u16, (u64, u64, u64, u64)> = TableDefinition::new("issue_1108");
let db = Database::create(tmpfile.path()).unwrap();
for i in 0..2730u16 {
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
for j in 0..24u16 {
let key: u16 = i * 24 + j;
let value = key as u64;
table.insert(key, (value, value, value, value)).unwrap();
}
}
txn.commit().unwrap();
}
}
#[test]
fn regression26() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, (&str, &[u8])> = TableDefinition::new("issue_1117");
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
table.insert(0, ("name", &[0u8][..])).unwrap();
}
txn.commit().unwrap();
{
let txn = db.begin_write().unwrap();
let mut table = txn.open_table(table_def).unwrap();
let mut access = table.get_mut(&0).unwrap().unwrap();
let name = access.value().0.to_string();
let large_value = vec![1u8; 8192];
access.insert((&name[..], large_value.as_slice())).unwrap();
drop(access);
drop(table);
txn.commit().unwrap();
}
}
#[test]
fn regression27() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let def: TableDefinition<&str, &[u8]> = TableDefinition::new("x");
let value = "world";
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(def).unwrap();
let mut reserved = table.insert_reserve("hello", value.len()).unwrap();
reserved.as_mut().copy_from_slice(value.as_bytes());
drop(reserved);
drop(table);
write_txn.commit().unwrap();
}
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(def).unwrap();
assert_eq!(
value.as_bytes(),
table.get("hello").unwrap().unwrap().value()
);
}
#[test]
fn check_integrity_clean() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<'static, u64, u64> = TableDefinition::new("x");
let mut db = Database::builder().create(tmpfile.path()).unwrap();
assert!(db.check_integrity().unwrap());
let txn = db.begin_write().unwrap();
let mut table = txn.open_table(table_def).unwrap();
for i in 0..10 {
table.insert(0, i).unwrap();
}
drop(table);
txn.commit().unwrap();
assert!(db.check_integrity().unwrap());
drop(db);
let mut db = Database::builder().create(tmpfile.path()).unwrap();
assert!(db.check_integrity().unwrap());
drop(db);
let mut db = Database::builder().open(tmpfile.path()).unwrap();
assert!(db.check_integrity().unwrap());
}
#[test]
fn check_integrity_clean_after_delete() {
let tmpfile = create_tempfile();
let table_def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
let mut db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(table_def).unwrap();
for i in 0..1000u64 {
table.insert(&i, [1u8; 20].as_slice()).unwrap();
}
}
write_txn.commit().unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(table_def).unwrap();
assert!(table.remove(&0u64).unwrap().is_some());
}
write_txn.commit().unwrap();
assert!(db.check_integrity().unwrap());
assert!(db.check_integrity().unwrap());
}
#[test]
fn check_integrity_preserves_non_durable_commit() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for k in 0..100u64 {
table.insert(&k, &(k * 10)).unwrap();
}
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(&7777u64, &42u64).unwrap();
table.remove(&0u64).unwrap();
}
txn.commit().unwrap();
{
let read = db.begin_read().unwrap();
let table = read.open_table(U64_TABLE).unwrap();
assert_eq!(table.get(&7777u64).unwrap().unwrap().value(), 42);
assert!(table.get(&0u64).unwrap().is_none());
}
assert!(db.check_integrity().unwrap());
let read = db.begin_read().unwrap();
let table = read.open_table(U64_TABLE).unwrap();
assert_eq!(
table.get(&7777u64).unwrap().map(|v| v.value()),
Some(42),
"check_integrity rolled back a non-durable commit (insert lost)"
);
assert!(
table.get(&0u64).unwrap().is_none(),
"check_integrity rolled back a non-durable commit (remove undone)"
);
drop(table);
drop(read);
drop(db);
let db = Database::open(tmpfile.path()).unwrap();
let read = db.begin_read().unwrap();
let table = read.open_table(U64_TABLE).unwrap();
assert_eq!(table.get(&7777u64).unwrap().unwrap().value(), 42);
assert!(table.get(&0u64).unwrap().is_none());
}
#[test]
fn check_integrity_with_live_savepoint_on_nondurable_commit() {
let table2: TableDefinition<u64, &[u8]> = TableDefinition::new("x2");
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(U64_TABLE).unwrap();
t.insert(0, 0).unwrap();
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut t = txn.open_table(U64_TABLE).unwrap();
for k in 1..100u64 {
t.insert(k, k * 10).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
txn.abort().unwrap();
assert!(matches!(
db.check_integrity(),
Err(DatabaseError::TransactionInProgress)
));
let mut txn = db.begin_write().unwrap();
txn.restore_savepoint(&savepoint).unwrap();
txn.commit().unwrap();
drop(savepoint);
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table2).unwrap();
let big = vec![0xABu8; 4096];
for k in 0..100u64 {
t.insert(k, big.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let read = db.begin_read().unwrap();
let t = read.open_table(U64_TABLE).unwrap();
for k in 1..100u64 {
let v = t.get(k).unwrap();
assert_eq!(v.map(|x| x.value()), Some(k * 10), "key {k} corrupted");
}
drop(t);
drop(read);
assert!(db.check_integrity().unwrap());
}
#[test]
fn multimap_stats() {
let tmpfile = create_tempfile();
let db = Database::builder().create(tmpfile.path()).unwrap();
let table_def: MultimapTableDefinition<u128, u128> = MultimapTableDefinition::new("x");
let mut last_size = 0;
for i in 0..1000 {
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
let mut table = txn.open_multimap_table(table_def).unwrap();
if i == 0 {
assert_eq!(table.stats().unwrap().leaf_pages(), 0);
}
table.insert(0, i).unwrap();
assert!(table.stats().unwrap().leaf_pages() > 0);
drop(table);
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let bytes = txn.stats().unwrap().stored_bytes();
assert!(bytes > last_size, "{i}");
last_size = bytes;
}
let read_txn = db.begin_read().unwrap();
let typed_stats = read_txn
.open_multimap_table(table_def)
.unwrap()
.stats()
.unwrap();
let untyped_stats = read_txn
.open_untyped_multimap_table(table_def)
.unwrap()
.stats()
.unwrap();
assert_eq!(typed_stats.leaf_pages(), untyped_stats.leaf_pages());
assert_eq!(typed_stats.stored_bytes(), untyped_stats.stored_bytes());
}
#[test]
fn no_downgrade_durability_with_savepoint() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let mut tx = db.begin_write().unwrap();
tx.persistent_savepoint().unwrap();
assert!(matches!(
tx.set_durability(Durability::None),
Err(SetDurabilityError::PersistentSavepointModified)
));
assert!(matches!(tx.set_durability(Durability::Immediate), Ok(())));
}
#[test]
fn no_savepoint_resurrection() {
let tmpfile = create_tempfile();
let db = Database::builder()
.set_cache_size(41178283)
.create(tmpfile.path())
.unwrap();
let tx = db.begin_write().unwrap();
let persistent_savepoint = tx.persistent_savepoint().unwrap();
tx.commit().unwrap();
let tx = db.begin_write().unwrap();
let savepoint2 = tx.ephemeral_savepoint().unwrap();
tx.delete_persistent_savepoint(persistent_savepoint)
.unwrap();
tx.commit().unwrap();
let mut tx = db.begin_write().unwrap();
tx.restore_savepoint(&savepoint2).unwrap();
tx.delete_persistent_savepoint(persistent_savepoint)
.unwrap();
tx.commit().unwrap();
}
#[test]
fn non_durable_read_isolation() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let mut write_txn = db.begin_write().unwrap();
write_txn.set_durability(Durability::None).unwrap();
{
let mut table = write_txn.open_table(STR_TABLE).unwrap();
table.insert("hello", "world").unwrap();
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let read_table = read_txn.open_table(STR_TABLE).unwrap();
assert_eq!("world", read_table.get("hello").unwrap().unwrap().value());
let mut write_txn = db.begin_write().unwrap();
write_txn.set_durability(Durability::None).unwrap();
{
let mut table = write_txn.open_table(STR_TABLE).unwrap();
table.remove("hello").unwrap();
table.insert("hello2", "world2").unwrap();
table.insert("hello3", "world3").unwrap();
}
write_txn.commit().unwrap();
let read_txn2 = db.begin_read().unwrap();
let read_table2 = read_txn2.open_table(STR_TABLE).unwrap();
assert!(read_table2.get("hello").unwrap().is_none());
assert_eq!(
"world2",
read_table2.get("hello2").unwrap().unwrap().value()
);
assert_eq!(
"world3",
read_table2.get("hello3").unwrap().unwrap().value()
);
assert_eq!(read_table2.len().unwrap(), 2);
assert_eq!("world", read_table.get("hello").unwrap().unwrap().value());
assert!(read_table.get("hello2").unwrap().is_none());
assert!(read_table.get("hello3").unwrap().is_none());
assert_eq!(read_table.len().unwrap(), 1);
}
#[test]
fn range_query() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..10 {
table.insert(&i, &i).unwrap();
}
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(U64_TABLE).unwrap();
let mut iter = table.range(3..7).unwrap();
for i in 3..7u64 {
let (key, value) = iter.next().unwrap().unwrap();
assert_eq!(i, key.value());
assert_eq!(i, value.value());
}
assert!(iter.next().is_none());
let mut iter = table.range(3..=7).unwrap();
for i in 3..=7u64 {
let (key, value) = iter.next().unwrap().unwrap();
assert_eq!(i, key.value());
assert_eq!(i, value.value());
}
assert!(iter.next().is_none());
let total: u64 = table
.range(1..=3)
.unwrap()
.map(|item| item.unwrap().1.value())
.sum();
assert_eq!(total, 6);
}
#[test]
fn range_query_reversed() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..10 {
table.insert(&i, &i).unwrap();
}
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(U64_TABLE).unwrap();
let mut iter = table.range(3..7).unwrap().rev();
for i in (3..7u64).rev() {
let (key, value) = iter.next().unwrap().unwrap();
assert_eq!(i, key.value());
assert_eq!(i, value.value());
}
assert!(iter.next().is_none());
let mut iter = table.range(3..7).unwrap();
let (key, _) = iter.next().unwrap().unwrap();
assert_eq!(3, key.value());
let mut iter = iter.rev();
let (key, _) = iter.next().unwrap().unwrap();
assert_eq!(6, key.value());
let (key, _) = iter.next().unwrap().unwrap();
assert_eq!(5, key.value());
let mut iter = iter.rev();
let (key, _) = iter.next().unwrap().unwrap();
assert_eq!(4, key.value());
assert!(iter.next().is_none());
}
#[test]
fn range_mixed_direction_no_duplicates_multi_page() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(U64_TABLE).unwrap();
for i in 0..5_000u64 {
table.insert(&i, &i).unwrap();
}
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(U64_TABLE).unwrap();
let mut iter = table.range(100..4_900).unwrap();
let mut seen = vec![false; 5_000];
let mut count = 0;
for _ in 0..100 {
let (key, value) = iter.next_back().unwrap().unwrap();
assert_eq!(key.value(), value.value());
let key = key.value() as usize;
assert!(!seen[key], "duplicate key {key}");
seen[key] = true;
count += 1;
}
for item in iter {
let (key, value) = item.unwrap();
assert_eq!(key.value(), value.value());
let key = key.value() as usize;
assert!(!seen[key], "duplicate key {key}");
seen[key] = true;
count += 1;
}
assert_eq!(count, 4_800);
for (key, was_seen) in seen.iter().enumerate().take(4_900).skip(100) {
assert!(*was_seen, "missing key {key}");
}
}
#[test]
fn range_alternating_direction_meets_at_leaf_boundary() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let definition: TableDefinition<u64, [u8; 200]> = TableDefinition::new("x");
let mut table = write_txn.open_table(definition).unwrap();
for i in 0..5_000u64 {
table.insert(i, &[0u8; 200]).unwrap();
}
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let definition: TableDefinition<u64, [u8; 200]> = TableDefinition::new("x");
let table = read_txn.open_table(definition).unwrap();
let mut iter = table.range(0..5_000).unwrap();
for i in 0..2_500u64 {
let (key, _) = iter.next().unwrap().unwrap();
assert_eq!(key.value(), i);
let (key, _) = iter.next_back().unwrap().unwrap();
assert_eq!(key.value(), 4_999 - i);
}
assert!(iter.next().is_none());
assert!(iter.next_back().is_none());
}
#[test]
fn alias_table() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
let table = write_txn.open_table(STR_TABLE).unwrap();
let result = write_txn.open_table(STR_TABLE);
assert!(matches!(
result.err().unwrap(),
TableError::TableAlreadyOpen(_, _)
));
drop(table);
}
#[test]
fn delete_table() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let y_def: MultimapTableDefinition<&str, &str> = MultimapTableDefinition::new("y");
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(STR_TABLE).unwrap();
table.insert("hello", "world").unwrap();
let mut multitable = write_txn.open_multimap_table(y_def).unwrap();
multitable.insert("hello2", "world2").unwrap();
}
write_txn.commit().unwrap();
let write_txn = db.begin_write().unwrap();
assert!(write_txn.delete_table(STR_TABLE).unwrap());
assert!(!write_txn.delete_table(STR_TABLE).unwrap());
assert!(write_txn.delete_multimap_table(y_def).unwrap());
assert!(!write_txn.delete_multimap_table(y_def).unwrap());
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let result = read_txn.open_table(STR_TABLE);
assert!(result.is_err());
let result = read_txn.open_multimap_table(y_def);
assert!(result.is_err());
}
#[test]
fn delete_all_tables() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let x_def: TableDefinition<&str, &str> = TableDefinition::new("x");
let y_def: TableDefinition<&str, &str> = TableDefinition::new("y");
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(x_def).unwrap();
table.insert("hello", "world").unwrap();
let mut table = write_txn.open_table(y_def).unwrap();
table.insert("hello", "world").unwrap();
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
assert_eq!(2, read_txn.list_tables().unwrap().count());
let write_txn = db.begin_write().unwrap();
for table in write_txn.list_tables().unwrap() {
write_txn.delete_table(table).unwrap();
}
write_txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
assert_eq!(0, read_txn.list_tables().unwrap().count());
}
#[test]
fn dropped_write() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let write_txn = db.begin_write().unwrap();
{
let mut table = write_txn.open_table(STR_TABLE).unwrap();
table.insert("hello", "world").unwrap();
}
drop(write_txn);
let read_txn = db.begin_read().unwrap();
let result = read_txn.open_table(STR_TABLE);
assert!(matches!(result, Err(TableError::TableDoesNotExist(_))));
}
#[test]
fn non_page_size_multiple() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
let key = vec![0u8; 1024];
let value = vec![0u8; 1];
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
table.insert(key.as_slice(), value.as_slice()).unwrap();
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(SLICE_TABLE).unwrap();
assert_eq!(table.len().unwrap(), 1);
}
#[test]
fn does_not_exist() {
let tmpfile = create_tempfile();
fs::remove_file(tmpfile.path()).unwrap();
let result = Database::open(tmpfile.path());
if let Err(DatabaseError::Storage(StorageError::Io(e))) = result {
assert!(matches!(e.kind(), ErrorKind::NotFound));
} else {
panic!();
}
let tmpfile = create_tempfile();
let result = Database::open(tmpfile.path());
if let Err(DatabaseError::Storage(StorageError::Io(e))) = result {
assert!(matches!(e.kind(), ErrorKind::InvalidData));
} else {
panic!();
}
}
#[test]
fn invalid_database_file() {
let mut tmpfile = create_tempfile();
tmpfile.write_all(b"hi").unwrap();
let result = Database::open(tmpfile.path());
if let Err(DatabaseError::Storage(StorageError::Io(e))) = result {
assert!(matches!(e.kind(), ErrorKind::InvalidData));
} else {
panic!();
}
let result = Database::create(tmpfile.path());
if let Err(DatabaseError::Storage(StorageError::Io(e))) = result {
assert!(matches!(e.kind(), ErrorKind::InvalidData));
} else {
panic!();
}
}
#[test]
fn wrong_types() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, u32> = TableDefinition::new("x");
let wrong_definition: TableDefinition<u64, u64> = TableDefinition::new("x");
let txn = db.begin_write().unwrap();
txn.open_table(definition).unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
assert!(matches!(
txn.open_table(wrong_definition),
Err(TableError::TableTypeMismatch { .. })
));
txn.abort().unwrap();
let txn = db.begin_read().unwrap();
txn.open_table(definition).unwrap();
assert!(matches!(
txn.open_table(wrong_definition),
Err(TableError::TableTypeMismatch { .. })
));
}
#[test]
fn tree_balance() {
const EXPECTED_ORDER: usize = 9;
fn expected_height(mut elements: usize) -> u32 {
let mut height = 1;
elements /= 2;
height += 1;
height += (elements as f32).log((EXPECTED_ORDER / 2) as f32) as usize + 1;
height.try_into().unwrap()
}
let tmpfile = create_tempfile();
let num_internal_entries = 2;
let key_size = 410;
let db = Database::builder().create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
let elements = (EXPECTED_ORDER / 2).pow(2) - num_internal_entries;
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in (0..elements).rev() {
let mut key = vec![0u8; key_size];
key[0..8].copy_from_slice(&(i as u64).to_le_bytes());
table.insert(key.as_slice(), b"".as_slice()).unwrap();
}
}
txn.commit().unwrap();
let expected = expected_height(elements + num_internal_entries);
let txn = db.begin_write().unwrap();
let height = txn.stats().unwrap().tree_height();
assert!(height <= expected, "height={height} expected={expected}",);
let reduce_to = EXPECTED_ORDER / 2 - num_internal_entries;
{
let mut table = txn.open_table(SLICE_TABLE).unwrap();
for i in 0..(elements - reduce_to) {
let mut key = vec![0u8; key_size];
key[0..8].copy_from_slice(&(i as u64).to_le_bytes());
table.remove(key.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let expected = expected_height(reduce_to + num_internal_entries);
let txn = db.begin_write().unwrap();
let height = txn.stats().unwrap().tree_height();
txn.abort().unwrap();
assert!(height <= expected, "height={height} expected={expected}",);
}
#[cfg(not(target_os = "wasi"))]
#[test]
fn database_lock() {
let tmpfile = create_tempfile();
let result = Database::create(tmpfile.path());
assert!(result.is_ok());
let result2 = Database::open(tmpfile.path());
assert!(
matches!(result2, Err(DatabaseError::DatabaseAlreadyOpen)),
"{result2:?}",
);
drop(result);
let result = Database::open(tmpfile.path());
assert!(result.is_ok());
}
#[test]
fn persistent_savepoint() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &str> = TableDefinition::new("x");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.insert(&0, "hello").unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint_id = txn.persistent_savepoint().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.remove(&0).unwrap();
}
txn.commit().unwrap();
drop(db);
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
let savepoint = txn.get_persistent_savepoint(savepoint_id).unwrap();
txn.restore_savepoint(&savepoint).unwrap();
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(definition).unwrap();
assert_eq!(table.get(&0).unwrap().unwrap().value(), "hello");
}
#[test]
fn savepoint() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &str> = TableDefinition::new("x");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.insert(&0, "hello").unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.remove(&0).unwrap();
}
txn.commit().unwrap();
let mut txn = db.begin_write().unwrap();
let savepoint2 = txn.ephemeral_savepoint().unwrap();
txn.restore_savepoint(&savepoint).unwrap();
assert!(matches!(
txn.restore_savepoint(&savepoint2).err().unwrap(),
SavepointError::InvalidSavepoint
));
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(definition).unwrap();
assert_eq!(table.get(&0).unwrap().unwrap().value(), "hello");
let mut txn = db.begin_write().unwrap();
txn.restore_savepoint(&savepoint).unwrap();
txn.commit().unwrap();
}
#[test]
fn savepoint_restore_data_loss_pending_table_updates() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u64, &str> = TableDefinition::new("data_loss_test");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.insert(&1, "original").unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.insert(&2, "should_be_rolled_back").unwrap();
table.insert(&3, "also_should_be_rolled_back").unwrap();
}
let mut txn = txn;
txn.restore_savepoint(&savepoint).unwrap();
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(definition).unwrap();
assert_eq!(table.get(&1).unwrap().unwrap().value(), "original");
assert!(
table.get(&2).unwrap().is_none(),
"DATA LOSS BUG: key 2 should not exist after savepoint restore, \
but pending_table_updates re-applied stale modifications"
);
assert!(
table.get(&3).unwrap().is_none(),
"DATA LOSS BUG: key 3 should not exist after savepoint restore, \
but pending_table_updates re-applied stale modifications"
);
assert_eq!(
table.len().unwrap(),
1,
"DATA LOSS BUG: table should have 1 entry after savepoint restore, \
but has {} due to stale pending_table_updates",
table.len().unwrap()
);
}
#[test]
fn compaction() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &[u8]> = TableDefinition::new("x");
let big_value = vec![0u8; 100 * 1024];
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
for i in 0..100 {
table.insert(&i, big_value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
for i in 0..90 {
table.remove(&i).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
drop(db);
let file_size = tmpfile.as_file().metadata().unwrap().len();
let mut db = Database::open(tmpfile.path()).unwrap();
assert!(db.compact().unwrap());
drop(db);
let file_size2 = tmpfile.as_file().metadata().unwrap().len();
assert!(file_size2 < file_size);
}
#[test]
fn compact_after_non_durable_commit() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &[u8]> = TableDefinition::new("x");
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.insert(&0, [0; 1024].as_slice()).unwrap();
}
txn.commit().unwrap();
db.compact().unwrap();
}
#[test]
fn compact_after_post_commit_page_reuse() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &[u8]> = TableDefinition::new("x");
let big_value = vec![0u8; 100 * 1024];
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
for i in 0..25 {
table.insert(&i, big_value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
for i in 0..20 {
table.remove(&i).unwrap();
}
}
txn.commit().unwrap();
db.compact().unwrap();
}
#[test]
fn compact_does_not_grow_file() {
let tmpfile = create_tempfile();
let def: TableDefinition<u64, &[u8]> = TableDefinition::new("x");
{
let db = Database::create(tmpfile.path()).unwrap();
let value = vec![0u8; 3500];
for batch in 0..200 {
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(def).unwrap();
for i in 0..25u64 {
let k = batch * 1000 + i;
t.insert(&k, value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
}
}
let before = tmpfile.as_file().metadata().unwrap().len();
{
let mut db = Database::open(tmpfile.path()).unwrap();
db.compact().unwrap();
}
let after = tmpfile.as_file().metadata().unwrap().len();
assert!(
after <= before,
"compact() grew the file from {before} -> {after} bytes"
);
}
#[test]
fn compact_returns_persistent_savepoint_error() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(1, 1).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let _id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let err = db.compact().unwrap_err();
assert!(
matches!(err, CompactionError::PersistentSavepointExists),
"expected PersistentSavepointExists, got {err:?}"
);
}
#[test]
fn compact_returns_ephemeral_savepoint_error() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(1, 1).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
txn.commit().unwrap();
let err = db.compact().unwrap_err();
assert!(
matches!(err, CompactionError::EphemeralSavepointExists),
"expected EphemeralSavepointExists, got {err:?}"
);
drop(savepoint);
}
#[test]
fn compact_returns_persistent_savepoint_error_after_reopen() {
let tmpfile = create_tempfile();
{
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(1, 1).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let _id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
}
let mut db = Database::open(tmpfile.path()).unwrap();
let err = db.compact().unwrap_err();
assert!(
matches!(err, CompactionError::PersistentSavepointExists),
"expected PersistentSavepointExists, got {err:?}"
);
}
#[test]
fn compact_returns_transaction_in_progress_error() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(1, 1).unwrap();
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let err = db.compact().unwrap_err();
assert!(
matches!(err, CompactionError::TransactionInProgress),
"expected TransactionInProgress, got {err:?}"
);
drop(read_txn);
db.compact().unwrap();
}
#[test]
fn compact_does_not_block_on_held_write_transaction_with_savepoint() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(1, 1).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
let (sender, receiver) = std::sync::mpsc::channel();
let handle = std::thread::spawn(move || {
let result = db.compact();
sender.send(()).unwrap();
(db, result)
});
receiver
.recv_timeout(std::time::Duration::from_secs(30))
.expect("compact() deadlocked on the caller's own write transaction");
let (mut db, result) = handle.join().unwrap();
let err = result.unwrap_err();
assert!(
matches!(err, CompactionError::EphemeralSavepointExists),
"expected EphemeralSavepointExists, got {err:?}"
);
drop(savepoint);
txn.abort().unwrap();
db.compact().unwrap();
}
#[test]
fn compact_returns_persistent_savepoint_error_before_commit() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
table.insert(1, 1).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let _id = txn.persistent_savepoint().unwrap();
let (sender, receiver) = std::sync::mpsc::channel();
let handle = std::thread::spawn(move || {
let result = db.compact();
sender.send(()).unwrap();
(db, result)
});
receiver
.recv_timeout(std::time::Duration::from_secs(30))
.expect("compact() deadlocked on the caller's own write transaction");
let (mut db, result) = handle.join().unwrap();
let err = result.unwrap_err();
assert!(
matches!(err, CompactionError::PersistentSavepointExists),
"expected PersistentSavepointExists, got {err:?}"
);
txn.abort().unwrap();
db.compact().unwrap();
}
fn require_send<T: Send>(_: &T) {}
fn require_sync<T: Sync + Send>(_: &T) {}
#[test]
fn is_send() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &[u8]> = TableDefinition::new("x");
let txn = db.begin_write().unwrap();
{
let table = txn.open_table(definition).unwrap();
require_send(&table);
require_sync(&txn);
}
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_table(definition).unwrap();
require_sync(&table);
require_sync(&txn);
}
struct DelegatingTable<K: Key + 'static, V: Value + 'static, T: ReadableTable<K, V>> {
inner: T,
_key: PhantomData<K>,
_value: PhantomData<V>,
}
impl<K: Key + 'static, V: Value + 'static, T: ReadableTable<K, V>> ReadableTable<K, V>
for DelegatingTable<K, V, T>
{
fn get<'a>(
&self,
key: impl Borrow<K::SelfType<'a>>,
) -> redb::Result<Option<AccessGuard<'_, V>>> {
self.inner.get(key)
}
#[cfg(feature = "experimental-api-5")]
fn range<'a>(&self, range: impl KeyRange<'a, K>) -> redb::Result<Range<'_, K, V>> {
self.inner.range(range)
}
#[cfg(not(feature = "experimental-api-5"))]
fn range<'a, KR>(&self, range: impl RangeBounds<KR> + 'a) -> redb::Result<Range<'_, K, V>>
where
KR: Borrow<K::SelfType<'a>> + 'a,
{
self.inner.range(range)
}
fn first(&self) -> redb::Result<Option<(AccessGuard<'_, K>, AccessGuard<'_, V>)>> {
self.inner.first()
}
fn last(&self) -> redb::Result<Option<(AccessGuard<'_, K>, AccessGuard<'_, V>)>> {
self.inner.last()
}
#[cfg(feature = "experimental-api-5")]
fn lower_bound<'a>(
&self,
bound: std::ops::Bound<impl Borrow<K::SelfType<'a>>>,
) -> redb::Result<redb::Cursor<'_, K, V>> {
self.inner.lower_bound(bound)
}
#[cfg(feature = "experimental-api-5")]
fn upper_bound<'a>(
&self,
bound: std::ops::Bound<impl Borrow<K::SelfType<'a>>>,
) -> redb::Result<redb::Cursor<'_, K, V>> {
self.inner.upper_bound(bound)
}
}
impl<K: Key + 'static, V: Value + 'static, T: ReadableTable<K, V>> ReadableTableMetadata
for DelegatingTable<K, V, T>
{
fn stats(&self) -> redb::Result<TableStats> {
self.inner.stats()
}
fn len(&self) -> redb::Result<u64> {
self.inner.len()
}
}
struct DelegatingMultimapTable<K: Key + 'static, V: Key + 'static, T: ReadableMultimapTable<K, V>> {
inner: T,
_key: PhantomData<K>,
_value: PhantomData<V>,
}
impl<K: Key + 'static, V: Key + 'static, T: ReadableMultimapTable<K, V>> ReadableMultimapTable<K, V>
for DelegatingMultimapTable<K, V, T>
{
fn get<'a>(&self, key: impl Borrow<K::SelfType<'a>>) -> redb::Result<MultimapValue<'_, V>> {
self.inner.get(key)
}
#[cfg(feature = "experimental-api-5")]
fn range<'a>(&self, range: impl KeyRange<'a, K>) -> redb::Result<MultimapRange<'_, K, V>> {
self.inner.range(range)
}
#[cfg(not(feature = "experimental-api-5"))]
fn range<'a, KR>(
&self,
range: impl RangeBounds<KR> + 'a,
) -> redb::Result<MultimapRange<'_, K, V>>
where
KR: Borrow<K::SelfType<'a>> + 'a,
{
self.inner.range(range)
}
#[cfg(feature = "experimental-api-5")]
fn lower_bound<'a>(
&self,
bound: std::ops::Bound<impl Borrow<K::SelfType<'a>>>,
) -> redb::Result<redb::MultimapCursor<'_, K, V>> {
self.inner.lower_bound(bound)
}
#[cfg(feature = "experimental-api-5")]
fn upper_bound<'a>(
&self,
bound: std::ops::Bound<impl Borrow<K::SelfType<'a>>>,
) -> redb::Result<redb::MultimapCursor<'_, K, V>> {
self.inner.upper_bound(bound)
}
}
impl<K: Key + 'static, V: Key + 'static, T: ReadableMultimapTable<K, V>> ReadableTableMetadata
for DelegatingMultimapTable<K, V, T>
{
fn stats(&self) -> redb::Result<TableStats> {
self.inner.stats()
}
fn len(&self) -> redb::Result<u64> {
self.inner.len()
}
}
#[test]
fn custom_table_type() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let definition: TableDefinition<u32, &str> = TableDefinition::new("x");
let definition_multimap: MultimapTableDefinition<u32, &str> =
MultimapTableDefinition::new("multi");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(definition).unwrap();
table.insert(0, "hello").unwrap();
let mut table = txn.open_multimap_table(definition_multimap).unwrap();
table.insert(1, "world").unwrap();
}
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = DelegatingTable {
inner: txn.open_table(definition).unwrap(),
_key: Default::default(),
_value: Default::default(),
};
assert_eq!("hello", table.get(0).unwrap().unwrap().value());
let table = DelegatingMultimapTable {
inner: txn.open_multimap_table(definition_multimap).unwrap(),
_key: Default::default(),
_value: Default::default(),
};
assert_eq!(
"world",
table.get(1).unwrap().next().unwrap().unwrap().value()
);
let txn = db.begin_write().unwrap();
let table = DelegatingTable {
inner: txn.open_table(definition).unwrap(),
_key: Default::default(),
_value: Default::default(),
};
assert_eq!("hello", table.get(0).unwrap().unwrap().value());
let table = DelegatingMultimapTable {
inner: txn.open_multimap_table(definition_multimap).unwrap(),
_key: Default::default(),
_value: Default::default(),
};
assert_eq!(
"world",
table.get(1).unwrap().next().unwrap().unwrap().value()
);
}
#[test]
fn rename_table_with_modifications_data_loss() {
let tmpfile = create_tempfile();
const ORIGINAL_TABLE: TableDefinition<u64, &str> = TableDefinition::new("original");
const RENAMED_TABLE: TableDefinition<u64, &str> = TableDefinition::new("renamed");
{
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(ORIGINAL_TABLE).unwrap();
table.insert(0, "initial_value_0").unwrap();
table.insert(1, "initial_value_1").unwrap();
}
txn.commit().unwrap();
}
{
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(ORIGINAL_TABLE).unwrap();
table.insert(2, "new_value_2").unwrap();
table.insert(3, "new_value_3").unwrap();
txn.rename_table(table, RENAMED_TABLE).unwrap();
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(RENAMED_TABLE).unwrap();
assert_eq!(table.get(0).unwrap().unwrap().value(), "initial_value_0");
assert_eq!(table.get(1).unwrap().unwrap().value(), "initial_value_1");
assert_eq!(table.get(2).unwrap().unwrap().value(), "new_value_2");
assert_eq!(table.get(3).unwrap().unwrap().value(), "new_value_3");
assert_eq!(table.len().unwrap(), 4);
}
{
let mut db = Database::builder().create(tmpfile.path()).unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(RENAMED_TABLE).unwrap();
assert_eq!(table.get(0).unwrap().unwrap().value(), "initial_value_0");
assert_eq!(table.get(1).unwrap().unwrap().value(), "initial_value_1");
assert_eq!(table.get(2).unwrap().unwrap().value(), "new_value_2");
assert_eq!(table.get(3).unwrap().unwrap().value(), "new_value_3");
assert_eq!(table.len().unwrap(), 4);
drop(table);
drop(read_txn);
assert!(db.check_integrity().unwrap());
}
}
#[test]
fn savepoint_restore_data_loss_stale_freed_pages() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_a: TableDefinition<u64, &[u8]> = TableDefinition::new("table_a");
let table_b: TableDefinition<u64, &[u8]> = TableDefinition::new("table_b");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_a).unwrap();
for i in 0..100u64 {
let value = vec![0u8; 200];
table.insert(&i, value.as_slice()).unwrap();
}
}
txn.commit().unwrap();
{
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(table_a).unwrap();
for i in 0..100u64 {
let val = table.get(&i).unwrap().unwrap();
assert!(val.value().iter().all(|x| *x == 0));
}
}
let txn = db.begin_write().unwrap();
let savepoint = txn.ephemeral_savepoint().unwrap();
{
let mut table = txn.open_table(table_a).unwrap();
for i in 0..100u64 {
let overwrite = vec![0xFFu8; 200];
table.insert(&i, overwrite.as_slice()).unwrap();
}
}
let mut txn = txn;
txn.restore_savepoint(&savepoint).unwrap();
drop(savepoint);
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_b).unwrap();
for i in 0..300u64 {
let filler = vec![0xBBu8; 200];
table.insert(&i, filler.as_slice()).unwrap();
}
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(table_a).unwrap();
let mut data_loss_detected = false;
for i in 0..100u64 {
match table.get(&i) {
Ok(Some(val)) => {
let v = val.value();
if v.len() != 200 || v.iter().any(|x| *x != 0) {
data_loss_detected = true;
break;
}
}
Ok(None) => {
data_loss_detected = true;
break;
}
Err(_) => {
data_loss_detected = true;
break;
}
}
}
assert!(
!data_loss_detected,
"DATA LOSS: table_a data is corrupted after savepoint restore. \
restore_savepoint() did not clear freed_pages, causing the old pages \
(still referenced by the committed tree root) to be stored in DATA_FREED_TABLE \
and later freed by process_freed_pages(). When those pages were reused by table_b, \
table_a's B-tree became corrupted."
);
}
#[test]
fn restore_savepoint_from_foreign_database_panics() {
let tmpfile1 = create_tempfile();
let tmpfile2 = create_tempfile();
let db1 = Database::create(tmpfile1.path()).unwrap();
let db2 = Database::create(tmpfile2.path()).unwrap();
let txn1 = db1.begin_write().unwrap();
let foreign_savepoint = txn1.ephemeral_savepoint().unwrap();
txn1.commit().unwrap();
let mut txn2 = db2.begin_write().unwrap();
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
txn2.restore_savepoint(&foreign_savepoint)
}));
assert!(result.is_ok());
assert!(matches!(
result.unwrap(),
Err(SavepointError::InvalidSavepoint)
));
}
#[test]
fn delete_table_panic_after_modification() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_a: TableDefinition<u64, &[u8]> = TableDefinition::new("table_a");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_a).unwrap();
for i in 0..100u64 {
table.insert(&i, &vec![0u8; 200][..]).unwrap();
}
}
txn.commit().unwrap();
{
let read_txn = db.begin_read().unwrap();
let table = read_txn.open_table(table_a).unwrap();
assert_eq!(table.len().unwrap(), 100);
}
let txn = db.begin_write().unwrap();
let _savepoint = txn.ephemeral_savepoint().unwrap();
{
let mut table = txn.open_table(table_a).unwrap();
for i in 0..100u64 {
table.insert(&i, &vec![0xFFu8; 200][..]).unwrap();
}
}
let deleted = txn.delete_table(table_a).unwrap();
assert!(deleted);
let commit_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| txn.commit()));
assert!(commit_result.is_ok() && commit_result.as_ref().unwrap().is_ok());
}
#[test]
fn persistent_savepoint_abort_unbounded_leak() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table: TableDefinition<u64, u64> = TableDefinition::new("data");
{
let txn = db.begin_write().unwrap();
let id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.delete_persistent_savepoint(id).unwrap();
txn.commit().unwrap();
}
{
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table).unwrap();
for i in 0..20u64 {
t.insert(i, i).unwrap();
}
}
txn.commit().unwrap();
}
for _ in 0..3 {
db.begin_write().unwrap().commit().unwrap();
}
let txn = db.begin_write().unwrap();
let baseline = txn.stats().unwrap().allocated_pages();
txn.abort().unwrap();
const ITERATIONS: u64 = 20;
for round in 0..ITERATIONS {
{
let txn = db.begin_write().unwrap();
let _id = txn.persistent_savepoint().unwrap();
txn.abort().unwrap();
}
{
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table).unwrap();
t.insert(0, round).unwrap();
}
txn.commit().unwrap();
}
for _ in 0..3 {
db.begin_write().unwrap().commit().unwrap();
}
}
let txn = db.begin_write().unwrap();
let after = txn.stats().unwrap().allocated_pages();
txn.abort().unwrap();
assert_eq!(
baseline,
after,
"After {} iterations of persistent_savepoint+abort+modify, page usage grew \
from {} to {} ({} pages leaked). Because the leak per iteration is \
independent of N, running N iterations leaks O(N) pages.",
ITERATIONS,
baseline,
after,
after.saturating_sub(baseline),
);
}
#[test]
fn check_integrity_with_live_read_transaction() {
let tmpfile = create_tempfile();
let mut db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(U64_TABLE).unwrap();
for i in 0..10u64 {
table.insert(&i, &i).unwrap();
}
}
txn.commit().unwrap();
let read_txn = db.begin_read().unwrap();
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| db.check_integrity()));
assert!(result.is_ok(), "check_integrity() should not panic");
assert!(matches!(
result.unwrap(),
Err(DatabaseError::TransactionInProgress)
));
drop(read_txn);
assert!(db.check_integrity().unwrap());
}
#[test]
fn restore_savepoint_partial_revert_commits_as_data_loss() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let tab: TableDefinition<u64, u64> = TableDefinition::new("t");
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(tab).unwrap();
t.insert(1u64, 1u64).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let ps1 = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(tab).unwrap();
t.insert(2u64, 2u64).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
let _ps2 = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(tab).unwrap();
t.insert(3u64, 3u64).unwrap();
}
txn.commit().unwrap();
{
let rt = db.begin_read().unwrap();
let t = rt.open_table(tab).unwrap();
assert_eq!(t.get(1u64).unwrap().unwrap().value(), 1);
assert_eq!(t.get(2u64).unwrap().unwrap().value(), 2);
assert_eq!(t.get(3u64).unwrap().unwrap().value(), 3);
}
let mut txn = db.begin_write().unwrap();
txn.set_durability(Durability::None).unwrap();
let sp = txn.get_persistent_savepoint(ps1).unwrap();
let restore_result = txn.restore_savepoint(&sp);
assert!(
matches!(
restore_result,
Err(SavepointError::ImmediateDurabilityRequired)
),
"restore_savepoint() should fail with ImmediateDurabilityRequired when \
durability != Immediate and newer persistent savepoints exist, got: {restore_result:?}"
);
txn.commit().unwrap();
let rt = db.begin_read().unwrap();
let t = rt.open_table(tab).unwrap();
let k1 = t.get(1u64).unwrap().map(|v| v.value());
let k2 = t.get(2u64).unwrap().map(|v| v.value());
let k3 = t.get(3u64).unwrap().map(|v| v.value());
assert_eq!(k1, Some(1));
assert_eq!(k2, Some(2));
assert_eq!(k3, Some(3));
}
#[test]
fn restore_savepoint_abort_unbounded_leak() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table: TableDefinition<u64, u64> = TableDefinition::new("data");
{
let txn = db.begin_write().unwrap();
let id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
txn.delete_persistent_savepoint(id).unwrap();
txn.commit().unwrap();
}
{
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table).unwrap();
for i in 0..50u64 {
t.insert(i, i).unwrap();
}
}
txn.commit().unwrap();
}
for _ in 0..3 {
db.begin_write().unwrap().commit().unwrap();
}
let txn = db.begin_write().unwrap();
let baseline = txn.stats().unwrap().allocated_pages();
txn.abort().unwrap();
let older = {
let txn = db.begin_write().unwrap();
let id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
id
};
{
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table).unwrap();
t.insert(0, u64::MAX).unwrap();
}
txn.commit().unwrap();
}
let newer = {
let txn = db.begin_write().unwrap();
let id = txn.persistent_savepoint().unwrap();
txn.commit().unwrap();
id
};
{
let mut txn = db.begin_write().unwrap();
let sp = txn.get_persistent_savepoint(older).unwrap();
txn.restore_savepoint(&sp).unwrap();
drop(sp);
txn.abort().unwrap();
}
{
let txn = db.begin_write().unwrap();
txn.delete_persistent_savepoint(older).unwrap();
txn.commit().unwrap();
}
const ITERATIONS: u64 = 100;
for round in 0..ITERATIONS {
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table).unwrap();
t.insert(0, round).unwrap();
}
txn.commit().unwrap();
}
drop(db);
let db = Database::create(tmpfile.path()).unwrap();
{
let mut txn = db.begin_write().unwrap();
let sp = txn.get_persistent_savepoint(newer).unwrap();
txn.restore_savepoint(&sp).unwrap();
drop(sp);
txn.commit().unwrap();
}
{
let txn = db.begin_write().unwrap();
txn.delete_persistent_savepoint(newer).unwrap();
txn.commit().unwrap();
}
for _ in 0..3 {
db.begin_write().unwrap().commit().unwrap();
}
let txn = db.begin_write().unwrap();
let after = txn.stats().unwrap().allocated_pages();
txn.abort().unwrap();
assert_eq!(
baseline,
after,
"After {} iterations of insert+commit between an aborted \
restore_savepoint and the next restore_savepoint, page usage grew \
from {} to {} ({} pages leaked). The leak per iteration is independent \
of N, so running N iterations leaks O(N) pages.",
ITERATIONS,
baseline,
after,
after.saturating_sub(baseline),
);
}
#[test]
fn restore_savepoint_abort_after_ephemeral_drop() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table: TableDefinition<u64, u64> = TableDefinition::new("data");
{
let txn = db.begin_write().unwrap();
{
let mut t = txn.open_table(table).unwrap();
for i in 0..20u64 {
t.insert(i, i).unwrap();
}
}
txn.commit().unwrap();
}
let older = {
let txn = db.begin_write().unwrap();
let sp = txn.ephemeral_savepoint().unwrap();
txn.commit().unwrap();
sp
};
let newer = {
let txn = db.begin_write().unwrap();
let sp = txn.ephemeral_savepoint().unwrap();
txn.commit().unwrap();
sp
};
{
let mut txn = db.begin_write().unwrap();
txn.restore_savepoint(&older).unwrap();
drop(newer);
txn.abort().unwrap();
}
{
let mut txn = db.begin_write().unwrap();
txn.restore_savepoint(&older).unwrap();
txn.commit().unwrap();
}
drop(older);
for _ in 0..3 {
db.begin_write().unwrap().commit().unwrap();
}
assert!(
db.begin_read()
.unwrap()
.open_table(table)
.unwrap()
.get(&0u64)
.unwrap()
.is_some()
);
}
#[test]
fn extract_if_next_then_next_back_panic() {
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let table_def: TableDefinition<u64, u64> = TableDefinition::new("t");
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
table.insert(&1u64, &10u64).unwrap();
}
txn.commit().unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_table(table_def).unwrap();
let mut iter = table.extract_if(|_, _| true).unwrap();
let first = iter.next();
assert!(first.is_some());
let first = first.unwrap();
assert!(first.is_ok());
assert!(iter.next().is_none());
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| iter.next_back()));
assert!(result.is_ok());
}
txn.abort().unwrap();
}
#[test]
fn multimap_value_next_back_does_not_update_len() {
const TABLE: MultimapTableDefinition<u32, u32> = MultimapTableDefinition::new("m");
let tmpfile = create_tempfile();
let db = Database::create(tmpfile.path()).unwrap();
let txn = db.begin_write().unwrap();
{
let mut table = txn.open_multimap_table(TABLE).unwrap();
for i in 0..5u32 {
table.insert(&1u32, &i).unwrap();
}
}
txn.commit().unwrap();
let txn = db.begin_read().unwrap();
let table = txn.open_multimap_table(TABLE).unwrap();
let mut iter = table.get(&1u32).unwrap();
assert_eq!(iter.len(), 5);
assert!(!iter.is_empty());
let _ = iter.next_back().unwrap().unwrap();
assert_eq!(iter.len(), 4);
for expected_remaining in (0..4).rev() {
iter.next_back().unwrap().unwrap();
assert_eq!(iter.len(), expected_remaining as u64);
}
assert!(iter.is_empty());
assert!(iter.next_back().is_none());
}
#[test]
#[should_panic(expected = "assertion failed: !name.is_empty()")]
fn table_definition_new_panics_on_empty_name() {
let name = String::new();
let _def: TableDefinition<u64, u64> = TableDefinition::new(&name);
}
#[test]
#[should_panic(expected = "assertion failed: !name.is_empty()")]
fn multimap_table_definition_new_panics_on_empty_name() {
let name = String::new();
let _def: MultimapTableDefinition<u64, u64> = MultimapTableDefinition::new(&name);
}