use crate::storage_engine::{DatabaseObject, Operations, StorageEngine};
use dragonfly_client_core::{
error::{ErrorType, OrErr},
Error, Result,
};
use rocksdb::WriteOptions;
use std::{
ops::Deref,
path::{Path, PathBuf},
};
use tracing::{info, warn};
pub struct RocksdbStorageEngine {
inner: rocksdb::DB,
}
impl Deref for RocksdbStorageEngine {
type Target = rocksdb::DB;
fn deref(&self) -> &Self::Target {
&self.inner
}
}
impl RocksdbStorageEngine {
const DEFAULT_DIR_NAME: &'static str = "metadata";
const DEFAULT_MEMTABLE_MEMORY_BUDGET: usize = 64 * 1024 * 1024;
const DEFAULT_PREFIX_CF_MEMTABLE_MEMORY_BUDGET: usize = 512 * 1024 * 1024;
const DEFAULT_MAX_BACKGROUND_JOBS: i32 = 2;
const DEFAULT_BLOCK_SIZE: usize = 64 * 1024;
const DEFAULT_CACHE_SIZE: usize = 1024 * 1024 * 1024;
const DEFAULT_LOG_MAX_SIZE: usize = 64 * 1024 * 1024;
const DEFAULT_LOG_MAX_FILES: usize = 10;
const DEFAULT_BYTES_PER_SYNC: u64 = 2 * 1024 * 1024;
pub fn open(
dir: &Path,
log_dir: &PathBuf,
cf_names: &[&str],
prefix_cf_names: &[&str],
keep: bool,
) -> Result<Self> {
info!("initializing metadata directory: {:?} {:?}", dir, cf_names);
let mut options = rocksdb::Options::default();
options.create_if_missing(true);
options.create_missing_column_families(true);
options.set_compression_type(rocksdb::DBCompressionType::Lz4);
options.set_bottommost_compression_type(rocksdb::DBCompressionType::Lz4);
let parallelism = std::thread::available_parallelism().map_or(1, |n| n.get()) as i32;
options.increase_parallelism(parallelism);
options.set_max_background_jobs(std::cmp::max(
parallelism,
Self::DEFAULT_MAX_BACKGROUND_JOBS,
));
options.set_use_fsync(false);
options.set_bytes_per_sync(Self::DEFAULT_BYTES_PER_SYNC);
options.set_db_log_dir(log_dir);
options.set_log_level(rocksdb::LogLevel::Info);
options.set_max_log_file_size(Self::DEFAULT_LOG_MAX_SIZE);
options.set_keep_log_file_num(Self::DEFAULT_LOG_MAX_FILES);
let mut block_options = rocksdb::BlockBasedOptions::default();
block_options.set_block_cache(&rocksdb::Cache::new_lru_cache(Self::DEFAULT_CACHE_SIZE));
block_options.set_block_size(Self::DEFAULT_BLOCK_SIZE);
block_options.set_cache_index_and_filter_blocks(true);
block_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
block_options.set_bloom_filter(10.0, false);
options.set_block_based_table_factory(&block_options);
let mut prefix_cf_options = rocksdb::Options::default();
prefix_cf_options.set_prefix_extractor(rocksdb::SliceTransform::create_fixed_prefix(64));
prefix_cf_options.set_memtable_prefix_bloom_ratio(0.25);
prefix_cf_options
.optimize_level_style_compaction(Self::DEFAULT_PREFIX_CF_MEMTABLE_MEMORY_BUDGET);
prefix_cf_options.set_block_based_table_factory(&block_options);
let mut cf_options = rocksdb::Options::default();
cf_options.optimize_level_style_compaction(Self::DEFAULT_MEMTABLE_MEMORY_BUDGET);
cf_options.set_block_based_table_factory(&block_options);
let cfs = cf_names
.iter()
.map(|name| (name.to_string(), cf_options.clone()))
.chain(
prefix_cf_names
.iter()
.map(|name| (name.to_string(), prefix_cf_options.clone())),
)
.collect::<Vec<_>>();
let dir = dir.join(Self::DEFAULT_DIR_NAME);
if !keep {
rocksdb::DB::destroy(&options, &dir).unwrap_or_else(|err| {
warn!("destroy {:?} failed: {}", dir, err);
});
}
let db =
rocksdb::DB::open_cf_with_opts(&options, &dir, cfs).or_err(ErrorType::StorageError)?;
Ok(Self { inner: db })
}
}
impl Operations for RocksdbStorageEngine {
fn get<O: DatabaseObject>(&self, key: &[u8]) -> Result<Option<O>> {
let cf = cf_handle::<O>(self)?;
let value = self.get_cf(cf, key).or_err(ErrorType::StorageError)?;
match value {
Some(value) => Ok(Some(O::deserialize_from(&value)?)),
None => Ok(None),
}
}
fn multi_get<O: DatabaseObject>(&self, keys: &[&[u8]]) -> Result<Vec<Option<O>>> {
let cf = cf_handle::<O>(self)?;
self.batched_multi_get_cf(cf, keys, false)
.into_iter()
.map(|value| match value.or_err(ErrorType::StorageError)? {
Some(value) => Ok(Some(O::deserialize_from(&value)?)),
None => Ok(None),
})
.collect()
}
fn exists<O: DatabaseObject>(&self, key: &[u8]) -> Result<bool> {
let cf = cf_handle::<O>(self)?;
Ok(self
.get_cf(cf, key)
.or_err(ErrorType::StorageError)?
.is_some())
}
fn put<O: DatabaseObject>(&self, key: &[u8], value: &O) -> Result<()> {
let cf = cf_handle::<O>(self)?;
let mut options = rocksdb::WriteOptions::default();
options.set_sync(false);
self.put_cf_opt(cf, key, value.serialized()?, &options)
.or_err(ErrorType::StorageError)?;
Ok(())
}
fn delete<O: DatabaseObject>(&self, key: &[u8]) -> Result<()> {
let cf = cf_handle::<O>(self)?;
let mut options = WriteOptions::default();
options.set_sync(false);
self.delete_cf_opt(cf, key, &options)
.or_err(ErrorType::StorageError)?;
Ok(())
}
fn iter<O: DatabaseObject>(&self) -> Result<impl Iterator<Item = Result<(Box<[u8]>, O)>>> {
let cf = cf_handle::<O>(self)?;
let mut options = rocksdb::ReadOptions::default();
options.fill_cache(false);
let iter = self.iterator_cf_opt(cf, options, rocksdb::IteratorMode::Start);
Ok(iter.map(|ele| {
let (key, value) = ele.or_err(ErrorType::StorageError)?;
Ok((key, O::deserialize_from(&value)?))
}))
}
fn iter_raw<O: DatabaseObject>(
&self,
) -> Result<impl Iterator<Item = Result<(Box<[u8]>, Box<[u8]>)>>> {
let cf = cf_handle::<O>(self)?;
let mut options = rocksdb::ReadOptions::default();
options.fill_cache(false);
Ok(self
.iterator_cf_opt(cf, options, rocksdb::IteratorMode::Start)
.map(|ele| {
let (key, value) = ele.or_err(ErrorType::StorageError)?;
Ok((key, value))
}))
}
fn prefix_iter<O: DatabaseObject>(
&self,
prefix: &[u8],
) -> Result<impl Iterator<Item = Result<(Box<[u8]>, O)>>> {
let cf = cf_handle::<O>(self)?;
let iter = self.prefix_iterator_cf(cf, prefix);
Ok(iter.map(|ele| {
let (key, value) = ele.or_err(ErrorType::StorageError)?;
Ok((key, O::deserialize_from(&value)?))
}))
}
fn prefix_iter_raw<O: DatabaseObject>(
&self,
prefix: &[u8],
) -> Result<impl Iterator<Item = Result<(Box<[u8]>, Box<[u8]>)>>> {
let cf = cf_handle::<O>(self)?;
Ok(self.prefix_iterator_cf(cf, prefix).map(|ele| {
let (key, value) = ele.or_err(ErrorType::StorageError)?;
Ok((key, value))
}))
}
fn batch_delete<O: DatabaseObject>(&self, keys: Vec<&[u8]>) -> Result<()> {
let cf = cf_handle::<O>(self)?;
let mut batch = rocksdb::WriteBatch::default();
for key in keys {
batch.delete_cf(cf, key);
}
let mut options = WriteOptions::default();
options.set_sync(false);
Ok(self
.write_opt(batch, &options)
.or_err(ErrorType::StorageError)?)
}
}
impl StorageEngine<'_> for RocksdbStorageEngine {}
fn cf_handle<T>(db: &rocksdb::DB) -> Result<&rocksdb::ColumnFamily>
where
T: DatabaseObject,
{
let cf_name = T::NAMESPACE;
db.cf_handle(cf_name)
.ok_or_else(|| Error::ColumnFamilyNotFound(cf_name.to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
use serde::{Deserialize, Serialize};
use tempfile::tempdir;
#[derive(Debug, Serialize, Deserialize, PartialEq, Clone)]
struct Object {
id: String,
value: i32,
}
impl DatabaseObject for Object {
const NAMESPACE: &'static str = "object";
}
#[derive(Debug, Serialize, Deserialize, PartialEq)]
struct UnregisteredObject {
data: String,
}
impl DatabaseObject for UnregisteredObject {
const NAMESPACE: &'static str = "unregistered";
}
fn open(dir: &Path, keep: bool) -> RocksdbStorageEngine {
RocksdbStorageEngine::open(dir, &dir.to_path_buf(), &[], &[Object::NAMESPACE], keep)
.unwrap()
}
#[test]
fn get_and_exists_follow_put_and_delete() {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
let key = b"1";
assert_eq!(engine.get::<Object>(key).unwrap(), None);
assert!(!engine.exists::<Object>(key).unwrap());
let object = Object {
id: "1".to_string(),
value: 42,
};
engine.put(key, &object).unwrap();
assert_eq!(engine.get::<Object>(key).unwrap(), Some(object));
assert!(engine.exists::<Object>(key).unwrap());
let object = Object {
id: "1".to_string(),
value: 43,
};
engine.put(key, &object).unwrap();
assert_eq!(engine.get::<Object>(key).unwrap(), Some(object));
engine.delete::<Object>(key).unwrap();
assert_eq!(engine.get::<Object>(key).unwrap(), None);
assert!(!engine.exists::<Object>(key).unwrap());
engine.delete::<Object>(key).unwrap();
assert!(!engine.exists::<Object>(key).unwrap());
}
#[test]
fn multi_get_returns_values_in_key_order() {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
for (id, value) in [("1", 1), ("2", 2)] {
engine
.put(
id.as_bytes(),
&Object {
id: id.to_string(),
value,
},
)
.unwrap();
}
let test_cases = vec![
(vec![], vec![]),
(
vec!["2", "missing", "1"],
vec![
Some(Object {
id: "2".to_string(),
value: 2,
}),
None,
Some(Object {
id: "1".to_string(),
value: 1,
}),
],
),
];
for (ids, expected) in test_cases {
let keys: Vec<&[u8]> = ids.iter().map(|id| id.as_bytes()).collect();
assert_eq!(engine.multi_get::<Object>(&keys).unwrap(), expected);
}
}
#[test]
fn batch_delete_removes_only_the_given_keys() {
let test_cases = vec![
(vec![], vec!["1", "2", "3"]),
(vec!["missing"], vec!["1", "2", "3"]),
(vec!["1", "3"], vec!["2"]),
(vec!["1", "2", "3"], vec![]),
];
for (deleted_ids, expected_ids) in test_cases {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
for (id, value) in [("1", 1), ("2", 2), ("3", 3)] {
engine
.put(
id.as_bytes(),
&Object {
id: id.to_string(),
value,
},
)
.unwrap();
}
let keys: Vec<&[u8]> = deleted_ids.iter().map(|id| id.as_bytes()).collect();
engine.batch_delete::<Object>(keys).unwrap();
let stored_ids: Vec<String> = engine
.iter::<Object>()
.unwrap()
.map(|ele| {
let (_, object) = ele.unwrap();
object.id
})
.collect();
assert_eq!(stored_ids, expected_ids);
}
}
#[test]
fn iter_yields_objects_in_key_order() {
let test_cases = vec![
(vec![], vec![]),
(
vec![
Object {
id: "3".to_string(),
value: 30,
},
Object {
id: "1".to_string(),
value: 10,
},
Object {
id: "2".to_string(),
value: 20,
},
],
vec![
Object {
id: "1".to_string(),
value: 10,
},
Object {
id: "2".to_string(),
value: 20,
},
Object {
id: "3".to_string(),
value: 30,
},
],
),
];
for (objects, expected) in test_cases {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
for object in &objects {
engine.put(object.id.as_bytes(), object).unwrap();
}
let expected: Vec<(Box<[u8]>, Object)> = expected
.iter()
.map(|object| (object.id.as_bytes().into(), object.clone()))
.collect();
let iterated = engine
.iter::<Object>()
.unwrap()
.collect::<Result<Vec<_>>>()
.unwrap();
let raw_iterated: Vec<(Box<[u8]>, Object)> = engine
.iter_raw::<Object>()
.unwrap()
.map(|ele| {
let (key, value) = ele.unwrap();
(key, Object::deserialize_from(&value).unwrap())
})
.collect();
assert_eq!(iterated, expected);
assert_eq!(raw_iterated, expected);
}
}
#[test]
fn prefix_iter_yields_only_objects_under_the_prefix() {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
for (prefix_char, suffix, id, value) in [
("a", "_suffix1", "a1", 100),
("a", "_suffix2", "a2", 200),
("b", "_suffix1", "b1", 300),
("b", "_suffix2", "b2", 400),
] {
engine
.put(
format!("{}{suffix}", prefix_char.repeat(64)).as_bytes(),
&Object {
id: id.to_string(),
value,
},
)
.unwrap();
}
let test_cases = vec![
(
"a",
vec![
Object {
id: "a1".to_string(),
value: 100,
},
Object {
id: "a2".to_string(),
value: 200,
},
],
),
(
"b",
vec![
Object {
id: "b1".to_string(),
value: 300,
},
Object {
id: "b2".to_string(),
value: 400,
},
],
),
("0", vec![]),
];
for (prefix_char, expected) in test_cases {
let prefix = prefix_char.repeat(64);
let iterated = engine
.prefix_iter::<Object>(prefix.as_bytes())
.unwrap()
.collect::<Result<Vec<_>>>()
.unwrap();
let raw_iterated: Vec<(Box<[u8]>, Object)> = engine
.prefix_iter_raw::<Object>(prefix.as_bytes())
.unwrap()
.map(|ele| {
let (key, value) = ele.unwrap();
(key, Object::deserialize_from(&value).unwrap())
})
.collect();
assert_eq!(raw_iterated, iterated);
assert!(iterated
.iter()
.all(|(key, _)| key.starts_with(prefix.as_bytes())));
let objects: Vec<Object> = iterated.into_iter().map(|(_, object)| object).collect();
assert_eq!(objects, expected);
}
}
#[test]
fn operations_fail_on_an_unregistered_column_family() {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
let test_cases: Vec<fn(&RocksdbStorageEngine) -> Result<()>> = vec![
|engine| engine.get::<UnregisteredObject>(b"1").map(|_| ()),
|engine| {
engine
.multi_get::<UnregisteredObject>(&[b"1".as_slice()])
.map(|_| ())
},
|engine| engine.exists::<UnregisteredObject>(b"1").map(|_| ()),
|engine| {
engine.put(
b"1",
&UnregisteredObject {
data: "1".to_string(),
},
)
},
|engine| engine.delete::<UnregisteredObject>(b"1"),
|engine| engine.iter::<UnregisteredObject>().map(|_| ()),
|engine| engine.iter_raw::<UnregisteredObject>().map(|_| ()),
|engine| engine.prefix_iter::<UnregisteredObject>(b"1").map(|_| ()),
|engine| {
engine
.prefix_iter_raw::<UnregisteredObject>(b"1")
.map(|_| ())
},
|engine| engine.batch_delete::<UnregisteredObject>(vec![b"1".as_slice()]),
];
for run in test_cases {
let result = run(&engine);
assert!(matches!(
result,
Err(Error::ColumnFamilyNotFound(ref name)) if name == UnregisteredObject::NAMESPACE
));
}
}
#[test]
fn open_keeps_or_destroys_the_existing_data() {
let test_cases = vec![
(
true,
Some(Object {
id: "1".to_string(),
value: 42,
}),
),
(false, None),
];
for (keep, expected) in test_cases {
let dir = tempdir().unwrap();
let engine = open(dir.path(), false);
engine
.put(
b"1",
&Object {
id: "1".to_string(),
value: 42,
},
)
.unwrap();
drop(engine);
let engine = open(dir.path(), keep);
assert_eq!(engine.get::<Object>(b"1").unwrap(), expected);
}
}
}