use crate::inner::Inner;
use crate::options::DatabaseOptions;
#[cfg(not(target_arch = "wasm32"))]
use crate::options::{DEFAULT_CLEANUP_INTERVAL, DEFAULT_GC_INTERVAL};
#[cfg(not(target_arch = "wasm32"))]
use crate::persistence::Persistence;
use crate::pool::Pool;
use crate::pool::DEFAULT_POOL_SIZE;
use crate::tx::Transaction;
use std::ops::Deref;
use std::sync::atomic::Ordering;
use std::sync::Arc;
#[cfg(not(target_arch = "wasm32"))]
use std::time::Duration;
pub struct Database {
inner: Arc<Inner>,
pool: Arc<Pool>,
#[cfg(not(target_arch = "wasm32"))]
persistence: Option<Persistence>,
#[cfg(not(target_arch = "wasm32"))]
gc_interval: Duration,
#[cfg(not(target_arch = "wasm32"))]
cleanup_interval: Duration,
}
impl Default for Database {
fn default() -> Self {
let inner = Arc::new(Inner::default());
let pool = Pool::new(inner.clone(), DEFAULT_POOL_SIZE);
Database {
inner,
pool,
#[cfg(not(target_arch = "wasm32"))]
persistence: None,
#[cfg(not(target_arch = "wasm32"))]
gc_interval: DEFAULT_GC_INTERVAL,
#[cfg(not(target_arch = "wasm32"))]
cleanup_interval: DEFAULT_CLEANUP_INTERVAL,
}
}
}
impl Drop for Database {
fn drop(&mut self) {
self.shutdown();
}
}
impl Deref for Database {
type Target = Inner;
fn deref(&self) -> &Self::Target {
&self.inner
}
}
impl Database {
pub fn new() -> Self {
Self::new_with_options(DatabaseOptions::default())
}
pub fn new_with_options(opts: DatabaseOptions) -> Self {
let inner = Arc::new(Inner::new(&opts));
let pool = Pool::new(inner.clone(), opts.pool_size);
let db = Database {
inner,
pool,
#[cfg(not(target_arch = "wasm32"))]
persistence: None,
#[cfg(not(target_arch = "wasm32"))]
gc_interval: opts.gc_interval,
#[cfg(not(target_arch = "wasm32"))]
cleanup_interval: opts.cleanup_interval,
};
#[cfg(not(target_arch = "wasm32"))]
{
if opts.enable_cleanup {
db.initialise_cleanup_worker();
}
if opts.enable_gc {
db.initialise_garbage_worker();
}
}
db
}
#[cfg(not(target_arch = "wasm32"))]
pub fn new_with_persistence(
opts: DatabaseOptions,
persistence_opts: crate::PersistenceOptions,
) -> std::io::Result<Self> {
let inner = Arc::new(Inner::new(&opts));
let pool = Pool::new(inner.clone(), opts.pool_size);
let persist = Persistence::new_with_options(persistence_opts, inner.clone())
.map_err(std::io::Error::other)?;
inner.persistence.write().replace(Arc::new(persist.clone()));
let db = Database {
inner,
pool,
persistence: Some(persist),
gc_interval: opts.gc_interval,
cleanup_interval: opts.cleanup_interval,
};
if opts.enable_cleanup {
db.initialise_cleanup_worker();
}
if opts.enable_gc {
db.initialise_garbage_worker();
}
Ok(db)
}
pub fn transaction(&self, write: bool) -> Transaction {
self.pool.get(write)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn persistence(&self) -> Option<&Persistence> {
self.persistence.as_ref()
}
pub fn run_cleanup(&self) {
self.inner.cleanup_commit_queue();
}
pub fn run_gc(&self) {
if let Some(cleanup_ts) = self.compute_cleanup_ts() {
self.gc_candidates.pin().clear();
self.run_gc_full(cleanup_ts);
}
}
pub fn run_gc_tracked(&self) {
if self.gc_candidates.pin().is_empty() {
return;
}
if let Some(cleanup_ts) = self.compute_cleanup_ts() {
self.inner.run_gc_tracked(cleanup_ts);
}
}
fn shutdown(&self) {
#[cfg(not(target_arch = "wasm32"))]
{
if let Some(ref persistence) = self.persistence {
persistence.background_threads_enabled.store(false, Ordering::Release);
if let Some(handle) = persistence.fsync_handle.write().take() {
handle.thread().unpark();
let _ = handle.join();
}
if let Some(handle) = persistence.snapshot_handle.write().take() {
handle.thread().unpark();
let _ = handle.join();
}
if let Some(handle) = persistence.appender_handle.write().take() {
handle.thread().unpark();
let _ = handle.join();
}
}
}
self.background_threads_enabled.store(false, Ordering::Relaxed);
#[cfg(not(target_arch = "wasm32"))]
{
if let Some(handle) = self.transaction_cleanup_handle.write().take() {
handle.thread().unpark();
let _ = handle.join();
}
if let Some(handle) = self.garbage_collection_handle.write().take() {
handle.thread().unpark();
let _ = handle.join();
}
}
}
#[cfg(not(target_arch = "wasm32"))]
fn initialise_cleanup_worker(&self) {
let db = self.inner.clone();
if db.transaction_cleanup_handle.read().is_none() {
let interval = self.cleanup_interval;
let handle = std::thread::spawn(move || {
while db.background_threads_enabled.load(Ordering::Relaxed) {
std::thread::park_timeout(interval);
if !db.background_threads_enabled.load(Ordering::Relaxed) {
break;
}
db.cleanup_commit_queue();
}
});
*self.inner.transaction_cleanup_handle.write() = Some(handle);
}
}
#[cfg(not(target_arch = "wasm32"))]
fn initialise_garbage_worker(&self) {
let db = self.inner.clone();
if db.garbage_collection_handle.read().is_none() {
let interval = self.gc_interval;
let handle = std::thread::spawn(move || {
while db.background_threads_enabled.load(Ordering::Relaxed) {
std::thread::park_timeout(interval);
if !db.background_threads_enabled.load(Ordering::Relaxed) {
break;
}
if db.gc_candidates.pin().is_empty() {
continue;
}
let Some(cleanup_ts) = db.compute_cleanup_ts() else {
continue;
};
db.run_gc_tracked(cleanup_ts);
}
});
*self.inner.garbage_collection_handle.write() = Some(handle);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn begin_tx() {
let db = Database::new();
db.transaction(false);
}
#[test]
fn finished_tx_not_writeable() {
let db = Database::new();
let mut tx = db.transaction(true);
let res = tx.cancel();
assert!(res.is_ok());
let res = tx.put("test", "something");
assert!(res.is_err());
let res = tx.set("test", "something");
assert!(res.is_err());
let res = tx.del("test");
assert!(res.is_err());
let res = tx.commit();
assert!(res.is_err());
let res = tx.cancel();
assert!(res.is_err());
}
#[test]
fn cancelled_tx_is_cancelled() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("test", "something").unwrap();
let res = tx.exists("test").unwrap();
assert!(res);
let res = tx.get("test").unwrap();
assert_eq!(res.as_deref(), Some(b"something" as &[u8]));
let res = tx.cancel();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.exists("test").unwrap();
assert!(!res);
let res = tx.get("test").unwrap();
assert_eq!(res, None);
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn committed_tx_is_committed() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("test", "something").unwrap();
let res = tx.exists("test").unwrap();
assert!(res);
let res = tx.get("test").unwrap();
assert_eq!(res.as_deref(), Some(b"something" as &[u8]));
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.exists("test").unwrap();
assert!(res);
let res = tx.get("test").unwrap();
assert_eq!(res.as_deref(), Some(b"something" as &[u8]));
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn multiple_concurrent_readers() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("test", "something").unwrap();
let res = tx.exists("test").unwrap();
assert!(res);
let res = tx.get("test").unwrap();
assert_eq!(res.as_deref(), Some(b"something" as &[u8]));
let res = tx.commit();
assert!(res.is_ok());
let mut tx1 = db.transaction(false);
let res = tx1.exists("test").unwrap();
assert!(res);
let res = tx1.exists("temp").unwrap();
assert!(!res);
let mut tx2 = db.transaction(false);
let res = tx2.exists("test").unwrap();
assert!(res);
let res = tx2.exists("temp").unwrap();
assert!(!res);
let res = tx1.cancel();
assert!(res.is_ok());
let res = tx2.cancel();
assert!(res.is_ok());
}
#[test]
fn multiple_concurrent_operators() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("test", "something").unwrap();
let res = tx.exists("test").unwrap();
assert!(res);
let res = tx.get("test").unwrap();
assert_eq!(res.as_deref(), Some(b"something" as &[u8]));
let res = tx.commit();
assert!(res.is_ok());
let mut tx1 = db.transaction(false);
let res = tx1.exists("test").unwrap();
assert!(res);
let res = tx1.exists("temp").unwrap();
assert!(!res);
let mut txw = db.transaction(true);
txw.put("temp", "other").unwrap();
let res = txw.exists("test").unwrap();
assert!(res);
let res = txw.exists("temp").unwrap();
assert!(res);
let res = txw.commit();
assert!(res.is_ok());
let mut tx2 = db.transaction(false);
let res = tx2.exists("test").unwrap();
assert!(res);
let res = tx2.exists("temp").unwrap();
assert!(res);
let res = tx1.exists("temp").unwrap();
assert!(!res);
let res = tx1.cancel();
assert!(res.is_ok());
let res = tx2.cancel();
assert!(res.is_ok());
}
#[test]
fn iterate_keys_forward() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "a").unwrap();
tx.put("b", "b").unwrap();
tx.put("c", "c").unwrap();
tx.put("d", "d").unwrap();
tx.put("e", "e").unwrap();
tx.put("f", "f").unwrap();
tx.put("g", "g").unwrap();
tx.put("h", "h").unwrap();
tx.put("i", "i").unwrap();
tx.put("j", "j").unwrap();
tx.put("k", "k").unwrap();
tx.put("l", "l").unwrap();
tx.put("m", "m").unwrap();
tx.put("n", "n").unwrap();
tx.put("o", "o").unwrap();
let res = tx.keys("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].as_ref(), b"c");
assert_eq!(res[1], "d");
assert_eq!(res[2], "e");
assert_eq!(res[3], "f");
assert_eq!(res[4], "g");
assert_eq!(res[5], "h");
assert_eq!(res[6], "i");
assert_eq!(res[7], "j");
assert_eq!(res[8], "k");
assert_eq!(res[9], "l");
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.keys("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].as_ref(), b"c");
assert_eq!(res[1], "d");
assert_eq!(res[2], "e");
assert_eq!(res[3], "f");
assert_eq!(res[4], "g");
assert_eq!(res[5], "h");
assert_eq!(res[6], "i");
assert_eq!(res[7], "j");
assert_eq!(res[8], "k");
assert_eq!(res[9], "l");
let res = tx.cancel();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.keys("c".."z", Some(3), Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0], "f");
assert_eq!(res[1], "g");
assert_eq!(res[2], "h");
assert_eq!(res[3], "i");
assert_eq!(res[4], "j");
assert_eq!(res[5], "k");
assert_eq!(res[6], "l");
assert_eq!(res[7], "m");
assert_eq!(res[8], "n");
assert_eq!(res[9], "o");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn iterate_keys_reverse() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "a").unwrap();
tx.put("b", "b").unwrap();
tx.put("c", "c").unwrap();
tx.put("d", "d").unwrap();
tx.put("e", "e").unwrap();
tx.put("f", "f").unwrap();
tx.put("g", "g").unwrap();
tx.put("h", "h").unwrap();
tx.put("i", "i").unwrap();
tx.put("j", "j").unwrap();
tx.put("k", "k").unwrap();
tx.put("l", "l").unwrap();
tx.put("m", "m").unwrap();
tx.put("n", "n").unwrap();
tx.put("o", "o").unwrap();
let res = tx.keys_reverse("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0], "o");
assert_eq!(res[1], "n");
assert_eq!(res[2], "m");
assert_eq!(res[3], "l");
assert_eq!(res[4], "k");
assert_eq!(res[5], "j");
assert_eq!(res[6], "i");
assert_eq!(res[7], "h");
assert_eq!(res[8], "g");
assert_eq!(res[9], "f");
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.keys_reverse("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0], "o");
assert_eq!(res[1], "n");
assert_eq!(res[2], "m");
assert_eq!(res[3], "l");
assert_eq!(res[4], "k");
assert_eq!(res[5], "j");
assert_eq!(res[6], "i");
assert_eq!(res[7], "h");
assert_eq!(res[8], "g");
assert_eq!(res[9], "f");
let res = tx.cancel();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.keys_reverse("c".."z", Some(3), Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0], "l");
assert_eq!(res[1], "k");
assert_eq!(res[2], "j");
assert_eq!(res[3], "i");
assert_eq!(res[4], "h");
assert_eq!(res[5], "g");
assert_eq!(res[6], "f");
assert_eq!(res[7], "e");
assert_eq!(res[8], "d");
assert_eq!(res[9], "c");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn iterate_keys_values_forward() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "a").unwrap();
tx.put("b", "b").unwrap();
tx.put("c", "c").unwrap();
tx.put("d", "d").unwrap();
tx.put("e", "e").unwrap();
tx.put("f", "f").unwrap();
tx.put("g", "g").unwrap();
tx.put("h", "h").unwrap();
tx.put("i", "i").unwrap();
tx.put("j", "j").unwrap();
tx.put("k", "k").unwrap();
tx.put("l", "l").unwrap();
tx.put("m", "m").unwrap();
tx.put("n", "n").unwrap();
tx.put("o", "o").unwrap();
let res = tx.scan("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].0.as_ref(), b"c");
assert_eq!(res[0].1.as_ref(), b"c");
assert_eq!(res[1].0.as_ref(), b"d");
assert_eq!(res[1].1.as_ref(), b"d");
assert_eq!(res[2].0.as_ref(), b"e");
assert_eq!(res[3].0.as_ref(), b"f");
assert_eq!(res[4].0.as_ref(), b"g");
assert_eq!(res[5].0.as_ref(), b"h");
assert_eq!(res[6].0.as_ref(), b"i");
assert_eq!(res[7].0.as_ref(), b"j");
assert_eq!(res[8].0.as_ref(), b"k");
assert_eq!(res[9].0.as_ref(), b"l");
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.scan("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].0.as_ref(), b"c");
assert_eq!(res[0].1.as_ref(), b"c");
assert_eq!(res[1].0.as_ref(), b"d");
assert_eq!(res[1].1.as_ref(), b"d");
assert_eq!(res[2].0.as_ref(), b"e");
assert_eq!(res[3].0.as_ref(), b"f");
assert_eq!(res[4].0.as_ref(), b"g");
assert_eq!(res[5].0.as_ref(), b"h");
assert_eq!(res[6].0.as_ref(), b"i");
assert_eq!(res[7].0.as_ref(), b"j");
assert_eq!(res[8].0.as_ref(), b"k");
assert_eq!(res[9].0.as_ref(), b"l");
let res = tx.cancel();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.scan("c".."z", Some(3), Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].0.as_ref(), b"f");
assert_eq!(res[1].0.as_ref(), b"g");
assert_eq!(res[2].0.as_ref(), b"h");
assert_eq!(res[3].0.as_ref(), b"i");
assert_eq!(res[4].0.as_ref(), b"j");
assert_eq!(res[5].0.as_ref(), b"k");
assert_eq!(res[6].0.as_ref(), b"l");
assert_eq!(res[7].0.as_ref(), b"m");
assert_eq!(res[8].0.as_ref(), b"n");
assert_eq!(res[9].0.as_ref(), b"o");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn iterate_keys_values_reverse() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "a").unwrap();
tx.put("b", "b").unwrap();
tx.put("c", "c").unwrap();
tx.put("d", "d").unwrap();
tx.put("e", "e").unwrap();
tx.put("f", "f").unwrap();
tx.put("g", "g").unwrap();
tx.put("h", "h").unwrap();
tx.put("i", "i").unwrap();
tx.put("j", "j").unwrap();
tx.put("k", "k").unwrap();
tx.put("l", "l").unwrap();
tx.put("m", "m").unwrap();
tx.put("n", "n").unwrap();
tx.put("o", "o").unwrap();
let res = tx.scan_reverse("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].0.as_ref(), b"o");
assert_eq!(res[1].0.as_ref(), b"n");
assert_eq!(res[2].0.as_ref(), b"m");
assert_eq!(res[3].0.as_ref(), b"l");
assert_eq!(res[4].0.as_ref(), b"k");
assert_eq!(res[5].0.as_ref(), b"j");
assert_eq!(res[6].0.as_ref(), b"i");
assert_eq!(res[7].0.as_ref(), b"h");
assert_eq!(res[8].0.as_ref(), b"g");
assert_eq!(res[9].0.as_ref(), b"f");
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.scan_reverse("c".."z", None, Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].0.as_ref(), b"o");
assert_eq!(res[1].0.as_ref(), b"n");
assert_eq!(res[2].0.as_ref(), b"m");
assert_eq!(res[3].0.as_ref(), b"l");
assert_eq!(res[4].0.as_ref(), b"k");
assert_eq!(res[5].0.as_ref(), b"j");
assert_eq!(res[6].0.as_ref(), b"i");
assert_eq!(res[7].0.as_ref(), b"h");
assert_eq!(res[8].0.as_ref(), b"g");
assert_eq!(res[9].0.as_ref(), b"f");
let res = tx.cancel();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.scan_reverse("c".."z", Some(3), Some(10)).unwrap();
assert_eq!(res.len(), 10);
assert_eq!(res[0].0.as_ref(), b"l");
assert_eq!(res[1].0.as_ref(), b"k");
assert_eq!(res[2].0.as_ref(), b"j");
assert_eq!(res[3].0.as_ref(), b"i");
assert_eq!(res[4].0.as_ref(), b"h");
assert_eq!(res[5].0.as_ref(), b"g");
assert_eq!(res[6].0.as_ref(), b"f");
assert_eq!(res[7].0.as_ref(), b"e");
assert_eq!(res[8].0.as_ref(), b"d");
assert_eq!(res[9].0.as_ref(), b"c");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn count_keys_values() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "a").unwrap();
tx.put("b", "b").unwrap();
tx.put("c", "c").unwrap();
tx.put("d", "d").unwrap();
tx.put("e", "e").unwrap();
tx.put("f", "f").unwrap();
tx.put("g", "g").unwrap();
tx.put("h", "h").unwrap();
tx.put("i", "i").unwrap();
tx.put("j", "j").unwrap();
tx.put("k", "k").unwrap();
tx.put("l", "l").unwrap();
tx.put("m", "m").unwrap();
tx.put("n", "n").unwrap();
tx.put("o", "o").unwrap();
let res = tx.total("c".."z", None, Some(10)).unwrap();
assert_eq!(res, 10);
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let res = tx.total("c".."z", Some(3), Some(10)).unwrap();
assert_eq!(res, 10);
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn cursor_forward_iteration() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let mut cursor = tx.cursor("a".."z").unwrap();
cursor.seek_to_first();
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"a");
assert_eq!(cursor.value().unwrap().as_ref(), b"1");
cursor.next();
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"b");
cursor.next(); cursor.next(); cursor.next(); assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"e");
cursor.next();
assert!(!cursor.valid());
drop(cursor);
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn cursor_reverse_iteration() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let mut cursor = tx.cursor("a".."z").unwrap();
cursor.seek_to_last();
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"e");
assert_eq!(cursor.value().unwrap().as_ref(), b"5");
cursor.prev();
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"d");
cursor.prev(); cursor.prev(); cursor.prev(); assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"a");
cursor.prev();
assert!(!cursor.valid());
drop(cursor);
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn cursor_seek_operations() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("c", "3").unwrap();
tx.put("e", "5").unwrap();
tx.put("g", "7").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let mut cursor = tx.cursor("a".."z").unwrap();
cursor.seek("c");
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"c");
cursor.seek("d");
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"e");
cursor.seek_for_prev("e");
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"c");
cursor.seek("z");
assert!(!cursor.valid());
drop(cursor);
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn cursor_bidirectional_switch() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let mut cursor = tx.cursor("a".."z").unwrap();
cursor.seek_to_first();
assert_eq!(cursor.key().unwrap().as_ref(), b"a");
cursor.next();
assert_eq!(cursor.key().unwrap().as_ref(), b"b");
cursor.next();
assert_eq!(cursor.key().unwrap().as_ref(), b"c");
cursor.prev();
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"b");
cursor.next();
assert!(cursor.valid());
assert_eq!(cursor.key().unwrap().as_ref(), b"c");
drop(cursor);
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn keys_iterator_forward() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let keys: Vec<_> = tx.keys_iter("b".."e").unwrap().collect();
assert_eq!(keys.len(), 3);
assert_eq!(keys[0].as_ref(), b"b");
assert_eq!(keys[1].as_ref(), b"c");
assert_eq!(keys[2].as_ref(), b"d");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn keys_iterator_reverse() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let keys: Vec<_> = tx.keys_iter_reverse("b".."e").unwrap().collect();
assert_eq!(keys.len(), 3);
assert_eq!(keys[0].as_ref(), b"d");
assert_eq!(keys[1].as_ref(), b"c");
assert_eq!(keys[2].as_ref(), b"b");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn keys_iterator_with_take() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
tx.put("f", "6").unwrap();
tx.put("g", "7").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let keys: Vec<_> = tx.keys_iter("a".."z").unwrap().take(3).collect();
assert_eq!(keys.len(), 3);
assert_eq!(keys[0].as_ref(), b"a");
assert_eq!(keys[1].as_ref(), b"b");
assert_eq!(keys[2].as_ref(), b"c");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn keys_iterator_with_skip() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let keys: Vec<_> = tx.keys_iter("a".."z").unwrap().skip(2).collect();
assert_eq!(keys.len(), 3);
assert_eq!(keys[0].as_ref(), b"c");
assert_eq!(keys[1].as_ref(), b"d");
assert_eq!(keys[2].as_ref(), b"e");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn scan_iterator_forward() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let pairs: Vec<_> = tx.scan_iter("a".."z").unwrap().collect();
assert_eq!(pairs.len(), 3);
assert_eq!(pairs[0].0.as_ref(), b"a");
assert_eq!(pairs[0].1.as_ref(), b"1");
assert_eq!(pairs[1].0.as_ref(), b"b");
assert_eq!(pairs[1].1.as_ref(), b"2");
assert_eq!(pairs[2].0.as_ref(), b"c");
assert_eq!(pairs[2].1.as_ref(), b"3");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn scan_iterator_reverse() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let pairs: Vec<_> = tx.scan_iter_reverse("a".."z").unwrap().collect();
assert_eq!(pairs.len(), 3);
assert_eq!(pairs[0].0.as_ref(), b"c");
assert_eq!(pairs[0].1.as_ref(), b"3");
assert_eq!(pairs[1].0.as_ref(), b"b");
assert_eq!(pairs[1].1.as_ref(), b"2");
assert_eq!(pairs[2].0.as_ref(), b"a");
assert_eq!(pairs[2].1.as_ref(), b"1");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn scan_iterator_with_take() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
tx.put("d", "4").unwrap();
tx.put("e", "5").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(false);
let pairs: Vec<_> = tx.scan_iter("a".."z").unwrap().take(2).collect();
assert_eq!(pairs.len(), 2);
assert_eq!(pairs[0].0.as_ref(), b"a");
assert_eq!(pairs[1].0.as_ref(), b"b");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn iterator_sees_uncommitted_writes() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
let keys: Vec<_> = tx.keys_iter("a".."z").unwrap().collect();
assert_eq!(keys.len(), 2);
assert_eq!(keys[0].as_ref(), b"a");
assert_eq!(keys[1].as_ref(), b"b");
let res = tx.commit();
assert!(res.is_ok());
}
#[test]
fn cursor_handles_deleted_entries() {
let db = Database::new();
let mut tx = db.transaction(true);
tx.put("a", "1").unwrap();
tx.put("b", "2").unwrap();
tx.put("c", "3").unwrap();
let res = tx.commit();
assert!(res.is_ok());
let mut tx = db.transaction(true);
tx.del("b").unwrap();
let keys: Vec<_> = tx.keys_iter("a".."z").unwrap().collect();
assert_eq!(keys.len(), 2);
assert_eq!(keys[0].as_ref(), b"a");
assert_eq!(keys[1].as_ref(), b"c");
let res = tx.cancel();
assert!(res.is_ok());
}
#[test]
fn cleanup_trims_commit_queue_when_idle() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..10 {
let mut tx = db.transaction(true);
tx.set(format!("key{i}"), "value").unwrap();
tx.commit().unwrap();
}
assert_eq!(db.transaction_commit_queue.len(), 10);
db.run_cleanup();
assert_eq!(db.transaction_commit_queue.len(), 1);
}
#[test]
fn cleanup_respects_active_reader() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..5 {
let mut tx = db.transaction(true);
tx.set(format!("pre{i}"), "value").unwrap();
tx.commit().unwrap();
}
let reader = db.transaction(false);
for i in 0..5 {
let mut tx = db.transaction(true);
tx.set(format!("post{i}"), "value").unwrap();
tx.commit().unwrap();
}
assert_eq!(db.transaction_commit_queue.len(), 10);
db.run_cleanup();
assert_eq!(db.transaction_commit_queue.len(), 6);
drop(reader);
db.run_cleanup();
assert_eq!(db.transaction_commit_queue.len(), 1);
}
#[test]
fn logical_clock_is_dense() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..10 {
let mut tx = db.transaction(true);
tx.set(format!("key{i}"), "value").unwrap();
tx.commit().unwrap();
}
assert_eq!(db.oracle.timestamp.load(std::sync::atomic::Ordering::SeqCst), 10);
let entry = db.datastore.get(b"key0".as_slice()).expect("key0 missing");
let guard = entry.value().read();
assert_eq!(guard.as_slice()[0].version, 1);
let entry = db.datastore.get(b"key9".as_slice()).expect("key9 missing");
let guard = entry.value().read();
assert_eq!(guard.as_slice()[0].version, 10);
}
#[test]
fn logical_clock_continues_above_seed() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
let seed = 1_700_000_000_000_000_000u64;
db.oracle.alloc.store(seed, std::sync::atomic::Ordering::SeqCst);
db.oracle.timestamp.store(seed, std::sync::atomic::Ordering::SeqCst);
db.merge_retire_id.store(seed, std::sync::atomic::Ordering::SeqCst);
let mut tx = db.transaction(true);
tx.set("key", "value").unwrap();
tx.commit().unwrap();
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
let guard = entry.value().read();
assert_eq!(guard.as_slice()[0].version, seed + 1);
let mut tx = db.transaction(false);
assert_eq!(tx.get("key").unwrap().as_deref(), Some(b"value" as &[u8]));
tx.cancel().unwrap();
}
#[test]
fn slot_lifecycle() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
assert_eq!(db.readers.len(), 0);
let tx = db.transaction(false);
assert_eq!(db.readers.len(), 1);
let first_id = *db.readers.front().unwrap().key();
drop(tx);
assert_eq!(db.readers.len(), 0);
let tx = db.transaction(false);
assert_eq!(db.readers.len(), 1);
let second_id = *db.readers.front().unwrap().key();
assert!(second_id > first_id);
drop(tx);
assert_eq!(db.readers.len(), 0);
}
#[test]
fn watermark_conservative_while_pinning() {
use crate::inner::Slot;
use std::sync::Arc;
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..5 {
let mut tx = db.transaction(true);
tx.set(format!("key{i}"), "value").unwrap();
tx.commit().unwrap();
}
db.readers.insert(u64::MAX, Arc::new(Slot::pinning()));
assert_eq!(db.compute_cleanup_ts(), None);
let before = db.transaction_commit_queue.len();
db.run_cleanup();
assert_eq!(db.transaction_commit_queue.len(), before);
db.readers.remove(&u64::MAX);
assert!(db.compute_cleanup_ts().is_some());
db.run_cleanup();
assert_eq!(db.transaction_commit_queue.len(), 1);
}
#[test]
fn commit_watermark_tracks_commits() {
use std::sync::atomic::Ordering;
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..5 {
let mut tx = db.transaction(true);
tx.set(format!("key{i}"), "value").unwrap();
tx.commit().unwrap();
}
assert_eq!(db.commit_watermark.load(Ordering::SeqCst), 5);
assert_eq!(db.transaction_commit_id.load(Ordering::SeqCst), 5);
let mut tx1 = db.transaction(true);
let mut tx2 = db.transaction(true);
tx1.set("clash", "one").unwrap();
tx2.set("clash", "two").unwrap();
tx1.commit().unwrap();
assert!(tx2.commit().is_err());
assert_eq!(
db.commit_watermark.load(Ordering::SeqCst),
db.transaction_commit_id.load(Ordering::SeqCst)
);
}
#[test]
fn inline_gc_trims_hot_key_at_commit() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..100 {
let mut tx = db.transaction(true);
tx.set("hotkey", format!("v{i}")).unwrap();
tx.commit().unwrap();
}
let entry = db.datastore.get(b"hotkey".as_slice()).expect("hotkey missing");
let guard = entry.value().read();
let chain = guard.as_slice();
assert_eq!(chain.len(), 1, "inline GC should trim superseded versions at commit");
assert_eq!(chain[0].value.as_deref(), Some(b"v99" as &[u8]));
}
#[test]
fn inline_gc_unlinks_deleted_key_at_commit() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
{
let mut tx = db.transaction(true);
tx.set("key", "value").unwrap();
tx.commit().unwrap();
}
assert!(db.datastore.get(b"key".as_slice()).is_some());
{
let mut tx = db.transaction(true);
tx.del("key").unwrap();
tx.commit().unwrap();
}
assert!(
db.datastore.get(b"key".as_slice()).is_none(),
"a delete with no readers should unlink the node at commit"
);
}
#[test]
fn safety_net_reclaims_after_reader_departs() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
{
let mut tx = db.transaction(true);
tx.set("key", "v0").unwrap();
tx.commit().unwrap();
}
let reader = db.transaction(false);
for i in 1..=50 {
let mut tx = db.transaction(true);
tx.set("key", format!("v{i}")).unwrap();
tx.commit().unwrap();
}
{
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
let guard = entry.value().read();
assert!(guard.as_slice().len() > 1, "the pinned reader should retain history");
}
drop(reader);
db.run_gc();
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
let guard = entry.value().read();
let chain = guard.as_slice();
assert_eq!(chain.len(), 1, "the safety-net sweep should reclaim departed-reader garbage");
assert_eq!(chain[0].value.as_deref(), Some(b"v50" as &[u8]));
}
#[test]
fn tracked_sweep_reclaims_departed_reader_garbage() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
{
let mut tx = db.transaction(true);
tx.set("key", "v0").unwrap();
tx.commit().unwrap();
}
let reader = db.transaction(false);
for i in 1..=50 {
let mut tx = db.transaction(true);
tx.set("key", format!("v{i}")).unwrap();
tx.commit().unwrap();
}
assert!(
db.gc_candidates.pin().contains(b"key".as_slice()),
"a commit that leaves garbage should track its key"
);
db.run_gc_tracked();
{
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
assert!(entry.value().read().as_slice().len() > 1);
assert!(
db.gc_candidates.pin().contains(b"key".as_slice()),
"a still-pinned chain should stay tracked after a sweep"
);
}
drop(reader);
db.run_gc_tracked();
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
let guard = entry.value().read();
let chain = guard.as_slice();
assert_eq!(chain.len(), 1, "the tracked sweep should reclaim departed-reader garbage");
assert_eq!(chain[0].value.as_deref(), Some(b"v50" as &[u8]));
drop(guard);
assert!(
!db.gc_candidates.pin().contains(b"key".as_slice()),
"a terminal chain should be untracked after the sweep"
);
}
#[test]
fn tracked_sweep_unlinks_pinned_tombstone() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
{
let mut tx = db.transaction(true);
tx.set("key", "value").unwrap();
tx.commit().unwrap();
}
let reader = db.transaction(false);
{
let mut tx = db.transaction(true);
tx.del("key").unwrap();
tx.commit().unwrap();
}
assert!(db.datastore.get(b"key".as_slice()).is_some());
assert!(db.gc_candidates.pin().contains(b"key".as_slice()));
drop(reader);
db.run_gc_tracked();
assert!(
db.datastore.get(b"key".as_slice()).is_none(),
"the tracked sweep should unlink a departed-reader tombstone"
);
assert!(!db.gc_candidates.pin().contains(b"key".as_slice()));
}
#[test]
fn full_scan_supersedes_candidate_set() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
{
let mut tx = db.transaction(true);
tx.set("key", "v0").unwrap();
tx.commit().unwrap();
}
let reader = db.transaction(false);
{
let mut tx = db.transaction(true);
tx.set("key", "v1").unwrap();
tx.commit().unwrap();
}
assert!(db.gc_candidates.pin().contains(b"key".as_slice()));
drop(reader);
db.run_gc();
assert_eq!(db.gc_candidates.pin().len(), 0, "the full scan should clear the candidate set");
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
assert_eq!(entry.value().read().as_slice().len(), 1);
}
#[cfg(not(target_arch = "wasm32"))]
#[test]
fn load_collapses_aol_replay_chains() {
let temp_dir = tempfile::TempDir::new().unwrap();
let persistence = crate::PersistenceOptions::new(temp_dir.path())
.with_aol_mode(crate::AolMode::SynchronousOnCommit)
.with_snapshot_mode(crate::SnapshotMode::Never)
.with_fsync_mode(crate::FsyncMode::EveryAppend);
{
let db = Database::new_with_persistence(
crate::DatabaseOptions::default(),
persistence.clone(),
)
.unwrap();
for i in 0..5 {
let mut tx = db.transaction(true);
tx.set("key", format!("v{i}")).unwrap();
tx.commit().unwrap();
}
}
let db =
Database::new_with_persistence(crate::DatabaseOptions::default(), persistence).unwrap();
let entry = db.datastore.get(b"key".as_slice()).expect("key missing after reload");
let guard = entry.value().read();
let chain = guard.as_slice();
assert_eq!(chain.len(), 1, "the load-time sweep should collapse replayed chains");
assert_eq!(chain[0].value.as_deref(), Some(b"v4" as &[u8]));
}
#[test]
fn full_scan_retracks_pinned_chains() {
let db = Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
{
let mut tx = db.transaction(true);
tx.set("key", "v0").unwrap();
tx.commit().unwrap();
}
let reader = db.transaction(false);
{
let mut tx = db.transaction(true);
tx.set("key", "v1").unwrap();
tx.commit().unwrap();
}
db.run_gc();
assert!(
db.gc_candidates.pin().contains(b"key".as_slice()),
"the full scan should re-track a chain it could not trim"
);
drop(reader);
db.run_gc_tracked();
let entry = db.datastore.get(b"key".as_slice()).expect("key missing");
assert_eq!(entry.value().read().as_slice().len(), 1);
}
#[test]
fn pin_then_read_stress() {
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
let db = Arc::new(Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
));
{
let mut tx = db.transaction(true);
tx.set("key", "v0").unwrap();
tx.commit().unwrap();
}
let stop = Arc::new(AtomicBool::new(false));
let none_reads = Arc::new(AtomicUsize::new(0));
let total_reads = Arc::new(AtomicUsize::new(0));
let writer = {
let db = Arc::clone(&db);
let stop = Arc::clone(&stop);
thread::spawn(move || {
let mut counter: u64 = 0;
while !stop.load(Ordering::Relaxed) {
let mut tx = db.transaction(true);
tx.set("key", format!("v{counter}")).unwrap();
tx.commit().unwrap();
counter = counter.wrapping_add(1);
}
})
};
let sweeper = {
let db = Arc::clone(&db);
let stop = Arc::clone(&stop);
thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
db.run_cleanup();
db.run_gc();
db.run_gc_tracked();
}
})
};
let mut readers = Vec::new();
for _ in 0..6 {
let db = Arc::clone(&db);
let stop = Arc::clone(&stop);
let none_reads = Arc::clone(&none_reads);
let total_reads = Arc::clone(&total_reads);
readers.push(thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
let mut tx = db.transaction(false);
let value = tx.get("key").unwrap();
total_reads.fetch_add(1, Ordering::Relaxed);
if value.is_none() {
none_reads.fetch_add(1, Ordering::Relaxed);
}
tx.cancel().unwrap();
}
}));
}
let started = Instant::now();
while started.elapsed() < Duration::from_millis(500) {
thread::sleep(Duration::from_millis(10));
}
stop.store(true, Ordering::Relaxed);
writer.join().unwrap();
sweeper.join().unwrap();
for r in readers {
r.join().unwrap();
}
let nones = none_reads.load(Ordering::Relaxed);
let total = total_reads.load(Ordering::Relaxed);
assert_eq!(
nones, 0,
"reader observed `None` for a key that always has a committed \
value ({nones} of {total} reads): the pin-then-read protocol \
failed to protect a registering reader from a sweep"
);
}
}