use exonum_profiler::ProfilerSpan;
use rocksdb::{self, WriteBatch, DBIterator};
use rocksdb::utils::get_cf_names;
use std::mem;
use std::sync::Arc;
use std::path::Path;
use std::fmt;
use std::error::Error;
use std::iter::Peekable;
use storage::{self, Database, Iterator, Iter, Snapshot, Patch};
use storage::db::Change;
pub use rocksdb::{Options as RocksDBOptions, WriteOptions as RocksDBWriteOptions};
pub use rocksdb::BlockBasedOptions as RocksBlockOptions;
impl From<rocksdb::Error> for storage::Error {
fn from(err: rocksdb::Error) -> storage::Error {
storage::Error::new(err.description())
}
}
pub struct RocksDB {
db: Arc<rocksdb::DB>,
}
pub struct RocksDBSnapshot {
snapshot: rocksdb::Snapshot<'static>,
_db: Arc<rocksdb::DB>,
}
struct RocksDBIterator {
iter: Peekable<DBIterator>,
key: Option<Box<[u8]>>,
value: Option<Box<[u8]>>,
}
impl RocksDB {
pub fn open<P: AsRef<Path>>(path: P, options: &RocksDBOptions) -> storage::Result<RocksDB> {
let db = {
if let Ok(names) = get_cf_names(&path) {
let cf_names = names.iter().map(|name| name.as_str()).collect::<Vec<_>>();
rocksdb::DB::open_cf(options, path, cf_names.as_ref())?
} else {
rocksdb::DB::open(options, path)?
}
};
Ok(RocksDB { db: Arc::new(db) })
}
fn do_merge(&self, patch: Patch, w_opts: &RocksDBWriteOptions) -> storage::Result<()> {
let _p = ProfilerSpan::new("RocksDB::merge");
let mut batch = WriteBatch::default();
for (cf_name, changes) in patch {
let cf = match self.db.cf_handle(&cf_name) {
Some(cf) => cf,
None => {
self.db
.create_cf(&cf_name, &RocksDBOptions::default())
.unwrap()
}
};
for (key, change) in changes {
match change {
Change::Put(ref value) => batch.put_cf(cf, key.as_ref(), value)?,
Change::Delete => batch.delete_cf(cf, &key)?,
}
}
}
self.db.write_opt(batch, w_opts).map_err(Into::into)
}
}
impl Database for RocksDB {
fn snapshot(&self) -> Box<Snapshot> {
let _p = ProfilerSpan::new("RocksDB::snapshot");
Box::new(RocksDBSnapshot {
snapshot: unsafe { mem::transmute(self.db.snapshot()) },
_db: Arc::clone(&self.db),
})
}
fn merge(&self, patch: Patch) -> storage::Result<()> {
let w_opts = RocksDBWriteOptions::default();
self.do_merge(patch, &w_opts)
}
fn merge_sync(&self, patch: Patch) -> storage::Result<()> {
let mut w_opts = RocksDBWriteOptions::default();
w_opts.set_sync(true);
self.do_merge(patch, &w_opts)
}
}
impl Snapshot for RocksDBSnapshot {
fn get(&self, name: &str, key: &[u8]) -> Option<Vec<u8>> {
let _p = ProfilerSpan::new("RocksDBSnapshot::get");
if let Some(cf) = self._db.cf_handle(name) {
match self.snapshot.get_cf(cf, key) {
Ok(value) => value.map(|v| v.to_vec()),
Err(e) => panic!(e),
}
} else {
None
}
}
fn iter<'a>(&'a self, name: &str, from: &[u8]) -> Iter<'a> {
use rocksdb::{IteratorMode, Direction};
let _p = ProfilerSpan::new("RocksDBSnapshot::iter");
let iter = match self._db.cf_handle(name) {
Some(cf) => {
self.snapshot
.iterator_cf(cf, IteratorMode::From(from, Direction::Forward))
.unwrap()
}
None => self.snapshot.iterator(IteratorMode::Start),
};
Box::new(RocksDBIterator {
iter: iter.peekable(),
key: None,
value: None,
})
}
}
impl Iterator for RocksDBIterator {
fn next(&mut self) -> Option<(&[u8], &[u8])> {
let _p = ProfilerSpan::new("RocksDBIterator::next");
if let Some((key, value)) = self.iter.next() {
self.key = Some(key);
self.value = Some(value);
Some((self.key.as_ref().unwrap(), self.value.as_ref().unwrap()))
} else {
None
}
}
fn peek(&mut self) -> Option<(&[u8], &[u8])> {
let _p = ProfilerSpan::new("RocksDBIterator::peek");
if let Some(&(ref key, ref value)) = self.iter.peek() {
Some((key, value))
} else {
None
}
}
}
impl From<RocksDB> for Arc<Database> {
fn from(db: RocksDB) -> Arc<Database> {
Arc::from(Box::new(db) as Box<Database>)
}
}
impl fmt::Debug for RocksDB {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(f, "RocksDB(..)")
}
}
impl fmt::Debug for RocksDBSnapshot {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
write!(f, "RocksDBSnapshot(..)")
}
}