mirror of
https://github.com/bitcoinresearchkit/brk.git
synced 2026-08-12 02:03:09 -07:00
global: fjall+lsm-tree+vecdb part 3
This commit is contained in:
@@ -111,9 +111,7 @@ pub fn main() -> anyhow::Result<()> {
|
||||
|
||||
Mimalloc::collect();
|
||||
|
||||
computer.compute(&indexer, &exit)?;
|
||||
|
||||
indexer.advance_safe_lengths()?;
|
||||
computer.compute(&mut indexer, &exit)?;
|
||||
|
||||
info!("Total time: {:?}", total_start.elapsed());
|
||||
info!("Waiting for new blocks...");
|
||||
|
||||
+1202
-5064
File diff suppressed because it is too large
Load Diff
@@ -53,7 +53,7 @@ pub fn main() -> color_eyre::Result<()> {
|
||||
|
||||
Mimalloc::collect();
|
||||
|
||||
computer.compute(&indexer, &exit)?;
|
||||
computer.compute(&mut indexer, &exit)?;
|
||||
dbg!(i.elapsed());
|
||||
sleep(Duration::from_secs(10));
|
||||
}
|
||||
|
||||
@@ -50,7 +50,7 @@ pub fn main() -> Result<()> {
|
||||
Mimalloc::collect();
|
||||
|
||||
let i = Instant::now();
|
||||
computer.compute(&indexer, &exit)?;
|
||||
computer.compute(&mut indexer, &exit)?;
|
||||
info!("Done in {:?}", i.elapsed());
|
||||
|
||||
// We want to benchmark the drop too
|
||||
|
||||
@@ -66,7 +66,7 @@ pub fn main() -> color_eyre::Result<()> {
|
||||
Mimalloc::collect();
|
||||
|
||||
let i = Instant::now();
|
||||
computer.compute(&indexer, &exit)?;
|
||||
computer.compute(&mut indexer, &exit)?;
|
||||
info!("Done in {:?}", i.elapsed());
|
||||
|
||||
sleep(Duration::from_secs(60));
|
||||
|
||||
@@ -335,7 +335,7 @@ impl Computer {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn compute(&mut self, indexer: &Indexer, exit: &Exit) -> Result<()> {
|
||||
pub fn compute(&mut self, indexer: &mut Indexer, exit: &Exit) -> Result<()> {
|
||||
internal::cache_clear_all();
|
||||
|
||||
let compute_start = Instant::now();
|
||||
@@ -536,6 +536,8 @@ impl Computer {
|
||||
exit,
|
||||
)?;
|
||||
|
||||
indexer.advance_safe_lengths()?;
|
||||
|
||||
info!("Total compute time: {:?}", compute_start.elapsed());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
@@ -41,6 +41,7 @@ fn main() -> color_eyre::Result<()> {
|
||||
loop {
|
||||
let i = Instant::now();
|
||||
indexer.checked_index(&reader, &client, &exit)?;
|
||||
indexer.advance_safe_lengths()?;
|
||||
info!("Done in {:?}", i.elapsed());
|
||||
|
||||
Mimalloc::collect();
|
||||
|
||||
@@ -48,6 +48,7 @@ fn main() -> Result<()> {
|
||||
|
||||
let i = Instant::now();
|
||||
indexer.index(&reader, &client, &exit)?;
|
||||
indexer.advance_safe_lengths()?;
|
||||
info!("Done in {:?}", i.elapsed());
|
||||
|
||||
sleep(Duration::from_secs(60));
|
||||
|
||||
@@ -49,6 +49,7 @@ fn main() -> Result<()> {
|
||||
loop {
|
||||
let i = Instant::now();
|
||||
indexer.index(&reader, &client, &exit)?;
|
||||
indexer.advance_safe_lengths()?;
|
||||
info!("Done in {:?}", i.elapsed());
|
||||
|
||||
Mimalloc::collect();
|
||||
|
||||
@@ -7,7 +7,7 @@ use std::{
|
||||
ops::Deref,
|
||||
sync::{
|
||||
Arc,
|
||||
atomic::{AtomicU64, Ordering},
|
||||
atomic::{AtomicU64, Ordering, fence},
|
||||
},
|
||||
};
|
||||
|
||||
@@ -88,7 +88,15 @@ unsafe impl Sync for ByteView {}
|
||||
|
||||
impl Clone for ByteView {
|
||||
fn clone(&self) -> Self {
|
||||
self.slice(..)
|
||||
if !self.is_inline() {
|
||||
self.get_heap_region()
|
||||
.ref_count
|
||||
.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
|
||||
// SAFETY: Inline views own no external resource. Heap views share their
|
||||
// allocation, whose reference count was incremented above.
|
||||
unsafe { std::ptr::read(self) }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -100,9 +108,10 @@ impl Drop for ByteView {
|
||||
|
||||
let heap_region = self.get_heap_region();
|
||||
|
||||
if heap_region.ref_count.fetch_sub(1, Ordering::AcqRel) != 1 {
|
||||
if heap_region.ref_count.fetch_sub(1, Ordering::Release) != 1 {
|
||||
return;
|
||||
}
|
||||
fence(Ordering::Acquire);
|
||||
|
||||
unsafe {
|
||||
let header_size = std::mem::size_of::<HeapAllocationHeader>();
|
||||
@@ -121,35 +130,27 @@ impl Eq for ByteView {}
|
||||
impl std::cmp::PartialEq for ByteView {
|
||||
fn eq(&self, other: &Self) -> bool {
|
||||
unsafe {
|
||||
let src_ptr = (self as *const Self).cast::<u8>();
|
||||
let other_ptr: *const u8 = (other as *const Self).cast::<u8>();
|
||||
|
||||
let a = *src_ptr.cast::<u64>();
|
||||
let b = *other_ptr.cast::<u64>();
|
||||
let a = std::ptr::from_ref(self).cast::<u64>().read_unaligned();
|
||||
let b = std::ptr::from_ref(other).cast::<u64>().read_unaligned();
|
||||
|
||||
if a != b {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
// NOTE: At this point we know
|
||||
// both strings must have the same prefix and same length
|
||||
//
|
||||
// If we are inlined, the other string must be inlined too,
|
||||
// so checking the short slice is enough
|
||||
if self.is_inline() {
|
||||
self.get_short_slice() == other.get_short_slice()
|
||||
} else {
|
||||
self.get_long_slice() == other.get_long_slice()
|
||||
}
|
||||
// The first word contains the length and cached four-byte prefix.
|
||||
// Compare only the bytes that were not already checked.
|
||||
self.get(PREFIX_SIZE..).unwrap_or_default() == other.get(PREFIX_SIZE..).unwrap_or_default()
|
||||
}
|
||||
}
|
||||
|
||||
impl std::cmp::Ord for ByteView {
|
||||
fn cmp(&self, other: &Self) -> std::cmp::Ordering {
|
||||
self.prefix()
|
||||
.cmp(other.prefix())
|
||||
.then_with(|| self.deref().cmp(&**other))
|
||||
self.prefix().cmp(other.prefix()).then_with(|| {
|
||||
self.get(PREFIX_SIZE..)
|
||||
.unwrap_or_default()
|
||||
.cmp(other.get(PREFIX_SIZE..).unwrap_or_default())
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -239,14 +240,8 @@ impl ByteView {
|
||||
let slice_ptr: &[u8] = &*self;
|
||||
let slice_ptr = slice_ptr.as_ptr();
|
||||
|
||||
// Zero out prefix
|
||||
(*self.trailer.long).prefix[0] = 0;
|
||||
(*self.trailer.long).prefix[1] = 0;
|
||||
(*self.trailer.long).prefix[2] = 0;
|
||||
(*self.trailer.long).prefix[3] = 0;
|
||||
|
||||
let prefix = (*self.trailer.long).prefix.as_mut_ptr();
|
||||
std::ptr::copy_nonoverlapping(slice_ptr, prefix, self.len().min(4));
|
||||
std::ptr::copy_nonoverlapping(slice_ptr, prefix, PREFIX_SIZE);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -594,7 +589,7 @@ impl ByteView {
|
||||
} else {
|
||||
// IMPORTANT: Increase ref count
|
||||
let heap_region = self.get_heap_region();
|
||||
heap_region.ref_count.fetch_add(1, Ordering::Release);
|
||||
heap_region.ref_count.fetch_add(1, Ordering::Relaxed);
|
||||
|
||||
let mut child = Self {
|
||||
// SAFETY: self.data must be defined
|
||||
@@ -625,17 +620,22 @@ impl ByteView {
|
||||
/// Returns `true` if `needle` is a prefix of the slice or equal to the slice.
|
||||
pub fn starts_with<T: AsRef<[u8]>>(&self, needle: T) -> bool {
|
||||
let needle = needle.as_ref();
|
||||
let prefix_len = PREFIX_SIZE.min(needle.len());
|
||||
|
||||
unsafe {
|
||||
let len = PREFIX_SIZE.min(needle.len());
|
||||
let needle_prefix: &[u8] = needle.get_unchecked(..len);
|
||||
let needle_prefix: &[u8] = needle.get_unchecked(..prefix_len);
|
||||
|
||||
if !self.prefix().starts_with(needle_prefix) {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
self.deref().starts_with(needle)
|
||||
if needle.len() <= PREFIX_SIZE {
|
||||
true
|
||||
} else {
|
||||
self.get(PREFIX_SIZE..)
|
||||
.is_some_and(|bytes| bytes.starts_with(&needle[PREFIX_SIZE..]))
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns `true` if the slice is empty.
|
||||
|
||||
@@ -399,14 +399,10 @@ impl Keyspace {
|
||||
pub fn iter_standard(
|
||||
&self,
|
||||
) -> impl DoubleEndedIterator<Item = crate::Result<crate::KvPair>> + Send + 'static {
|
||||
let nonce = self.supervisor.snapshot_tracker.open();
|
||||
let range = ..;
|
||||
self.tree
|
||||
.create_range::<&[u8], _>(&range, nonce.instant, None)
|
||||
.map(move |item| {
|
||||
let _keep_snapshot_alive = &nonce;
|
||||
item.map_err(Into::into)
|
||||
})
|
||||
.create_range_exclusive::<&[u8], _>(&range)
|
||||
.map(|item| item.map_err(Into::into))
|
||||
}
|
||||
|
||||
/// Returns an iterator over a range of items.
|
||||
@@ -440,13 +436,9 @@ impl Keyspace {
|
||||
&self,
|
||||
range: R,
|
||||
) -> impl DoubleEndedIterator<Item = crate::Result<crate::KvPair>> + Send + 'static {
|
||||
let nonce = self.supervisor.snapshot_tracker.open();
|
||||
self.tree
|
||||
.create_range(&range, nonce.instant, None)
|
||||
.map(move |item| {
|
||||
let _keep_snapshot_alive = &nonce;
|
||||
item.map_err(Into::into)
|
||||
})
|
||||
.create_range_exclusive(&range)
|
||||
.map(|item| item.map_err(Into::into))
|
||||
}
|
||||
|
||||
/// Returns an iterator over a prefixed set of items.
|
||||
@@ -480,13 +472,9 @@ impl Keyspace {
|
||||
&self,
|
||||
prefix: K,
|
||||
) -> impl DoubleEndedIterator<Item = crate::Result<crate::KvPair>> + Send + 'static {
|
||||
let nonce = self.supervisor.snapshot_tracker.open();
|
||||
self.tree
|
||||
.create_prefix(prefix, nonce.instant, None)
|
||||
.map(move |item| {
|
||||
let _keep_snapshot_alive = &nonce;
|
||||
item.map_err(Into::into)
|
||||
})
|
||||
.create_prefix_exclusive(prefix)
|
||||
.map(|item| item.map_err(Into::into))
|
||||
}
|
||||
|
||||
/// Approximates the number of items in the keyspace.
|
||||
@@ -646,7 +634,7 @@ impl Keyspace {
|
||||
&self,
|
||||
key: K,
|
||||
) -> crate::Result<Option<lsm_tree::UserValue>> {
|
||||
Ok(self.tree.get(key, SeqNo::MAX)?)
|
||||
Ok(self.tree.get_exclusive(key)?)
|
||||
}
|
||||
|
||||
/// Retrieves the size of an item from the keyspace.
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use criterion::{Criterion, criterion_group, criterion_main};
|
||||
use rand::{Rng, RngCore};
|
||||
use rand::RngExt;
|
||||
|
||||
// Not really worth it anymore on new CPUs...?
|
||||
fn fast_block_index(c: &mut Criterion) {
|
||||
@@ -36,7 +36,7 @@ fn standard_filter_construction(c: &mut Criterion) {
|
||||
|
||||
b.iter(|| {
|
||||
let mut key = [0; 16];
|
||||
rng.fill_bytes(&mut key);
|
||||
rng.fill(&mut key);
|
||||
|
||||
filter.set_with_hash(Builder::get_hash(&key));
|
||||
});
|
||||
@@ -47,7 +47,7 @@ fn standard_filter_construction(c: &mut Criterion) {
|
||||
|
||||
b.iter(|| {
|
||||
let mut key = [0; 16];
|
||||
rng.fill_bytes(&mut key);
|
||||
rng.fill(&mut key);
|
||||
|
||||
filter.set_with_hash(Builder::get_hash(&key));
|
||||
});
|
||||
|
||||
@@ -100,6 +100,7 @@ impl TreeIter {
|
||||
guard: IterState,
|
||||
range: R,
|
||||
seqno: SeqNo,
|
||||
include_memtables: bool,
|
||||
) -> Self {
|
||||
Self::new(guard, |lock| {
|
||||
let lo = match range.start_bound() {
|
||||
@@ -195,33 +196,30 @@ impl TreeIter {
|
||||
}
|
||||
}
|
||||
|
||||
// Sealed memtables
|
||||
for memtable in lock.version.sealed_memtables.iter() {
|
||||
let iter = memtable.range(range.clone());
|
||||
if include_memtables {
|
||||
for memtable in lock.version.sealed_memtables.iter() {
|
||||
let iter = memtable.range(range.clone());
|
||||
|
||||
iters.push(Box::new(
|
||||
iter.filter(move |item| seqno_filter(item.key.seqno, seqno))
|
||||
.map(Ok),
|
||||
));
|
||||
}
|
||||
iters.push(Box::new(
|
||||
iter.filter(move |item| seqno_filter(item.key.seqno, seqno))
|
||||
.map(Ok),
|
||||
));
|
||||
}
|
||||
|
||||
// Active memtable
|
||||
{
|
||||
let iter = lock.version.active_memtable.range(range.clone());
|
||||
|
||||
iters.push(Box::new(
|
||||
iter.filter(move |item| seqno_filter(item.key.seqno, seqno))
|
||||
.map(Ok),
|
||||
));
|
||||
}
|
||||
|
||||
if let Some((mt, seqno)) = &lock.ephemeral {
|
||||
let iter = Box::new(
|
||||
mt.range(range)
|
||||
.filter(move |item| seqno_filter(item.key.seqno, *seqno))
|
||||
.map(Ok),
|
||||
);
|
||||
iters.push(iter);
|
||||
if let Some((mt, seqno)) = &lock.ephemeral {
|
||||
let iter = Box::new(
|
||||
mt.range(range)
|
||||
.filter(move |item| seqno_filter(item.key.seqno, *seqno))
|
||||
.map(Ok),
|
||||
);
|
||||
iters.push(iter);
|
||||
}
|
||||
}
|
||||
|
||||
let merged = Merger::new(iters);
|
||||
|
||||
@@ -164,23 +164,33 @@ impl Table {
|
||||
seqno: SeqNo,
|
||||
key_hash: u64,
|
||||
) -> crate::Result<Option<InternalValue>> {
|
||||
self.get_with(key, seqno, key_hash, Self::point_read)
|
||||
let mut key_hash = Some(key_hash);
|
||||
self.get_with(key, seqno, &mut key_hash, Self::point_read)
|
||||
}
|
||||
|
||||
pub(crate) fn get_value(
|
||||
&self,
|
||||
key: &[u8],
|
||||
seqno: SeqNo,
|
||||
key_hash: u64,
|
||||
key_hash: &mut Option<u64>,
|
||||
) -> crate::Result<Option<PointReadValue>> {
|
||||
self.get_with(key, seqno, key_hash, Self::point_read_value)
|
||||
}
|
||||
|
||||
pub(crate) fn get_lazy(
|
||||
&self,
|
||||
key: &[u8],
|
||||
seqno: SeqNo,
|
||||
key_hash: &mut Option<u64>,
|
||||
) -> crate::Result<Option<InternalValue>> {
|
||||
self.get_with(key, seqno, key_hash, Self::point_read)
|
||||
}
|
||||
|
||||
fn get_with<T>(
|
||||
&self,
|
||||
key: &[u8],
|
||||
seqno: SeqNo,
|
||||
key_hash: u64,
|
||||
key_hash: &mut Option<u64>,
|
||||
point_read: impl FnOnce(&Self, &[u8], SeqNo) -> crate::Result<Option<T>>,
|
||||
) -> crate::Result<Option<T>> {
|
||||
// Translate seqno to "our" seqno
|
||||
@@ -226,6 +236,9 @@ impl Table {
|
||||
};
|
||||
|
||||
if let Some(filter_block) = &filter_block {
|
||||
let key_hash = *key_hash.get_or_insert_with(|| {
|
||||
crate::table::filter::standard_bloom::Builder::get_hash(key)
|
||||
});
|
||||
if !filter_block.maybe_contains_hash(key_hash)? {
|
||||
return Ok(None);
|
||||
}
|
||||
|
||||
+104
-18
@@ -618,19 +618,7 @@ impl AbstractTree for Tree {
|
||||
return Ok(ignore_tombstone_value(entry).map(|entry| entry.value));
|
||||
}
|
||||
|
||||
let key_hash = crate::table::filter::standard_bloom::Builder::get_hash(key);
|
||||
for table in super_version
|
||||
.version
|
||||
.iter_levels()
|
||||
.flat_map(|level| level.iter())
|
||||
.filter_map(|run| run.get_for_key(key))
|
||||
{
|
||||
if let Some(item) = table.get_value(key, seqno, key_hash)? {
|
||||
return Ok((!item.value_type.is_tombstone()).then_some(item.value));
|
||||
}
|
||||
}
|
||||
|
||||
Ok(None)
|
||||
Self::get_value_from_tables(&super_version.version, key, seqno)
|
||||
}
|
||||
|
||||
fn insert<K: Into<UserKey>, V: Into<UserValue>>(
|
||||
@@ -681,7 +669,36 @@ impl Tree {
|
||||
|
||||
let iter_state = { IterState { version, ephemeral } };
|
||||
|
||||
TreeIter::create_range(iter_state, bounds, seqno)
|
||||
TreeIter::create_range(iter_state, bounds, seqno, true)
|
||||
}
|
||||
|
||||
fn create_internal_range_exclusive<'a, K: AsRef<[u8]> + 'a, R: RangeBounds<K> + 'a>(
|
||||
version: SuperVersion,
|
||||
range: &'a R,
|
||||
) -> impl DoubleEndedIterator<Item = crate::Result<InternalValue>> + 'static + use<K, R> {
|
||||
use crate::range::{IterState, TreeIter};
|
||||
use std::ops::Bound::{Excluded, Included, Unbounded};
|
||||
|
||||
let lo: std::ops::Bound<UserKey> = match range.start_bound() {
|
||||
Included(x) => Included(x.as_ref().into()),
|
||||
Excluded(x) => Excluded(x.as_ref().into()),
|
||||
Unbounded => Unbounded,
|
||||
};
|
||||
let hi: std::ops::Bound<UserKey> = match range.end_bound() {
|
||||
Included(x) => Included(x.as_ref().into()),
|
||||
Excluded(x) => Excluded(x.as_ref().into()),
|
||||
Unbounded => Unbounded,
|
||||
};
|
||||
|
||||
TreeIter::create_range(
|
||||
IterState {
|
||||
version,
|
||||
ephemeral: None,
|
||||
},
|
||||
(lo, hi),
|
||||
SeqNo::MAX,
|
||||
false,
|
||||
)
|
||||
}
|
||||
|
||||
pub(crate) fn get_internal_entry_from_version(
|
||||
@@ -709,16 +726,14 @@ impl Tree {
|
||||
key: &[u8],
|
||||
seqno: SeqNo,
|
||||
) -> crate::Result<Option<InternalValue>> {
|
||||
// NOTE: Create key hash for hash sharing
|
||||
// https://fjall-rs.github.io/post/bloom-filter-hash-sharing/
|
||||
let key_hash = crate::table::filter::standard_bloom::Builder::get_hash(key);
|
||||
let mut key_hash = None;
|
||||
|
||||
for table in version
|
||||
.iter_levels()
|
||||
.flat_map(|lvl| lvl.iter())
|
||||
.filter_map(|run| run.get_for_key(key))
|
||||
{
|
||||
if let Some(item) = table.get(key, seqno, key_hash)? {
|
||||
if let Some(item) = table.get_lazy(key, seqno, &mut key_hash)? {
|
||||
return Ok(ignore_tombstone_value(item));
|
||||
}
|
||||
}
|
||||
@@ -726,6 +741,26 @@ impl Tree {
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
fn get_value_from_tables(
|
||||
version: &Version,
|
||||
key: &[u8],
|
||||
seqno: SeqNo,
|
||||
) -> crate::Result<Option<UserValue>> {
|
||||
let mut key_hash = None;
|
||||
|
||||
for table in version
|
||||
.iter_levels()
|
||||
.flat_map(|level| level.iter())
|
||||
.filter_map(|run| run.get_for_key(key))
|
||||
{
|
||||
if let Some(item) = table.get_value(key, seqno, &mut key_hash)? {
|
||||
return Ok((!item.value_type.is_tombstone()).then_some(item.value));
|
||||
}
|
||||
}
|
||||
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
fn get_internal_entry_from_sealed_memtables(
|
||||
super_version: &SuperVersion,
|
||||
key: &[u8],
|
||||
@@ -832,6 +867,43 @@ impl Tree {
|
||||
})
|
||||
}
|
||||
|
||||
/// Reads the latest on-disk value without probing memtables.
|
||||
///
|
||||
/// The caller must ensure all writes use exclusive ingestion.
|
||||
#[doc(hidden)]
|
||||
pub fn get_exclusive<K: AsRef<[u8]>>(&self, key: K) -> crate::Result<Option<UserValue>> {
|
||||
let key = key.as_ref();
|
||||
#[expect(clippy::expect_used, reason = "lock is expected to not be poisoned")]
|
||||
let super_version = self
|
||||
.version_history
|
||||
.read()
|
||||
.expect("lock is poisoned")
|
||||
.latest_version_arc();
|
||||
|
||||
Self::get_value_from_tables(&super_version.version, key, SeqNo::MAX)
|
||||
}
|
||||
|
||||
/// Creates a latest-version range without adding empty memtable readers.
|
||||
///
|
||||
/// The caller must ensure all writes use exclusive ingestion.
|
||||
#[doc(hidden)]
|
||||
pub fn create_range_exclusive<'a, K: AsRef<[u8]> + 'a, R: RangeBounds<K> + 'a>(
|
||||
&self,
|
||||
range: &'a R,
|
||||
) -> impl DoubleEndedIterator<Item = crate::Result<KvPair>> + 'static + use<K, R> {
|
||||
#[expect(clippy::expect_used, reason = "lock is expected to not be poisoned")]
|
||||
let super_version = self
|
||||
.version_history
|
||||
.read()
|
||||
.expect("lock is poisoned")
|
||||
.latest_version();
|
||||
|
||||
Self::create_internal_range_exclusive(super_version, range).map(|item| match item {
|
||||
Ok(kv) => Ok((kv.key.user_key, kv.value)),
|
||||
Err(e) => Err(e),
|
||||
})
|
||||
}
|
||||
|
||||
#[doc(hidden)]
|
||||
pub fn create_prefix<'a, K: AsRef<[u8]> + 'a>(
|
||||
&self,
|
||||
@@ -845,6 +917,20 @@ impl Tree {
|
||||
self.create_range(&range, seqno, ephemeral)
|
||||
}
|
||||
|
||||
/// Creates a latest-version prefix iterator without memtable readers.
|
||||
///
|
||||
/// The caller must ensure all writes use exclusive ingestion.
|
||||
#[doc(hidden)]
|
||||
pub fn create_prefix_exclusive<K: AsRef<[u8]>>(
|
||||
&self,
|
||||
prefix: K,
|
||||
) -> impl DoubleEndedIterator<Item = crate::Result<KvPair>> + 'static + use<K> {
|
||||
use crate::range::prefix_to_range;
|
||||
|
||||
let range = prefix_to_range(prefix.as_ref());
|
||||
self.create_range_exclusive(&range)
|
||||
}
|
||||
|
||||
/// Adds an item to the active memtable.
|
||||
///
|
||||
/// Returns the added item's size and new size of the memtable.
|
||||
|
||||
@@ -158,12 +158,19 @@ impl SuperVersions {
|
||||
pub fn latest_version(&self) -> SuperVersion {
|
||||
#[expect(clippy::expect_used, reason = "SuperVersion is expected to exist")]
|
||||
self.0
|
||||
.iter()
|
||||
.last()
|
||||
.back()
|
||||
.map(|version| version.as_ref().clone())
|
||||
.expect("should always have a SuperVersion")
|
||||
}
|
||||
|
||||
pub(crate) fn latest_version_arc(&self) -> Arc<SuperVersion> {
|
||||
#[expect(clippy::expect_used, reason = "SuperVersion is expected to exist")]
|
||||
self.0
|
||||
.back()
|
||||
.cloned()
|
||||
.expect("should always have a SuperVersion")
|
||||
}
|
||||
|
||||
pub(crate) fn get_version_arc_for_snapshot(&self, seqno: SeqNo) -> Arc<SuperVersion> {
|
||||
if seqno == 0 {
|
||||
#[expect(clippy::expect_used, reason = "SuperVersion is expected to exist")]
|
||||
|
||||
+23
-12
@@ -183,16 +183,20 @@ impl Region {
|
||||
self.write_with(data, Some(at), false)
|
||||
}
|
||||
|
||||
/// Writes (offset, value) pairs directly to the mmap within region bounds.
|
||||
/// Writes ascending (offset, value) pairs directly to the mmap within region bounds.
|
||||
#[inline]
|
||||
pub fn batch_write_each<T, F>(
|
||||
pub fn batch_write_ordered<T, F>(
|
||||
&self,
|
||||
iter: impl Iterator<Item = (usize, T)>,
|
||||
mut iter: impl Iterator<Item = (usize, T)>,
|
||||
value_len: usize,
|
||||
mut write_fn: F,
|
||||
) where
|
||||
F: FnMut(&T, &mut [u8]),
|
||||
{
|
||||
let Some((first_offset, first_value)) = iter.next() else {
|
||||
return;
|
||||
};
|
||||
|
||||
let meta = self.meta();
|
||||
let region_start = meta.start();
|
||||
let region_len = meta.len();
|
||||
@@ -202,10 +206,19 @@ impl Region {
|
||||
let mmap = db.mmap();
|
||||
let ptr = mmap.as_ptr() as *mut u8;
|
||||
|
||||
let mut dirty_start = usize::MAX;
|
||||
let mut dirty_end = 0usize;
|
||||
let first_end = first_offset
|
||||
.checked_add(value_len)
|
||||
.expect("offset + value_len overflow");
|
||||
assert!(first_end <= region_len);
|
||||
let first_abs_offset = region_start + first_offset;
|
||||
let first_slice =
|
||||
unsafe { std::slice::from_raw_parts_mut(ptr.add(first_abs_offset), value_len) };
|
||||
write_fn(&first_value, first_slice);
|
||||
|
||||
let mut previous_offset = first_offset;
|
||||
let mut dirty_end = first_end;
|
||||
for (offset, value) in iter {
|
||||
debug_assert!(offset >= previous_offset, "batch offsets must be ordered");
|
||||
let end_offset = offset
|
||||
.checked_add(value_len)
|
||||
.expect("offset + value_len overflow");
|
||||
@@ -214,15 +227,13 @@ impl Region {
|
||||
let abs_offset = region_start + offset;
|
||||
let slice = unsafe { std::slice::from_raw_parts_mut(ptr.add(abs_offset), value_len) };
|
||||
write_fn(&value, slice);
|
||||
dirty_start = dirty_start.min(offset);
|
||||
dirty_end = dirty_end.max(end_offset);
|
||||
previous_offset = offset;
|
||||
dirty_end = end_offset;
|
||||
}
|
||||
|
||||
if dirty_start < dirty_end {
|
||||
let mut bounds = self.0.dirty_bounds.lock();
|
||||
bounds.0 = bounds.0.min(dirty_start);
|
||||
bounds.1 = bounds.1.max(dirty_end);
|
||||
}
|
||||
let mut bounds = self.0.dirty_bounds.lock();
|
||||
bounds.0 = bounds.0.min(first_offset);
|
||||
bounds.1 = bounds.1.max(dirty_end);
|
||||
}
|
||||
|
||||
pub fn truncate(&self, from: usize) -> Result<()> {
|
||||
|
||||
@@ -59,6 +59,10 @@ pub trait AnyStoredVec: AnyVec {
|
||||
/// Prefixed with `any_` to avoid conflict with `WritableVec::stamped_write_with_changes`.
|
||||
fn any_stamped_write_with_changes(&mut self, stamp: Stamp) -> Result<()>;
|
||||
|
||||
/// Advances the in-memory rollback baseline to the current stored state.
|
||||
#[doc(hidden)]
|
||||
fn any_save_rollback_state(&mut self);
|
||||
|
||||
/// Flushes with the given stamp, optionally saving changes for rollback.
|
||||
#[inline]
|
||||
fn any_stamped_write_maybe_with_changes(
|
||||
@@ -69,7 +73,11 @@ pub trait AnyStoredVec: AnyVec {
|
||||
if with_changes {
|
||||
self.any_stamped_write_with_changes(stamp)
|
||||
} else {
|
||||
self.stamped_write(stamp)
|
||||
self.stamped_write(stamp)?;
|
||||
if self.saved_stamped_changes() > 0 {
|
||||
self.any_save_rollback_state();
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -182,7 +182,11 @@ where
|
||||
if with_changes {
|
||||
self.stamped_write_with_changes(stamp)
|
||||
} else {
|
||||
self.stamped_write(stamp)
|
||||
self.stamped_write(stamp)?;
|
||||
if self.saved_stamped_changes() > 0 {
|
||||
self.save_rollback_state();
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -200,6 +200,10 @@ where
|
||||
<Self as WritableVec<I, T>>::stamped_write_with_changes(self, stamp)
|
||||
}
|
||||
|
||||
fn any_save_rollback_state(&mut self) {
|
||||
<Self as WritableVec<I, T>>::save_rollback_state(self)
|
||||
}
|
||||
|
||||
fn remove(self) -> Result<()> {
|
||||
Self::remove(self)
|
||||
}
|
||||
|
||||
@@ -64,6 +64,10 @@ where
|
||||
self.0.stamped_write_with_changes(stamp)
|
||||
}
|
||||
|
||||
fn any_save_rollback_state(&mut self) {
|
||||
self.0.save_rollback_state()
|
||||
}
|
||||
|
||||
fn remove(self) -> Result<()> {
|
||||
self.0.remove()
|
||||
}
|
||||
|
||||
@@ -179,6 +179,10 @@ macro_rules! impl_vec_wrapper {
|
||||
$crate::WritableVec::stamped_write_with_changes(&mut self.0, stamp)
|
||||
}
|
||||
|
||||
fn any_save_rollback_state(&mut self) {
|
||||
$crate::WritableVec::save_rollback_state(&mut self.0)
|
||||
}
|
||||
|
||||
fn remove(self) -> $crate::Result<()> {
|
||||
self.0.remove()
|
||||
}
|
||||
|
||||
@@ -109,7 +109,7 @@ where
|
||||
}
|
||||
} else {
|
||||
// Normal case: write directly to mmap, no intermediate allocations
|
||||
region.batch_write_each(
|
||||
region.batch_write_ordered(
|
||||
updated
|
||||
.into_iter()
|
||||
.map(|(index, value)| (index * Self::SIZE_OF_T + HEADER_OFFSET, value)),
|
||||
@@ -154,6 +154,10 @@ where
|
||||
<Self as WritableVec<I, T>>::stamped_write_with_changes(self, stamp)
|
||||
}
|
||||
|
||||
fn any_save_rollback_state(&mut self) {
|
||||
<Self as WritableVec<I, T>>::save_rollback_state(self)
|
||||
}
|
||||
|
||||
fn remove(self) -> Result<()> {
|
||||
Self::remove(self)
|
||||
}
|
||||
|
||||
@@ -1471,6 +1471,32 @@ mod raw_rollback {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn run_rollback_after_untracked_checkpoint<V>() -> Result<()>
|
||||
where
|
||||
V: RollbackVec,
|
||||
V::Target: RollbackOps + AnyStoredVec,
|
||||
{
|
||||
let (db, _temp) = setup_db()?;
|
||||
let (mut vec, _) = V::import_with_changes(&db, "test", 10)?;
|
||||
|
||||
for i in 0..100 {
|
||||
vec.push(i);
|
||||
}
|
||||
|
||||
AnyStoredVec::any_stamped_write_maybe_with_changes(vec.deref_mut(), Stamp::new(1), false)?;
|
||||
|
||||
vec.deref_mut().update(65, 999)?;
|
||||
AnyStoredVec::any_stamped_write_maybe_with_changes(vec.deref_mut(), Stamp::new(2), true)?;
|
||||
|
||||
vec.deref_mut().rollback()?;
|
||||
|
||||
assert_eq!(vec.len(), 100);
|
||||
assert_eq!(vec.deref_mut().collect()[65], 65);
|
||||
assert_eq!(RollbackOps::stamp(vec.deref_mut()), Stamp::new(1));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// Test instantiation for each raw vec type
|
||||
// ============================================================================
|
||||
@@ -1549,6 +1575,10 @@ mod raw_rollback {
|
||||
fn rollback_after_rollback_with_delete() -> Result<()> {
|
||||
run_rollback_after_rollback_with_delete::<V>()
|
||||
}
|
||||
#[test]
|
||||
fn rollback_after_untracked_checkpoint() -> Result<()> {
|
||||
run_rollback_after_untracked_checkpoint::<V>()
|
||||
}
|
||||
}
|
||||
|
||||
mod bytes {
|
||||
@@ -1624,11 +1654,87 @@ mod raw_rollback {
|
||||
fn rollback_after_rollback_with_delete() -> Result<()> {
|
||||
run_rollback_after_rollback_with_delete::<V>()
|
||||
}
|
||||
#[test]
|
||||
fn rollback_after_untracked_checkpoint() -> Result<()> {
|
||||
run_rollback_after_untracked_checkpoint::<V>()
|
||||
}
|
||||
}
|
||||
} // end mod raw_rollback
|
||||
|
||||
// ============================================================================
|
||||
// PART 3: Comprehensive Integration Test
|
||||
// PART 3: Checkpoint Rollback Tests (ALL vec types)
|
||||
// ============================================================================
|
||||
|
||||
mod checkpoint_rollback {
|
||||
use super::*;
|
||||
|
||||
fn run<V>() -> Result<()>
|
||||
where
|
||||
V: StoredVec<I = usize, T = u32>,
|
||||
{
|
||||
let (db, _temp) = setup_db()?;
|
||||
let options = ImportOptions::new(&db, "test", Version::TWO).with_saved_stamped_changes(10);
|
||||
let mut vec = V::forced_import_with(options)?;
|
||||
|
||||
for i in 0..100 {
|
||||
vec.push(i);
|
||||
}
|
||||
AnyStoredVec::any_stamped_write_maybe_with_changes(&mut vec, Stamp::new(1), false)?;
|
||||
|
||||
vec.push(100);
|
||||
AnyStoredVec::any_stamped_write_maybe_with_changes(&mut vec, Stamp::new(2), true)?;
|
||||
WritableVec::rollback(&mut vec)?;
|
||||
|
||||
assert_eq!(vec.collect(), (0..100).collect::<Vec<_>>());
|
||||
assert_eq!(AnyStoredVec::stamp(&vec), Stamp::new(1));
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn bytes() -> Result<()> {
|
||||
run::<vecdb::BytesVec<usize, u32>>()
|
||||
}
|
||||
|
||||
#[cfg(feature = "zerocopy")]
|
||||
#[test]
|
||||
fn zerocopy() -> Result<()> {
|
||||
run::<vecdb::ZeroCopyVec<usize, u32>>()
|
||||
}
|
||||
|
||||
#[cfg(feature = "pco")]
|
||||
#[test]
|
||||
fn pco() -> Result<()> {
|
||||
run::<vecdb::PcoVec<usize, u32>>()
|
||||
}
|
||||
|
||||
#[cfg(feature = "lz4")]
|
||||
#[test]
|
||||
fn lz4() -> Result<()> {
|
||||
run::<vecdb::LZ4Vec<usize, u32>>()
|
||||
}
|
||||
|
||||
#[cfg(feature = "zstd")]
|
||||
#[test]
|
||||
fn zstd() -> Result<()> {
|
||||
run::<vecdb::ZstdVec<usize, u32>>()
|
||||
}
|
||||
|
||||
#[cfg(feature = "zerocopy")]
|
||||
#[test]
|
||||
fn eager_zerocopy() -> Result<()> {
|
||||
run::<vecdb::EagerVec<vecdb::ZeroCopyVec<usize, u32>>>()
|
||||
}
|
||||
|
||||
#[cfg(feature = "pco")]
|
||||
#[test]
|
||||
fn eager_pco() -> Result<()> {
|
||||
run::<vecdb::EagerVec<vecdb::PcoVec<usize, u32>>>()
|
||||
}
|
||||
}
|
||||
|
||||
// ============================================================================
|
||||
// PART 4: Comprehensive Integration Test
|
||||
// ============================================================================
|
||||
// Complex rollback + flush + reopen test with file integrity verification.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user