pub mod btree;
pub mod buffer_pool;
pub mod pager;
use crate::config::Config;
use crate::wal::WriteAheadLog;
use btree::BTree;
use buffer_pool::BufferPool;
use pager::Pager;
use parking_lot::Mutex;
pub struct Database {
config: Config,
wal: Option<Mutex<WriteAheadLog>>,
index: BTree,
}
impl Database {
pub fn open(config: Config) -> anyhow::Result<Self> {
if !config.storage.path.exists() {
std::fs::create_dir_all(&config.storage.path)?;
}
let data_path = config.storage.path.join("magnum.data");
let pager = Pager::open(&data_path)?;
let checkpoint_lsn = {
let mut tmp_pager = Pager::open(&data_path)?;
tmp_pager.read_checkpoint_lsn().unwrap_or(0)
};
let buffer_pool = BufferPool::new(pager, 1024)
.with_sync_interval(config.storage.sync_interval);
let mut index = BTree::new(buffer_pool)?;
let wal = if config.wal.enabled {
Some(Mutex::new(WriteAheadLog::open(&config.storage.path)?))
} else {
None
};
if let Some(wal_mutex) = &wal {
let mut wal_ref = wal_mutex.lock();
let entries = wal_ref.recover(checkpoint_lsn)?;
for entry in entries {
match entry {
crate::wal::WalEntry::Put(k, v) => {
index.insert(&k, &v)?;
}
crate::wal::WalEntry::Delete(k) => {
index.delete(&k)?;
}
}
}
if !entries_is_empty_hint(&*wal_ref) {
index.flush_and_sync()?;
}
}
Ok(Self { config, wal, index })
}
pub fn put(&self, key: &[u8], value: &[u8]) -> anyhow::Result<()> {
self.put_with_tx(0, key, value)
}
pub fn put_with_tx(&self, tx_id: u64, key: &[u8], value: &[u8]) -> anyhow::Result<()> {
if let Some(wal_mutex) = &self.wal {
let mut wal = wal_mutex.lock();
wal.append_tx_put(tx_id, key, value)?;
if self.config.wal.sync_on_write {
wal.sync()?;
}
}
self.index.insert(key, value)?;
Ok(())
}
pub fn get(&self, key: &[u8]) -> anyhow::Result<Option<Vec<u8>>> {
self.index.search(key)
}
pub fn scan(&self) -> anyhow::Result<Vec<(Vec<u8>, Vec<u8>)>> {
self.index.scan()
}
pub fn scan_prefix(&self, prefix: &[u8]) -> anyhow::Result<Vec<(Vec<u8>, Vec<u8>)>> {
self.index.scan_prefix(prefix)
}
pub fn delete(&self, key: &[u8]) -> anyhow::Result<()> {
self.delete_with_tx(0, key)
}
pub fn delete_with_tx(&self, tx_id: u64, key: &[u8]) -> anyhow::Result<()> {
if let Some(wal_mutex) = &self.wal {
let mut wal = wal_mutex.lock();
wal.append_tx_delete(tx_id, key)?;
if self.config.wal.sync_on_write {
wal.sync()?;
}
}
self.index.delete(key)?;
Ok(())
}
pub fn begin_tx(&self, tx_id: u64) -> anyhow::Result<()> {
if let Some(wal_mutex) = &self.wal {
let mut wal = wal_mutex.lock();
wal.append_begin(tx_id)?;
if self.config.wal.sync_on_write {
wal.sync()?;
}
}
Ok(())
}
pub fn commit_tx(&self, tx_id: u64) -> anyhow::Result<()> {
if let Some(wal_mutex) = &self.wal {
let mut wal = wal_mutex.lock();
wal.append_commit(tx_id)?;
wal.sync()?;
}
self.index.flush_and_sync()?;
Ok(())
}
pub fn rollback_tx(&self, tx_id: u64) -> anyhow::Result<()> {
if let Some(wal_mutex) = &self.wal {
let mut wal = wal_mutex.lock();
wal.append_rollback(tx_id)?;
if self.config.wal.sync_on_write {
wal.sync()?;
}
}
Ok(())
}
pub fn close(mut self) -> anyhow::Result<()> {
self.index.flush_all()?;
self.index.sync()?;
if let Some(wal_mutex) = &self.wal {
let mut wal = wal_mutex.lock();
let current_lsn = wal.current_lsn();
self.index
.buffer_pool()
.write_checkpoint_lsn(current_lsn)?;
self.index.sync()?;
wal.checkpoint()?;
}
Ok(())
}
}
fn entries_is_empty_hint(_wal: &WriteAheadLog) -> bool {
false
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::tempdir;
fn get_config(path: std::path::PathBuf, wal_enabled: bool) -> Config {
let mut config = Config::default();
config.storage.path = path;
config.wal.enabled = wal_enabled;
config.storage.sync_interval = 0; config
}
#[test]
fn test_database_close_wal_disabled() {
let dir = tempdir().unwrap();
let config = get_config(dir.path().to_path_buf(), false);
{
let mut db = Database::open(config.clone()).unwrap();
db.put(b"hello", b"world").unwrap();
db.close().unwrap(); }
{
let mut db = Database::open(config).unwrap();
let val = db.get(b"hello").unwrap().unwrap();
assert_eq!(val, b"world");
}
}
#[test]
fn test_database_wal_checkpoint() {
let dir = tempdir().unwrap();
let config = get_config(dir.path().to_path_buf(), true);
let wal_path = config.storage.path.join("magnum.wal");
{
let mut db = Database::open(config.clone()).unwrap();
for i in 0..100 {
let k = format!("k{}", i);
let v = format!("v{}", i);
db.put(k.as_bytes(), v.as_bytes()).unwrap();
}
db.close().unwrap(); }
assert_eq!(std::fs::metadata(&wal_path).unwrap().len(), 0);
{
let mut db = Database::open(config).unwrap();
for i in 0..100 {
let k = format!("k{}", i);
let v = format!("v{}", i);
let val = db.get(k.as_bytes()).unwrap().unwrap();
assert_eq!(val, v.as_bytes());
}
}
}
#[test]
fn test_database_tx_commit_durable() {
let dir = tempdir().unwrap();
let config = get_config(dir.path().to_path_buf(), true);
{
let mut db = Database::open(config.clone()).unwrap();
db.begin_tx(1).unwrap();
db.put_with_tx(1, b"txkey", b"txval").unwrap();
db.commit_tx(1).unwrap();
}
{
let mut db = Database::open(config).unwrap();
let val = db.get(b"txkey").unwrap();
assert!(val.is_some());
assert_eq!(val.unwrap(), b"txval");
}
}
}