mod mmdb;
pub(crate) use self::mmdb::MmDB as Engine;
pub(crate) use self::mmdb::{EngineSizing, write_file_durable};
type DbIter = self::mmdb::MmdbIter;
use crate::common::{
PREFIX_SIZE, PreBytes, RawKey, RawValue,
error::{Result, VsdbError},
namespace::{DEFAULT_NS_ID, Namespace},
};
use serde::{Deserialize, Serialize, de};
use std::{
borrow::Cow,
fmt,
marker::PhantomData,
ops::{Deref, DerefMut, RangeBounds},
result::Result as StdResult,
sync::OnceLock,
};
const MAPX_META_MAGIC: &[u8; 8] = b"VSMAPX01";
const MAPX_META_LEN: usize = MAPX_META_MAGIC.len() + PREFIX_SIZE;
const MAPX_META_NS_LEN: usize = MAPX_META_LEN + size_of::<u64>();
pub trait BatchTrait {
fn insert(&mut self, key: &[u8], value: &[u8]);
fn remove(&mut self, key: &[u8]);
fn commit(&mut self) -> Result<()>;
}
#[derive(Debug)]
pub(crate) struct Mapx {
prefix: Prefix,
ns: Namespace,
}
#[derive(Debug)]
enum Prefix {
Recovered(PreBytes),
Created(OnceLock<PreBytes>),
}
impl Mapx {
#[inline(always)]
fn prefix_bytes(&self) -> PreBytes {
match &self.prefix {
Prefix::Recovered(bytes) => *bytes,
Prefix::Created(cell) => *cell.get_or_init(|| {
let prefix = self.ns.engine().alloc_prefix();
let prefix_bytes = prefix.to_le_bytes();
debug_assert!(self.ns.engine().iter(prefix_bytes).next().is_none());
prefix_bytes
}),
}
}
fn materialize(&mut self) -> PreBytes {
let b = self.prefix_bytes();
if matches!(self.prefix, Prefix::Created(_)) {
self.prefix = Prefix::Recovered(b);
}
b
}
}
impl Mapx {
pub(crate) unsafe fn shadow(&self) -> Self {
Self {
prefix: Prefix::Recovered(self.prefix_bytes()),
ns: self.ns.clone(),
}
}
#[inline(always)]
pub(crate) fn new() -> Self {
Self::new_in(&Namespace::current())
}
#[inline(always)]
pub(crate) fn new_in(ns: &Namespace) -> Self {
Self {
prefix: Prefix::Created(OnceLock::new()),
ns: ns.clone(),
}
}
#[inline(always)]
pub(crate) fn namespace(&self) -> Namespace {
self.ns.clone()
}
#[inline(always)]
pub(crate) fn get(&self, key: &[u8]) -> Option<RawValue> {
self.ns.engine().get(self.prefix_bytes(), key)
}
#[inline(always)]
pub(crate) fn get_mut(&mut self, key: &[u8]) -> Option<ValueMut<'_>> {
let v = self.ns.engine().get(self.materialize(), key)?;
Some(ValueMut {
key: key.to_vec(),
value: v,
dirty: false,
hdr: self,
})
}
#[inline(always)]
pub(crate) fn mock_value_mut(
&mut self,
key: RawValue,
value: RawValue,
) -> ValueMut<'_> {
ValueMut {
key,
value,
dirty: true,
hdr: self,
}
}
#[inline(always)]
pub(crate) fn iter(&self) -> MapxIter<'_> {
MapxIter {
db_iter: self.ns.engine().iter(self.prefix_bytes()),
_marker: PhantomData,
}
}
#[inline(always)]
pub(crate) fn iter_mut(&mut self) -> MapxIterMut<'_> {
MapxIterMut {
db_iter: self.ns.engine().iter(self.materialize()),
hdr: self,
}
}
#[inline(always)]
pub(crate) fn range<'a, R: RangeBounds<Cow<'a, [u8]>>>(
&'a self,
bounds: R,
) -> MapxIter<'a> {
MapxIter {
db_iter: self.ns.engine().range(self.prefix_bytes(), bounds),
_marker: PhantomData,
}
}
#[inline(always)]
pub(crate) fn range_detached<'a, R: RangeBounds<Cow<'a, [u8]>>>(
&self,
bounds: R,
) -> MapxIter<'a> {
MapxIter {
db_iter: self.ns.engine().range(self.prefix_bytes(), bounds),
_marker: PhantomData,
}
}
#[inline(always)]
pub(crate) fn range_mut<'a, R: RangeBounds<Cow<'a, [u8]>>>(
&'a mut self,
bounds: R,
) -> MapxIterMut<'a> {
MapxIterMut {
db_iter: self.ns.engine().range(self.materialize(), bounds),
hdr: self,
}
}
#[inline(always)]
pub(crate) fn insert(&mut self, key: &[u8], value: &[u8]) {
let prefix = self.materialize();
self.ns.engine().insert(prefix, key, value);
}
#[inline(always)]
pub(crate) fn remove(&mut self, key: &[u8]) {
let prefix = self.materialize();
self.ns.engine().remove(prefix, key);
}
#[inline(always)]
pub(crate) fn lazy_delete(&self, key: &[u8]) {
self.ns.engine().lazy_delete(self.prefix_bytes(), key);
}
#[inline(always)]
pub(crate) fn lazy_delete_batch(
&self,
keys: impl IntoIterator<Item = impl AsRef<[u8]>>,
) {
self.ns
.engine()
.lazy_delete_batch(self.prefix_bytes(), keys);
}
#[inline(always)]
pub(crate) fn batch_begin(&mut self) -> Box<dyn BatchTrait + '_> {
let prefix = self.materialize();
self.ns.engine().batch_begin(prefix)
}
#[inline(always)]
pub(crate) fn batch_begin_wiped(&mut self) -> Box<dyn BatchTrait + '_> {
let prefix = self.materialize();
self.ns.engine().batch_begin_wiped(prefix)
}
#[inline(always)]
pub(crate) fn clear(&mut self) {
self.batch_begin_wiped()
.commit()
.expect("vsdb: batch delete failed during clear");
}
#[inline(always)]
pub(crate) unsafe fn from_prefix_slice_in(
ns: &Namespace,
s: impl AsRef<[u8]>,
) -> Self {
debug_assert_eq!(s.as_ref().len(), PREFIX_SIZE);
let mut prefix = PreBytes::default();
prefix.copy_from_slice(s.as_ref());
Self {
prefix: Prefix::Recovered(prefix),
ns: ns.clone(),
}
}
#[inline(always)]
pub(crate) unsafe fn from_prefix_slice(s: impl AsRef<[u8]>) -> Self {
unsafe { Self::from_prefix_slice_in(&Namespace::current(), s) }
}
pub(crate) fn from_prefix_meta(meta: &[u8]) -> Result<Self> {
let (prefix, ns_id) = Self::decode_prefix_meta(meta)?;
let ns = match ns_id {
None => Namespace::default_ns(),
Some(id) => Namespace::open(id)?,
};
if !ns.engine().reserve_recovered_prefix(prefix) {
return Err(VsdbError::Decode {
detail: format!(
"Mapx metadata prefix {} is outside the allocator-reserved range",
u64::from_le_bytes(prefix)
),
});
}
Ok(Self {
prefix: Prefix::Recovered(prefix),
ns,
})
}
#[inline(always)]
pub(crate) fn encode_prefix_meta(&self) -> Vec<u8> {
let mut meta = Vec::with_capacity(MAPX_META_NS_LEN);
meta.extend_from_slice(MAPX_META_MAGIC);
meta.extend_from_slice(&self.prefix_bytes());
let ns_id = self.ns.id();
if ns_id != DEFAULT_NS_ID {
meta.extend_from_slice(&ns_id.to_le_bytes());
}
meta
}
pub(crate) fn decode_prefix_meta(meta: &[u8]) -> Result<(PreBytes, Option<u64>)> {
if meta.len() != MAPX_META_LEN && meta.len() != MAPX_META_NS_LEN {
return Err(VsdbError::Decode {
detail: format!(
"invalid Mapx metadata length: expected {} or {}, got {}",
MAPX_META_LEN,
MAPX_META_NS_LEN,
meta.len()
),
});
}
if &meta[..MAPX_META_MAGIC.len()] != MAPX_META_MAGIC {
return Err(VsdbError::Decode {
detail: "invalid Mapx metadata magic".to_owned(),
});
}
let mut prefix = PreBytes::default();
prefix.copy_from_slice(&meta[MAPX_META_MAGIC.len()..MAPX_META_LEN]);
let ns_id = (meta.len() == MAPX_META_NS_LEN).then(|| {
let mut b = [0u8; 8];
b.copy_from_slice(&meta[MAPX_META_LEN..]);
u64::from_le_bytes(b)
});
Ok((prefix, ns_id))
}
#[inline(always)]
pub(crate) fn as_prefix_slice(&self) -> &PreBytes {
match &self.prefix {
Prefix::Recovered(bytes) => bytes,
Prefix::Created(cell) => {
self.prefix_bytes();
cell.get().expect("just initialized")
}
}
}
#[inline(always)]
pub fn is_the_same_instance(&self, other_hdr: &Self) -> bool {
self.ns.id() == other_hdr.ns.id()
&& self.prefix_bytes() == other_hdr.prefix_bytes()
}
}
impl Clone for Mapx {
fn clone(&self) -> Self {
const CLONE_CHUNK: usize = 4096;
let mut new_instance = Self::new_in(&self.ns);
let mut it = self.iter();
loop {
let mut batch = new_instance.batch_begin();
let mut wrote = false;
for (k, v) in it.by_ref().take(CLONE_CHUNK) {
batch.insert(&k, &v);
wrote = true;
}
if !wrote {
break;
}
batch.commit().expect("vsdb: clone failed — I/O error");
}
new_instance
}
}
impl PartialEq for Mapx {
fn eq(&self, other: &Mapx) -> bool {
if self.is_the_same_instance(other) {
return true;
}
let mut self_iter = self.iter();
let mut other_iter = other.iter();
loop {
match (self_iter.next(), other_iter.next()) {
(Some((k1, v1)), Some((k2, v2))) => {
if k1 != k2 || v1 != v2 {
return false;
}
}
(None, None) => return true,
_ => return false,
}
}
}
}
impl Eq for Mapx {}
pub(crate) struct SimpleVisitor;
impl<'de> de::Visitor<'de> for SimpleVisitor {
type Value = Vec<u8>;
fn expecting(&self, formatter: &mut fmt::Formatter) -> fmt::Result {
formatter.write_str("bytes")
}
fn visit_str<E>(self, v: &str) -> StdResult<Self::Value, E>
where
E: de::Error,
{
Ok(v.as_bytes().to_vec())
}
fn visit_string<E>(self, v: String) -> StdResult<Self::Value, E>
where
E: de::Error,
{
Ok(v.into_bytes())
}
fn visit_bytes<E>(self, v: &[u8]) -> StdResult<Self::Value, E>
where
E: de::Error,
{
Ok(v.to_vec())
}
fn visit_byte_buf<E>(self, v: Vec<u8>) -> StdResult<Self::Value, E>
where
E: de::Error,
{
Ok(v)
}
fn visit_seq<A>(self, mut seq: A) -> StdResult<Self::Value, A::Error>
where
A: de::SeqAccess<'de>,
{
let mut ret = vec![];
loop {
match seq.next_element() {
Ok(i) => {
if let Some(i) = i {
ret.push(i);
} else {
break;
}
}
Err(e) => {
return Err(de::Error::custom(e));
}
}
}
Ok(ret)
}
}
impl Serialize for Mapx {
fn serialize<S>(&self, serializer: S) -> StdResult<S::Ok, S::Error>
where
S: serde::Serializer,
{
serializer.serialize_bytes(&self.encode_prefix_meta())
}
}
impl<'de> Deserialize<'de> for Mapx {
fn deserialize<D>(deserializer: D) -> StdResult<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
deserializer
.deserialize_byte_buf(SimpleVisitor)
.and_then(|meta| Self::from_prefix_meta(&meta).map_err(de::Error::custom))
}
}
pub struct MapxIter<'a> {
db_iter: DbIter,
_marker: PhantomData<&'a ()>,
}
impl fmt::Debug for MapxIter<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_tuple("MapxIter").finish()
}
}
impl Iterator for MapxIter<'_> {
type Item = (RawKey, RawValue);
fn next(&mut self) -> Option<Self::Item> {
self.db_iter.next()
}
}
impl DoubleEndedIterator for MapxIter<'_> {
fn next_back(&mut self) -> Option<Self::Item> {
self.db_iter.next_back()
}
}
pub struct MapxIterMut<'a> {
db_iter: DbIter,
hdr: &'a mut Mapx,
}
impl fmt::Debug for MapxIterMut<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_tuple("MapxIterMut").field(&self.hdr).finish()
}
}
impl<'a> Iterator for MapxIterMut<'a> {
type Item = (RawKey, ValueIterMut<'a>);
fn next(&mut self) -> Option<Self::Item> {
let (k, v) = self.db_iter.next()?;
let vmut = ValueIterMut {
engine: self.hdr.ns.engine(),
prefix: self.hdr.prefix_bytes(),
key: k.clone(),
value: v,
dirty: false,
_marker: PhantomData,
};
Some((k, vmut))
}
}
impl<'a> DoubleEndedIterator for MapxIterMut<'a> {
fn next_back(&mut self) -> Option<Self::Item> {
let (k, v) = self.db_iter.next_back()?;
let vmut = ValueIterMut {
engine: self.hdr.ns.engine(),
prefix: self.hdr.prefix_bytes(),
key: k.clone(),
value: v,
dirty: false,
_marker: PhantomData,
};
Some((k, vmut))
}
}
pub struct ValueIterMut<'a> {
engine: &'static Engine,
prefix: PreBytes,
key: RawKey,
value: RawValue,
dirty: bool,
_marker: PhantomData<&'a mut ()>,
}
impl fmt::Debug for ValueIterMut<'_> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ValueIterMut")
.field("prefix", &self.prefix)
.field("key", &self.key)
.field("value", &self.value)
.field("dirty", &self.dirty)
.finish()
}
}
impl Drop for ValueIterMut<'_> {
fn drop(&mut self) {
if self.dirty {
self.engine.insert(self.prefix, &self.key, &self.value);
}
}
}
impl Deref for ValueIterMut<'_> {
type Target = RawValue;
fn deref(&self) -> &Self::Target {
&self.value
}
}
impl DerefMut for ValueIterMut<'_> {
fn deref_mut(&mut self) -> &mut Self::Target {
self.dirty = true;
&mut self.value
}
}
#[derive(Debug)]
pub struct ValueMut<'a> {
key: RawKey,
value: RawValue,
dirty: bool,
hdr: &'a mut Mapx,
}
impl Drop for ValueMut<'_> {
fn drop(&mut self) {
if self.dirty {
self.hdr.insert(&self.key[..], &self.value[..]);
}
}
}
impl Deref for ValueMut<'_> {
type Target = RawValue;
fn deref(&self) -> &Self::Target {
&self.value
}
}
impl DerefMut for ValueMut<'_> {
fn deref_mut(&mut self) -> &mut Self::Target {
self.dirty = true;
&mut self.value
}
}