use std::io;
use std::ops::Bound;
use std::sync::Arc;
use super::block::encoded_entry_size;
use super::block::{Block, decode_entry_at};
use super::block_cache::BlockCache;
use super::internal_key::{
INTERNAL_KEY_SUFFIX_LEN, VALUE_TYPE_DELETION, VALUE_TYPE_MERGE, VALUE_TYPE_VALUE,
compare_internal_keys, decode_internal_key, encode_internal_key, user_key_of,
};
use super::lookup_key::LookupKey;
use super::manifest::Version;
use super::memtable::MemTable;
use super::range_tombstone::{RangeTombstone, RangeTombstoneSet};
use super::sstable::{LiveSst, SsTableBlockCursor, SsTableReader};
use crate::DbSlice;
use crate::options::MergeOperator;
use crate::options::PrefixExtractor;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Direction {
Forward,
Reverse,
}
fn above_all_versions(user_key: &[u8]) -> Vec<u8> {
let mut buf = Vec::with_capacity(user_key.len() + INTERNAL_KEY_SUFFIX_LEN);
buf.extend_from_slice(user_key);
buf.extend_from_slice(&[0xff; INTERNAL_KEY_SUFFIX_LEN]);
buf
}
enum LevelIter {
Memtable(MemtableLevelIter),
SsTable(SsTableLevelIter),
LevelConcat(LevelConcatIter),
}
impl LevelIter {
fn seek_to_first(&mut self) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.seek_to_first();
Ok(())
}
Self::SsTable(it) => it.seek_to_first(),
Self::LevelConcat(it) => it.seek_to_first(),
}
}
fn seek_to_last(&mut self) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.seek_to_last();
Ok(())
}
Self::SsTable(it) => it.seek_to_last(),
Self::LevelConcat(it) => it.seek_to_last(),
}
}
fn seek(&mut self, target: &[u8]) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.seek(target);
Ok(())
}
Self::SsTable(it) => it.seek(target),
Self::LevelConcat(it) => it.seek(target),
}
}
fn seek_with_prefix_skip(&mut self, target: &[u8], prefix: &[u8]) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.seek(target);
Ok(())
}
Self::SsTable(it) => {
if !it.reader.may_have_prefix(prefix, &it.cache)? {
it.valid = false;
return Ok(());
}
it.seek(target)
}
Self::LevelConcat(it) => it.seek(target),
}
}
fn seek_for_prev(&mut self, target: &[u8]) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.seek_for_prev(target);
Ok(())
}
Self::SsTable(it) => it.seek_for_prev(target),
Self::LevelConcat(it) => it.seek_for_prev(target),
}
}
fn advance(&mut self) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.advance();
Ok(())
}
Self::SsTable(it) => it.advance(),
Self::LevelConcat(it) => it.advance(),
}
}
fn advance_backward(&mut self) -> io::Result<()> {
match self {
Self::Memtable(it) => {
it.advance_backward();
Ok(())
}
Self::SsTable(it) => it.advance_backward(),
Self::LevelConcat(it) => it.advance_backward(),
}
}
fn key(&self) -> Option<&[u8]> {
match self {
Self::Memtable(it) => it.curr.as_ref().map(|(k, _)| k.as_slice()),
Self::SsTable(it) => {
if it.valid {
Some(&it.cached_key[..])
} else {
None
}
}
Self::LevelConcat(it) => it.key(),
}
}
fn value(&self) -> Option<&[u8]> {
match self {
Self::Memtable(it) => it.curr.as_ref().map(|(_, v)| v.as_slice()),
Self::SsTable(it) => it.value(),
Self::LevelConcat(it) => it.value(),
}
}
fn value_slice(&self) -> Option<DbSlice> {
match self {
Self::Memtable(it) => it.curr.as_ref().map(|(_, v)| v.clone()),
Self::SsTable(it) => it.value_slice(),
Self::LevelConcat(it) => it.value_slice(),
}
}
}
struct MemtableLevelIter {
mt: Arc<MemTable>,
curr: Option<(DbSlice, DbSlice)>,
}
impl MemtableLevelIter {
fn new(mt: Arc<MemTable>) -> Self {
Self { mt, curr: None }
}
fn seek_to_first(&mut self) {
self.curr = self.mt.first_slice_from(Bound::Unbounded);
}
fn seek_to_last(&mut self) {
self.curr = self.mt.last_slice_before(Bound::Unbounded);
}
fn seek(&mut self, target: &[u8]) {
self.curr = self.mt.first_slice_from(Bound::Included(target));
}
fn seek_for_prev(&mut self, target: &[u8]) {
self.curr = self.mt.last_slice_before(Bound::Included(target));
}
fn advance(&mut self) {
let Some((key, _)) = self.curr.take() else {
return;
};
self.curr = self.mt.first_slice_from(Bound::Excluded(key.as_slice()));
}
fn advance_backward(&mut self) {
let Some((key, _)) = self.curr.take() else {
return;
};
self.curr = self.mt.last_slice_before(Bound::Excluded(key.as_slice()));
}
}
struct SsTableLevelIter {
reader: Arc<SsTableReader>,
cache: Arc<BlockCache>,
block: Option<Arc<Block>>,
block_cursor: Option<SsTableBlockCursor>,
entry_pos: usize,
next_entry_pos: usize,
cached_key: Vec<u8>,
cached_value_offset: usize,
cached_value_len: usize,
valid: bool,
current_entry_index: Option<usize>,
entry_offsets: Option<Vec<usize>>,
span: Vec<u8>,
span_start: u64,
span_cap: u64,
span_blocks: u64,
sequential: bool,
}
const SCAN_SPAN_BLOCKS: u64 = 16;
const SCAN_SPAN_MAX: u64 = 1 << 20;
const SCAN_SPAN_BLOCKS_INITIAL: u64 = 2;
impl SsTableLevelIter {
fn new(reader: Arc<SsTableReader>, cache: Arc<BlockCache>, span_cap: u64) -> Self {
Self {
reader,
cache,
block: None,
block_cursor: None,
entry_pos: 0,
next_entry_pos: 0,
cached_key: Vec::new(),
cached_value_offset: 0,
cached_value_len: 0,
valid: false,
current_entry_index: None,
entry_offsets: None,
span: Vec::new(),
span_start: 0,
span_cap,
span_blocks: SCAN_SPAN_BLOCKS_INITIAL,
sequential: false,
}
}
fn block_for(&mut self, handle: super::block::BlockHandle) -> io::Result<Arc<Block>> {
if !self.sequential || handle.size > self.span_cap {
return self.reader.read_block(handle, &self.cache);
}
if let Some(block) = self.cache.get(self.reader.file_id, handle.offset) {
return Ok(block);
}
let Some(handle_end) = handle.offset.checked_add(handle.size) else {
return self.reader.read_block(handle, &self.cache);
};
let span_end = self.span_start.saturating_add(self.span.len() as u64);
let held =
!self.span.is_empty() && handle.offset >= self.span_start && handle_end <= span_end;
if !held {
let want = handle
.size
.saturating_mul(self.span_blocks)
.clamp(handle.size, self.span_cap.max(handle.size));
self.span_blocks = (self.span_blocks * 2).min(SCAN_SPAN_BLOCKS);
self.span = self.reader.read_span(handle.offset, want)?;
self.span_start = handle.offset;
if (self.span.len() as u64) < handle.size {
self.span.clear();
return self.reader.read_block(handle, &self.cache);
}
}
let frame = usize::try_from(handle.offset - self.span_start)
.ok()
.zip(usize::try_from(handle.size).ok())
.and_then(|(lo, len)| lo.checked_add(len).map(|hi| (lo, hi)))
.and_then(|(lo, hi)| self.span.get(lo..hi));
let Some(frame) = frame else {
return self.reader.read_block(handle, &self.cache);
};
self.reader.decode_block_frame(handle, frame, &self.cache)
}
fn reset_readahead(&mut self) {
self.span = Vec::new();
self.span_start = 0;
self.span_blocks = SCAN_SPAN_BLOCKS_INITIAL;
self.sequential = false;
}
fn value(&self) -> Option<&[u8]> {
if !self.valid {
return None;
}
let data = self.block.as_ref()?.entry_data();
let end = self
.cached_value_offset
.checked_add(self.cached_value_len)?;
data.get(self.cached_value_offset..end)
}
fn value_slice(&self) -> Option<DbSlice> {
if !self.valid {
return None;
}
let block = Arc::clone(self.block.as_ref()?);
DbSlice::from_block(block, self.cached_value_offset, self.cached_value_len)
}
fn load_block(&mut self, cursor: SsTableBlockCursor) -> io::Result<()> {
let handle = self.reader.cursor_handle(&cursor, &self.cache)?;
self.block = Some(self.block_for(handle)?);
self.block_cursor = Some(cursor);
self.entry_pos = 0;
self.next_entry_pos = 0;
self.entry_offsets = None;
self.current_entry_index = None;
self.cached_key.clear();
Ok(())
}
fn decode_current(&mut self) {
let data = self.block.as_ref().unwrap().entry_data();
if self.entry_pos >= data.len() {
self.valid = false;
return;
}
let (consumed, val_off, val_len) =
decode_entry_at(data, self.entry_pos, &mut self.cached_key);
self.next_entry_pos = self.entry_pos + consumed;
self.cached_value_offset = val_off;
self.cached_value_len = val_len;
self.current_entry_index = self
.entry_offsets
.as_ref()
.and_then(|offsets| offsets.binary_search(&self.entry_pos).ok());
self.valid = true;
}
fn seek_to_first(&mut self) -> io::Result<()> {
self.reset_readahead();
let Some(cursor) = self.reader.first_block_cursor(&self.cache)? else {
self.valid = false;
return Ok(());
};
self.load_block(cursor)?;
self.decode_current();
Ok(())
}
fn seek(&mut self, target: &[u8]) -> io::Result<()> {
self.reset_readahead();
let cursor = match self.reader.seek_block_cursor(target, &self.cache)? {
Some(cursor) => cursor,
None => {
self.valid = false;
return Ok(());
}
};
self.load_block(cursor)?;
let data = self.block.as_ref().unwrap().entry_data();
self.entry_pos = 0;
self.cached_key.clear();
while self.entry_pos < data.len() {
let (consumed, val_off, val_len) =
decode_entry_at(data, self.entry_pos, &mut self.cached_key);
self.next_entry_pos = self.entry_pos + consumed;
self.cached_value_offset = val_off;
self.cached_value_len = val_len;
if compare_internal_keys(&self.cached_key, target).is_ge() {
self.valid = true;
return Ok(());
}
self.entry_pos = self.next_entry_pos;
}
let next = self
.reader
.next_block_cursor(self.block_cursor.as_ref().unwrap(), &self.cache)?;
let Some(next) = next else {
self.valid = false;
return Ok(());
};
self.sequential = true;
self.load_block(next)?;
self.decode_current();
Ok(())
}
fn seek_for_prev(&mut self, target: &[u8]) -> io::Result<()> {
self.reset_readahead();
let cursor = match self.reader.seek_block_cursor(target, &self.cache)? {
Some(cursor) => cursor,
None => match self.reader.last_block_cursor(&self.cache)? {
Some(cursor) => cursor,
None => {
self.valid = false;
return Ok(());
}
},
};
self.load_block(cursor)?;
self.build_entry_offsets();
let offsets = self.entry_offsets.as_ref().unwrap();
let data = self.block.as_ref().unwrap().entry_data();
let mut best: Option<usize> = None;
let mut temp_key = Vec::new();
for (i, &off) in offsets.iter().enumerate() {
let (_consumed, _vo, _vl) = decode_entry_at(data, off, &mut temp_key);
if compare_internal_keys(&temp_key, target).is_le() {
best = Some(i);
} else {
break;
}
}
match best {
Some(idx) => {
self.replay_key_to_index(idx);
self.valid = true;
}
None => {
let prev = self
.reader
.prev_block_cursor(self.block_cursor.as_ref().unwrap(), &self.cache)?;
let Some(prev) = prev else {
self.valid = false;
return Ok(());
};
self.load_block(prev)?;
self.build_entry_offsets();
let Some(last) = self.last_entry_index() else {
self.valid = false;
return Ok(());
};
self.replay_key_to_index(last);
self.valid = true;
}
}
Ok(())
}
fn advance(&mut self) -> io::Result<()> {
if !self.valid {
return Ok(());
}
self.entry_pos = self.next_entry_pos;
let data = self.block.as_ref().unwrap().entry_data();
if self.entry_pos >= data.len() {
let next = self
.reader
.next_block_cursor(self.block_cursor.as_ref().unwrap(), &self.cache)?;
self.sequential = true;
let Some(next) = next else {
self.valid = false;
return Ok(());
};
self.load_block(next)?;
}
self.decode_current();
Ok(())
}
fn seek_to_last(&mut self) -> io::Result<()> {
self.reset_readahead();
let Some(cursor) = self.reader.last_block_cursor(&self.cache)? else {
self.valid = false;
return Ok(());
};
self.load_block(cursor)?;
self.build_entry_offsets();
let Some(last) = self.last_entry_index() else {
self.valid = false;
return Ok(());
};
self.replay_key_to_index(last);
self.valid = true;
Ok(())
}
fn advance_backward(&mut self) -> io::Result<()> {
if !self.valid {
return Ok(());
}
if self.entry_offsets.is_none() {
self.build_entry_offsets();
}
let offsets = self.entry_offsets.as_ref().unwrap();
let cur_idx = self.current_entry_index.unwrap_or_else(|| {
offsets
.binary_search(&self.entry_pos)
.unwrap_or_else(|idx| idx.saturating_sub(1))
});
if cur_idx == 0 {
let prev = self
.reader
.prev_block_cursor(self.block_cursor.as_ref().unwrap(), &self.cache)?;
let Some(prev) = prev else {
self.valid = false;
return Ok(());
};
self.load_block(prev)?;
self.build_entry_offsets();
let Some(last) = self.last_entry_index() else {
self.valid = false;
return Ok(());
};
self.replay_key_to_index(last);
self.valid = true;
return Ok(());
}
self.replay_key_to_index(cur_idx - 1);
self.valid = true;
Ok(())
}
fn build_entry_offsets(&mut self) {
let data = self.block.as_ref().unwrap().entry_data();
let mut offsets = Vec::new();
let mut pos = 0;
while pos < data.len() {
offsets.push(pos);
pos += encoded_entry_size(data, pos);
}
self.entry_offsets = Some(offsets);
}
fn last_entry_index(&self) -> Option<usize> {
self.entry_offsets.as_ref()?.len().checked_sub(1)
}
fn replay_key_to_index(&mut self, target_idx: usize) {
let block = self.block.as_ref().unwrap();
let data = block.entry_data();
let Some(offsets) = self.entry_offsets.as_ref() else {
return;
};
let Some(&target_off) = offsets.get(target_idx) else {
return;
};
let (mut lo, mut hi) = (0usize, block.restart_count());
while lo < hi {
let mid = lo + (hi - lo) / 2;
if block.restart_offset(mid) <= target_off {
lo = mid + 1;
} else {
hi = mid;
}
}
let start_offset = if lo > 0 {
block.restart_offset(lo - 1)
} else {
0
};
let start_entry_idx = offsets.binary_search(&start_offset).unwrap_or(0);
self.cached_key.clear();
let mut pos = start_offset;
for idx in start_entry_idx..=target_idx {
let (consumed, val_off, val_len) = decode_entry_at(data, pos, &mut self.cached_key);
if idx == target_idx {
self.entry_pos = pos;
self.next_entry_pos = pos + consumed;
self.cached_value_offset = val_off;
self.cached_value_len = val_len;
self.current_entry_index = Some(target_idx);
}
pos += consumed;
}
}
}
struct LevelConcatIter {
files: Vec<Arc<LiveSst>>,
cache: Arc<BlockCache>,
file_idx: usize,
current: Option<SsTableLevelIter>,
span_cap: u64,
}
impl LevelConcatIter {
fn new(files: Vec<Arc<LiveSst>>, cache: Arc<BlockCache>, span_cap: u64) -> Self {
Self {
files,
cache,
file_idx: 0,
current: None,
span_cap,
}
}
fn open_current(&mut self) -> io::Result<()> {
if self.file_idx < self.files.len() {
self.current = Some(SsTableLevelIter::new(
Arc::clone(&self.files[self.file_idx].reader),
Arc::clone(&self.cache),
self.span_cap,
));
} else {
self.current = None;
}
Ok(())
}
fn seek_to_first(&mut self) -> io::Result<()> {
if self.files.is_empty() {
self.current = None;
return Ok(());
}
self.file_idx = 0;
loop {
self.open_current()?;
self.current.as_mut().unwrap().seek_to_first()?;
if self.current.as_ref().is_some_and(|it| it.valid) {
return Ok(());
}
self.file_idx += 1;
if self.file_idx >= self.files.len() {
self.current = None;
return Ok(());
}
}
}
fn seek_to_last(&mut self) -> io::Result<()> {
if self.files.is_empty() {
self.current = None;
return Ok(());
}
self.file_idx = self.files.len() - 1;
loop {
self.open_current()?;
self.current.as_mut().unwrap().seek_to_last()?;
if self.current.as_ref().is_some_and(|it| it.valid) {
return Ok(());
}
if self.file_idx == 0 {
self.current = None;
return Ok(());
}
self.file_idx -= 1;
}
}
fn seek(&mut self, target: &[u8]) -> io::Result<()> {
if self.files.is_empty() {
self.current = None;
return Ok(());
}
let uk = user_key_of(target);
let idx = self
.files
.partition_point(|f| f.meta.largest_key.as_slice() < uk);
if idx >= self.files.len() {
self.current = None;
return Ok(());
}
self.file_idx = idx;
loop {
self.open_current()?;
self.current.as_mut().unwrap().seek(target)?;
if self.current.as_ref().is_some_and(|it| it.valid) {
return Ok(());
}
self.file_idx += 1;
if self.file_idx >= self.files.len() {
self.current = None;
return Ok(());
}
}
}
fn seek_for_prev(&mut self, target: &[u8]) -> io::Result<()> {
if self.files.is_empty() {
self.current = None;
return Ok(());
}
let uk = user_key_of(target);
let idx = self
.files
.partition_point(|f| f.meta.largest_key.as_slice() < uk);
let idx = if idx >= self.files.len() {
self.files.len() - 1
} else {
idx
};
self.file_idx = idx;
loop {
self.open_current()?;
self.current.as_mut().unwrap().seek_for_prev(target)?;
if self.current.as_ref().is_some_and(|it| it.valid) {
return Ok(());
}
if self.file_idx == 0 {
self.current = None;
return Ok(());
}
self.file_idx -= 1;
}
}
fn advance(&mut self) -> io::Result<()> {
if let Some(ref mut it) = self.current {
it.advance()?;
if it.valid {
return Ok(());
}
}
self.file_idx += 1;
if self.file_idx >= self.files.len() {
self.current = None;
return Ok(());
}
loop {
self.open_current()?;
self.current.as_mut().unwrap().seek_to_first()?;
if self.current.as_ref().is_some_and(|it| it.valid) {
return Ok(());
}
self.file_idx += 1;
if self.file_idx >= self.files.len() {
self.current = None;
return Ok(());
}
}
}
fn advance_backward(&mut self) -> io::Result<()> {
if let Some(ref mut it) = self.current {
it.advance_backward()?;
if it.valid {
return Ok(());
}
}
if self.file_idx == 0 {
self.current = None;
return Ok(());
}
self.file_idx -= 1;
loop {
self.open_current()?;
self.current.as_mut().unwrap().seek_to_last()?;
if self.current.as_ref().is_some_and(|it| it.valid) {
return Ok(());
}
if self.file_idx == 0 {
self.current = None;
return Ok(());
}
self.file_idx -= 1;
}
}
fn key(&self) -> Option<&[u8]> {
self.current.as_ref().and_then(|it| {
if it.valid {
Some(it.cached_key.as_slice())
} else {
None
}
})
}
fn value(&self) -> Option<&[u8]> {
self.current.as_ref().and_then(|it| it.value())
}
fn value_slice(&self) -> Option<DbSlice> {
self.current.as_ref().and_then(|it| it.value_slice())
}
}
struct MergingIter {
levels: Vec<LevelIter>,
current_idx: Option<usize>,
}
impl MergingIter {
fn new(levels: Vec<LevelIter>) -> Self {
Self {
levels,
current_idx: None,
}
}
fn seek_to_first(&mut self) -> io::Result<()> {
for lvl in &mut self.levels {
lvl.seek_to_first()?;
}
self.pick_smallest();
Ok(())
}
fn seek_to_last(&mut self) -> io::Result<()> {
for lvl in &mut self.levels {
lvl.seek_to_last()?;
}
self.pick_largest();
Ok(())
}
fn seek(&mut self, target: &[u8]) -> io::Result<()> {
for lvl in &mut self.levels {
lvl.seek(target)?;
}
self.pick_smallest();
Ok(())
}
fn seek_with_prefix_skip(&mut self, target: &[u8], prefix: &[u8]) -> io::Result<()> {
for lvl in &mut self.levels {
lvl.seek_with_prefix_skip(target, prefix)?;
}
self.pick_smallest();
Ok(())
}
fn seek_for_prev(&mut self, target: &[u8]) -> io::Result<()> {
for lvl in &mut self.levels {
lvl.seek_for_prev(target)?;
}
self.pick_largest();
Ok(())
}
fn advance(&mut self) -> io::Result<()> {
if let Some(idx) = self.current_idx {
self.levels[idx].advance()?;
}
self.pick_smallest();
Ok(())
}
fn advance_backward(&mut self) -> io::Result<()> {
if let Some(idx) = self.current_idx {
self.levels[idx].advance_backward()?;
}
self.pick_largest();
Ok(())
}
fn key(&self) -> Option<&[u8]> {
self.current_idx.and_then(|i| self.levels[i].key())
}
fn value(&self) -> Option<&[u8]> {
self.current_idx.and_then(|i| self.levels[i].value())
}
fn value_slice(&self) -> Option<DbSlice> {
self.current_idx.and_then(|i| self.levels[i].value_slice())
}
fn pick_smallest(&mut self) {
let mut best: Option<usize> = None;
for (i, lvl) in self.levels.iter().enumerate() {
let Some(k) = lvl.key() else { continue };
match best {
None => best = Some(i),
Some(bi) => {
let bk = self.levels[bi].key().unwrap();
if compare_internal_keys(k, bk).is_lt() {
best = Some(i);
}
}
}
}
self.current_idx = best;
}
fn pick_largest(&mut self) {
let mut best: Option<usize> = None;
for (i, lvl) in self.levels.iter().enumerate() {
let Some(k) = lvl.key() else { continue };
match best {
None => best = Some(i),
Some(bi) => {
let bk = self.levels[bi].key().unwrap();
if compare_internal_keys(k, bk).is_gt() {
best = Some(i);
}
}
}
}
self.current_idx = best;
}
}
pub(crate) struct RegolithIterator {
inner: MergingIter,
snapshot_seq: u64,
direction: Direction,
valid_entry: bool,
positioned: bool,
curr_user_key: Vec<u8>,
reverse_curr: Option<(Vec<u8>, DbSlice)>,
merge_result: Option<(Vec<u8>, Vec<u8>)>,
range_tombstones: RangeTombstoneSet,
_version: Arc<Version>,
error: Option<io::Error>,
terminal_error: bool,
upper_bound: Option<Vec<u8>>,
prefix_extractor: Option<Arc<dyn PrefixExtractor>>,
merge_operator: Option<Arc<dyn MergeOperator>>,
pending_consume: bool,
last_forward_user_key: Option<Vec<u8>>,
}
fn prefix_upper_bound(prefix: &[u8]) -> Option<Vec<u8>> {
let mut out = prefix.to_vec();
while let Some(last) = out.last_mut() {
if *last != 0xff {
*last += 1;
return Some(out);
}
out.pop();
}
None
}
impl RegolithIterator {
pub(crate) fn new(
active: Arc<MemTable>,
frozen: Vec<Arc<MemTable>>,
version: Arc<Version>,
cache: Arc<BlockCache>,
snapshot_seq: u64,
prefix_extractor: Option<Arc<dyn PrefixExtractor>>,
merge_operator: Option<Arc<dyn MergeOperator>>,
) -> Self {
let mut levels: Vec<LevelIter> = Vec::new();
let mut range_tombstones: Vec<RangeTombstone> = Vec::new();
let sstable_sources =
version.levels[0].len() + version.levels[1..].iter().filter(|l| !l.is_empty()).count();
let span_cap = SCAN_SPAN_MAX / (sstable_sources.max(1) as u64);
range_tombstones.extend(active.clone_range_tombstones());
levels.push(LevelIter::Memtable(MemtableLevelIter::new(active)));
for mt in frozen.iter().rev() {
range_tombstones.extend(mt.clone_range_tombstones());
levels.push(LevelIter::Memtable(MemtableLevelIter::new(Arc::clone(mt))));
}
for file in version.levels[0].iter().rev() {
range_tombstones.extend(file.reader.range_tombstones().iter().cloned());
levels.push(LevelIter::SsTable(SsTableLevelIter::new(
Arc::clone(&file.reader),
Arc::clone(&cache),
span_cap,
)));
}
for level in 1..version.levels.len() {
if version.levels[level].is_empty() {
continue;
}
for file in &version.levels[level] {
range_tombstones.extend(file.reader.range_tombstones().iter().cloned());
}
let mut sorted: Vec<Arc<LiveSst>> = version.levels[level]
.iter()
.filter(|f| f.meta.num_entries > 0)
.map(Arc::clone)
.collect();
if sorted.is_empty() {
continue;
}
sorted.sort_by(|a, b| a.meta.smallest_key.cmp(&b.meta.smallest_key));
levels.push(LevelIter::LevelConcat(LevelConcatIter::new(
sorted,
Arc::clone(&cache),
span_cap,
)));
}
Self {
inner: MergingIter::new(levels),
snapshot_seq,
direction: Direction::Forward,
valid_entry: false,
positioned: false,
curr_user_key: Vec::new(),
reverse_curr: None,
merge_result: None,
range_tombstones: RangeTombstoneSet::from_vec(range_tombstones),
_version: version,
error: None,
terminal_error: false,
upper_bound: None,
prefix_extractor,
merge_operator,
pending_consume: false,
last_forward_user_key: None,
}
}
fn covering_rt_seq(&self, user_key: &[u8]) -> u64 {
if self.range_tombstones.is_empty() {
return 0;
}
self.range_tombstones
.max_covering_seq(user_key, self.snapshot_seq)
}
pub(crate) fn positioned(&self) -> bool {
self.positioned
}
pub(crate) fn seek_to_first(&mut self) {
self.positioned = true;
if self.terminal_error {
return;
}
self.error = None;
self.valid_entry = false;
self.last_forward_user_key = None;
self.merge_result = None;
self.reverse_curr = None;
self.pending_consume = false;
self.upper_bound = None;
self.direction = Direction::Forward;
if let Err(e) = self.inner.seek_to_first() {
self.error = Some(e);
return;
}
self.materialize_next_visible();
}
pub(crate) fn seek_to_last(&mut self) {
self.positioned = true;
if self.terminal_error {
return;
}
self.error = None;
self.valid_entry = false;
self.last_forward_user_key = None;
self.merge_result = None;
self.reverse_curr = None;
self.pending_consume = false;
self.upper_bound = None;
self.direction = Direction::Reverse;
if let Err(e) = self.inner.seek_to_last() {
self.error = Some(e);
return;
}
self.materialize_prev_visible();
}
pub(crate) fn seek(&mut self, target: &[u8]) {
self.positioned = true;
if self.terminal_error {
return;
}
self.error = None;
self.valid_entry = false;
self.last_forward_user_key = None;
self.merge_result = None;
self.reverse_curr = None;
self.pending_consume = false;
self.upper_bound = None;
self.direction = Direction::Forward;
let search_key = LookupKey::from_prefixed(target, u64::MAX);
if let Err(e) = self.inner.seek(search_key.internal()) {
self.error = Some(e);
return;
}
self.materialize_next_visible();
}
pub(crate) fn seek_for_prev(&mut self, target: &[u8]) {
self.positioned = true;
if self.terminal_error {
return;
}
self.error = None;
self.valid_entry = false;
self.last_forward_user_key = None;
self.merge_result = None;
self.reverse_curr = None;
self.pending_consume = false;
self.upper_bound = None;
self.direction = Direction::Reverse;
let probe = above_all_versions(target);
if let Err(e) = self.inner.seek_for_prev(&probe) {
self.error = Some(e);
return;
}
self.materialize_prev_visible();
}
pub(crate) fn seek_to_last_before(&mut self, exclusive_upper: &[u8]) {
self.positioned = true;
if self.terminal_error {
return;
}
self.error = None;
self.valid_entry = false;
self.last_forward_user_key = None;
self.merge_result = None;
self.reverse_curr = None;
self.pending_consume = false;
self.upper_bound = None;
self.direction = Direction::Reverse;
let probe = encode_internal_key(exclusive_upper, u64::MAX, VALUE_TYPE_DELETION);
if let Err(e) = self.inner.seek_for_prev(&probe) {
self.error = Some(e);
return;
}
while let Some(k) = self.inner.key() {
if user_key_of(k) < exclusive_upper {
break;
}
if let Err(e) = self.inner.advance_backward() {
self.error = Some(e);
return;
}
}
self.materialize_prev_visible();
}
pub(crate) fn seek_prefix(&mut self, prefix: &[u8]) {
self.positioned = true;
if self.terminal_error {
return;
}
self.error = None;
self.valid_entry = false;
self.last_forward_user_key = None;
self.merge_result = None;
self.reverse_curr = None;
self.pending_consume = false;
self.direction = Direction::Forward;
self.upper_bound = prefix_upper_bound(prefix);
let bloom_probe = self
.prefix_extractor
.as_ref()
.and_then(|ex| ex.extract_query(prefix).map(|p| p.to_vec()));
let search_key = LookupKey::from_prefixed(prefix, u64::MAX);
let res = if let Some(probe) = bloom_probe.as_deref() {
self.inner
.seek_with_prefix_skip(search_key.internal(), probe)
} else {
self.inner.seek(search_key.internal())
};
if let Err(e) = res {
self.error = Some(e);
return;
}
self.materialize_next_visible();
}
pub(crate) fn next(&mut self) {
if !self.valid() {
return;
}
if self.direction == Direction::Reverse {
self.flip_to_forward();
}
self.merge_result = None;
if self.pending_consume {
self.consume_curr_user_key_forward();
self.pending_consume = false;
}
self.materialize_next_visible();
}
pub(crate) fn prev(&mut self) {
if !self.valid() {
return;
}
if self.pending_consume {
self.consume_curr_user_key_forward();
self.pending_consume = false;
}
self.merge_result = None;
if self.direction == Direction::Forward {
self.flip_to_reverse();
}
self.reverse_curr = None;
self.materialize_prev_visible();
}
pub(crate) fn valid(&self) -> bool {
if self.error.is_some() {
return false;
}
match self.direction {
Direction::Forward => self.valid_entry || self.merge_result.is_some(),
Direction::Reverse => self.reverse_curr.is_some(),
}
}
pub(crate) fn key(&self) -> Option<&[u8]> {
match self.direction {
Direction::Forward => {
if let Some((k, _)) = &self.merge_result {
return Some(k.as_slice());
}
if !self.valid_entry {
return None;
}
self.inner.key().map(user_key_of)
}
Direction::Reverse => self.reverse_curr.as_ref().map(|(k, _)| k.as_slice()),
}
}
pub(crate) fn value(&self) -> Option<&[u8]> {
match self.direction {
Direction::Forward => {
if let Some((_, v)) = &self.merge_result {
return Some(v.as_slice());
}
if !self.valid_entry {
return None;
}
self.inner.value()
}
Direction::Reverse => self.reverse_curr.as_ref().map(|(_, v)| v.as_slice()),
}
}
pub(crate) fn value_slice(&self) -> Option<DbSlice> {
match self.direction {
Direction::Forward => {
if let Some((_, v)) = &self.merge_result {
return Some(DbSlice::from_vec(v.clone()));
}
if !self.valid_entry {
return None;
}
self.inner.value_slice()
}
Direction::Reverse => self.reverse_curr.as_ref().map(|(_, v)| v.clone()),
}
}
pub(crate) fn status(&self) -> io::Result<()> {
match &self.error {
Some(e) => Err(io::Error::new(e.kind(), e.to_string())),
None => Ok(()),
}
}
pub(crate) fn set_error(&mut self, err: io::Error) {
self.error = Some(err);
self.terminal_error = true;
self.valid_entry = false;
self.merge_result = None;
self.reverse_curr = None;
}
fn consume_curr_user_key_forward(&mut self) {
if let Err(e) = self.inner.advance() {
self.error = Some(e);
return;
}
loop {
let matches = {
let Some(ik) = self.inner.key() else {
return;
};
let (uk, _, _) = decode_internal_key(ik);
uk == self.curr_user_key.as_slice()
};
if !matches {
return;
}
if let Err(e) = self.inner.advance() {
self.error = Some(e);
return;
}
}
}
fn materialize_next_visible(&mut self) {
self.valid_entry = false;
loop {
let Some(ik) = self.inner.key() else {
return;
};
let (uk, seq, vt) = decode_internal_key(ik);
let went_backwards = self
.last_forward_user_key
.as_deref()
.is_some_and(|last| uk <= last);
if went_backwards {
self.error = Some(io::Error::new(
io::ErrorKind::InvalidData,
"SSTable iteration went backwards: file index or data block is corrupt",
));
return;
}
if let Some(ub) = self.upper_bound.as_deref()
&& uk >= ub
{
return;
}
if seq > self.snapshot_seq {
if let Err(e) = self.inner.advance() {
self.error = Some(e);
return;
}
continue;
}
let rt_seq = if self.range_tombstones.is_empty() {
0
} else {
self.covering_rt_seq(uk)
};
if rt_seq > seq {
self.curr_user_key.clear();
self.curr_user_key.extend_from_slice(uk);
self.consume_curr_user_key_forward();
continue;
}
match vt {
VALUE_TYPE_DELETION => {
self.curr_user_key.clear();
self.curr_user_key.extend_from_slice(uk);
self.consume_curr_user_key_forward();
continue;
}
VALUE_TYPE_MERGE => {
let uk_owned = uk.to_vec();
match self.collapse_merge_chain_forward(&uk_owned, rt_seq) {
Ok(Some(v)) => {
self.curr_user_key.clear();
self.curr_user_key.extend_from_slice(&uk_owned);
self.last_forward_user_key = Some(uk_owned.clone());
self.merge_result = Some((uk_owned, v));
self.pending_consume = false;
return;
}
Ok(None) => continue,
Err(e) => {
self.error = Some(e);
return;
}
}
}
_ => {
self.curr_user_key.clear();
self.curr_user_key.extend_from_slice(uk);
match self.last_forward_user_key.as_mut() {
Some(last) => {
last.clear();
last.extend_from_slice(uk);
}
None => self.last_forward_user_key = Some(uk.to_vec()),
}
self.valid_entry = true;
self.pending_consume = true;
return;
}
}
}
}
fn collapse_merge_chain_forward(
&mut self,
user_key: &[u8],
rt_seq: u64,
) -> io::Result<Option<Vec<u8>>> {
let merge_op = match self.merge_operator.clone() {
Some(op) => op,
None => {
self.consume_user_key_forward(user_key);
return Ok(None);
}
};
let mut operands_newest_first: Vec<Vec<u8>> = Vec::new();
let mut base: Option<Vec<u8>> = None;
let mut had_terminator = false;
#[allow(clippy::while_let_loop)]
loop {
let Some(ik) = self.inner.key() else { break };
let (uk, seq, vt) = decode_internal_key(ik);
if uk != user_key {
break;
}
if seq > self.snapshot_seq {
self.inner.advance()?;
continue;
}
if rt_seq > 0 && seq <= rt_seq {
base = None;
had_terminator = true;
self.consume_user_key_forward(user_key);
break;
}
let value = self.inner.value().map(|s| s.to_vec()).unwrap_or_default();
match vt {
VALUE_TYPE_MERGE => {
operands_newest_first.push(value);
self.inner.advance()?;
}
VALUE_TYPE_VALUE => {
base = Some(value);
had_terminator = true;
self.consume_user_key_forward(user_key);
break;
}
VALUE_TYPE_DELETION => {
base = None;
had_terminator = true;
self.consume_user_key_forward(user_key);
break;
}
_ => {
self.inner.advance()?;
}
}
}
let _ = had_terminator;
if operands_newest_first.is_empty() {
return Ok(base);
}
let operand_refs: Vec<&[u8]> = operands_newest_first
.iter()
.rev()
.map(|v| v.as_slice())
.collect();
match merge_op.full_merge(user_key, base.as_deref(), &operand_refs) {
Some(v) => Ok(Some(v)),
None => Err(io::Error::new(
io::ErrorKind::InvalidData,
format!("merge operator {} failed", merge_op.name()),
)),
}
}
fn materialize_prev_visible(&mut self) {
loop {
let Some(ik) = self.inner.key() else {
self.reverse_curr = None;
return;
};
let (uk, _, _) = decode_internal_key(ik);
let group = uk.to_vec();
let rt_seq = self.covering_rt_seq(&group);
let mut collected: Vec<(u64, u8, DbSlice)> = Vec::new();
while let Some(ik2) = self.inner.key() {
let (uk2, seq, vt) = decode_internal_key(ik2);
if uk2 != group.as_slice() {
break;
}
if seq <= self.snapshot_seq && (rt_seq == 0 || seq > rt_seq) {
let v = self.inner.value_slice().unwrap_or_else(DbSlice::empty);
collected.push((seq, vt, v));
}
if let Err(e) = self.inner.advance_backward() {
self.error = Some(e);
self.reverse_curr = None;
return;
}
}
if collected.is_empty() {
continue;
}
let mut terminator_idx: Option<usize> = None;
for (i, (_, vt, _)) in collected.iter().enumerate().rev() {
if *vt != VALUE_TYPE_MERGE {
terminator_idx = Some(i);
break;
}
}
let (base, operand_range_start) = match terminator_idx {
Some(i) => match collected[i].1 {
VALUE_TYPE_VALUE => (Some(collected[i].2.clone()), i + 1),
VALUE_TYPE_DELETION => (None, i + 1),
_ => (None, i + 1),
},
None => (None, 0),
};
let operand_slice = &collected[operand_range_start..];
if operand_slice.is_empty() {
match terminator_idx {
Some(i) if collected[i].1 == VALUE_TYPE_VALUE => {
self.reverse_curr = Some((group, base.unwrap()));
return;
}
_ => continue,
}
}
let Some(merge_op) = self.merge_operator.clone() else {
continue;
};
let operand_refs: Vec<&[u8]> = operand_slice.iter().map(|e| e.2.as_slice()).collect();
let base_bytes = base.as_ref().map(DbSlice::as_slice);
match merge_op.full_merge(&group, base_bytes, &operand_refs) {
Some(v) => {
self.reverse_curr = Some((group, DbSlice::from_vec(v)));
return;
}
None => {
self.error = Some(io::Error::new(
io::ErrorKind::InvalidData,
format!("merge operator {} failed", merge_op.name()),
));
self.reverse_curr = None;
return;
}
}
}
}
fn consume_user_key_forward(&mut self, user_key: &[u8]) {
loop {
let Some(ik) = self.inner.key() else { return };
let (uk, _, _) = decode_internal_key(ik);
if uk != user_key {
return;
}
if let Err(e) = self.inner.advance() {
self.error = Some(e);
return;
}
}
}
fn flip_to_forward(&mut self) {
let Some((uk, _)) = &self.reverse_curr else {
return;
};
let probe = above_all_versions(uk);
if let Err(e) = self.inner.seek(&probe) {
self.error = Some(e);
}
self.reverse_curr = None;
self.direction = Direction::Forward;
self.last_forward_user_key = None;
}
fn flip_to_reverse(&mut self) {
let uk: Vec<u8> = if let Some((k, _)) = &self.merge_result {
k.clone()
} else {
self.curr_user_key.clone()
};
if uk.is_empty() {
return;
}
let probe = LookupKey::from_prefixed(&uk, u64::MAX);
if let Err(e) = self.inner.seek_for_prev(probe.internal()) {
self.error = Some(e);
}
self.valid_entry = false;
self.merge_result = None;
self.direction = Direction::Reverse;
}
}