mod slot_type;
pub use slot_type::SlotType;
use crate::{
KeyEnDeOrdered,
common::{
error::Result,
staged::{StagedRows, prefix_successor},
},
};
use serde::{Deserialize, Serialize};
use std::{
borrow::Cow, collections::BTreeMap, fmt, marker::PhantomData, ops::Bound,
result::Result as StdResult,
};
use vsdb_core::basic::mapx_raw::MapxRaw;
pub(crate) type EntryCnt = u64;
type SkipNum = EntryCnt;
type TakeNum = EntryCnt;
type Distance = i128;
type PageSize = u16;
type PageIndex = u32;
const TAG_ENTRY: u8 = 0x00;
const TAG_LEVEL: u8 = 0x01;
const TAG_TOTAL: u8 = 0x02;
const LAYOUT_VERSION: u8 = 2;
pub struct SlotDex<S, K>
where
S: SlotType,
K: Clone + Ord + KeyEnDeOrdered,
{
store: MapxRaw,
tier_capacity: S,
swap_order: bool,
total_cache: EntryCnt,
levels: Vec<Level<S>>,
_p: PhantomData<K>,
}
struct Level<S> {
floor_base: S,
buckets: BTreeMap<S, EntryCnt>,
}
impl<S: SlotType> Level<S> {
fn new(level: u8, tier_capacity: &S) -> Self {
Self {
floor_base: floor_base_of(level, tier_capacity),
buckets: BTreeMap::new(),
}
}
}
fn floor_base_of<S: SlotType>(level: u8, tier_capacity: &S) -> S {
tier_capacity
.checked_pow(level as u32)
.filter(|v| *v != S::MIN)
.unwrap_or(S::MAX)
}
fn entry_key<S: SlotType, K: KeyEnDeOrdered>(slot: &S, k: &K) -> Vec<u8> {
let s = slot.to_bytes();
let kb = k.to_bytes();
let mut v = Vec::with_capacity(1 + s.len() + kb.len());
v.push(TAG_ENTRY);
v.extend_from_slice(&s);
v.extend_from_slice(&kb);
v
}
fn level_key<S: SlotType>(level: u8, floor: &S) -> Vec<u8> {
let f = floor.to_bytes();
let mut v = Vec::with_capacity(2 + f.len());
v.push(TAG_LEVEL);
v.push(level);
v.extend_from_slice(&f);
v
}
fn level_prefix(level: u8) -> Vec<u8> {
vec![TAG_LEVEL, level]
}
const TOTAL_KEY: [u8; 1] = [TAG_TOTAL];
fn encode_cnt(v: EntryCnt) -> [u8; 8] {
v.to_le_bytes()
}
fn decode_cnt(raw: &[u8]) -> EntryCnt {
let mut b = [0u8; 8];
b.copy_from_slice(&raw[..8]);
EntryCnt::from_le_bytes(b)
}
fn bound_to_raw(b: Bound<Vec<u8>>) -> Bound<Cow<'static, [u8]>> {
match b {
Bound::Included(v) => Bound::Included(Cow::Owned(v)),
Bound::Excluded(v) => Bound::Excluded(Cow::Owned(v)),
Bound::Unbounded => Bound::Unbounded,
}
}
fn entry_lower_bound<S: SlotType>(b: &Bound<S>) -> Option<Bound<Vec<u8>>> {
match b {
Bound::Unbounded => {
let mut v = vec![TAG_ENTRY];
v.extend_from_slice(&S::MIN.to_bytes());
Some(Bound::Included(v))
}
Bound::Included(s) => {
let mut v = vec![TAG_ENTRY];
v.extend_from_slice(&s.to_bytes());
Some(Bound::Included(v))
}
Bound::Excluded(s) => {
let succ = prefix_successor(&s.to_bytes())?;
let mut v = vec![TAG_ENTRY];
v.extend_from_slice(&succ);
Some(Bound::Included(v))
}
}
}
fn entry_upper_bound<S: SlotType>(b: &Bound<S>) -> Bound<Vec<u8>> {
match b {
Bound::Unbounded => Bound::Excluded(vec![TAG_ENTRY + 1]),
Bound::Included(s) => {
let mut v = vec![TAG_ENTRY];
v.extend_from_slice(&s.to_bytes());
Bound::Excluded(prefix_successor(&v).expect("prefix starts with 0x00"))
}
Bound::Excluded(s) => {
let mut v = vec![TAG_ENTRY];
v.extend_from_slice(&s.to_bytes());
Bound::Excluded(v)
}
}
}
fn decode_entry_row<S: SlotType, K: KeyEnDeOrdered>(raw_key: &[u8]) -> (S, K) {
let slot_len = S::MIN.to_bytes().len();
let slot = S::from_bytes(raw_key[1..1 + slot_len].to_vec())
.expect("SlotDex: corrupt entry-row slot bytes");
let k = K::from_bytes(raw_key[1 + slot_len..].to_vec())
.expect("SlotDex: corrupt entry-row key bytes");
(slot, k)
}
fn decode_level_row<S: SlotType>(raw_key: &[u8], raw_val: &[u8]) -> (S, EntryCnt) {
let floor = S::from_bytes(raw_key[2..].to_vec())
.expect("SlotDex: corrupt level-row floor bytes");
(floor, decode_cnt(raw_val))
}
impl<S, K> Serialize for SlotDex<S, K>
where
S: SlotType,
K: Clone + Ord + KeyEnDeOrdered,
{
fn serialize<Ser>(&self, serializer: Ser) -> StdResult<Ser::Ok, Ser::Error>
where
Ser: serde::Serializer,
{
crate::common::serialize_typed_handle_meta::<Self, Ser>(
&(
LAYOUT_VERSION,
&self.store,
&self.tier_capacity,
&self.swap_order,
),
serializer,
)
}
}
impl<'de, S, K> Deserialize<'de> for SlotDex<S, K>
where
S: SlotType,
K: Clone + Ord + KeyEnDeOrdered,
{
fn deserialize<D>(deserializer: D) -> StdResult<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
let (version, store, tier_capacity, swap_order) =
crate::common::deserialize_typed_handle_meta::<
Self,
(u8, MapxRaw, S, bool),
D,
>(deserializer)?;
if version != LAYOUT_VERSION {
return Err(serde::de::Error::custom(format!(
"SlotDex: unsupported layout version {version} (expected {LAYOUT_VERSION})"
)));
}
Ok(Self::hydrate(store, tier_capacity, swap_order))
}
}
impl<S, K> fmt::Debug for SlotDex<S, K>
where
S: SlotType,
K: Clone + Ord + KeyEnDeOrdered,
{
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("SlotDex")
.field("total", &self.total_cache)
.field("levels", &self.levels.len())
.field("tier_capacity", &self.tier_capacity)
.field("swap_order", &self.swap_order)
.finish()
}
}
impl<S, K> SlotDex<S, K>
where
S: SlotType,
K: Clone + Ord + KeyEnDeOrdered,
{
pub fn new(tier_capacity: S, swap_order: bool) -> Self {
assert!(
tier_capacity.as_i128() >= 2,
"SlotDex: tier_capacity must be >= 2"
);
Self {
store: MapxRaw::new(),
tier_capacity,
swap_order,
total_cache: 0,
levels: vec![],
_p: PhantomData,
}
}
fn hydrate(store: MapxRaw, tier_capacity: S, swap_order: bool) -> Self {
let total_cache = store.get(TOTAL_KEY).map(|v| decode_cnt(&v)).unwrap_or(0);
let mut levels = vec![];
for level in 1..=u8::MAX {
let prefix = level_prefix(level);
let lo = Bound::Included(Cow::Owned(prefix.clone()));
let hi = bound_to_raw(
prefix_successor(&prefix)
.map(Bound::Excluded)
.expect("prefix starts with 0x01"),
);
let mut cache = Level::new(level, &tier_capacity);
let mut empty = true;
for (rk, rv) in store.range((lo, hi)) {
let (floor, cnt) = decode_level_row::<S>(&rk, &rv);
cache.buckets.insert(floor, cnt);
empty = false;
}
if empty {
break;
}
levels.push(cache);
}
Self {
store,
tier_capacity,
swap_order,
total_cache,
levels,
_p: PhantomData,
}
}
#[inline(always)]
pub fn instance_id(&self) -> u64 {
self.store.instance_id()
}
pub fn save_meta(&self) -> Result<u64> {
let id = self.instance_id();
crate::common::save_instance_meta(id, self)?;
Ok(id)
}
pub fn from_meta(instance_id: u64) -> Result<Self> {
crate::common::load_instance_meta(instance_id)
}
pub fn insert(&mut self, slot: S, k: K) -> Result<()> {
let slot = self.to_storage_slot(slot);
let ekey = entry_key(&slot, &k);
if self.store.contains_key(&ekey) {
return Ok(());
}
let mut staged = StagedRows::new();
let grown = self.stage_level_growth(&mut staged);
staged.put(ekey, vec![]);
let slot_cnt = self.slot_entry_cnt(&slot);
staged.put(level_key(0, &slot), encode_cnt(slot_cnt + 1).to_vec());
let mut bumps: Vec<(usize, S, EntryCnt)> = Vec::with_capacity(self.levels.len());
for (i, lv) in self.levels.iter().enumerate() {
let floor = slot.floor_align(&lv.floor_base);
let base = lv.buckets.get(&floor).copied().unwrap_or(0);
staged.put(
level_key(i as u8 + 1, &floor),
encode_cnt(base + 1).to_vec(),
);
bumps.push((i, floor, base + 1));
}
if let Some(grown) = grown.as_ref() {
let floor = slot.floor_align(&grown.floor_base);
let base = grown.buckets.get(&floor).copied().unwrap_or(0);
staged.put(
level_key(self.levels.len() as u8 + 1, &floor),
encode_cnt(base + 1).to_vec(),
);
}
staged.put(
TOTAL_KEY.to_vec(),
encode_cnt(self.total_cache + 1).to_vec(),
);
staged.commit(&mut self.store)?;
if let Some(mut grown) = grown {
let floor = slot.floor_align(&grown.floor_base);
*grown.buckets.entry(floor).or_insert(0) += 1;
self.levels.push(grown);
}
for (i, floor, v) in bumps {
self.levels[i].buckets.insert(floor, v);
}
self.total_cache += 1;
Ok(())
}
pub fn insert_batch<I>(&mut self, items: I) -> Result<()>
where
I: IntoIterator<Item = (S, K)>,
{
let mut groups: BTreeMap<S, Vec<K>> = BTreeMap::new();
for (slot, k) in items {
groups
.entry(self.to_storage_slot(slot))
.or_default()
.push(k);
}
if groups.is_empty() {
return Ok(());
}
let mut staged = StagedRows::new();
let mut grown: Vec<Level<S>> = vec![];
let mut bumps: BTreeMap<(usize, S), EntryCnt> = BTreeMap::new();
let mut total_added: EntryCnt = 0;
for (slot, ks) in groups {
if let Some(lv) = self.stage_level_growth_over(&mut staged, &grown, &bumps) {
grown.push(lv);
}
let mut added: EntryCnt = 0;
for k in ks {
let ekey = entry_key(&slot, &k);
if staged.get_over(&self.store, &ekey).is_some() {
continue;
}
staged.put(ekey, vec![]);
added += 1;
}
if 0 == added {
continue;
}
let slot_cnt = self.slot_entry_cnt(&slot);
staged.put(level_key(0, &slot), encode_cnt(slot_cnt + added).to_vec());
for (i, lv) in self
.levels
.iter()
.map(|l| (&l.floor_base, &l.buckets))
.chain(grown.iter().map(|l| (&l.floor_base, &l.buckets)))
.enumerate()
{
let (floor_base, buckets) = lv;
let floor = slot.floor_align(floor_base);
let v = bumps
.get(&(i, floor.clone()))
.copied()
.unwrap_or_else(|| buckets.get(&floor).copied().unwrap_or(0))
+ added;
staged.put(level_key(i as u8 + 1, &floor), encode_cnt(v).to_vec());
bumps.insert((i, floor), v);
}
total_added += added;
}
if 0 == total_added {
return Ok(());
}
staged.put(
TOTAL_KEY.to_vec(),
encode_cnt(self.total_cache + total_added).to_vec(),
);
staged.commit(&mut self.store)?;
self.levels.extend(grown);
for ((i, floor), v) in bumps {
self.levels[i].buckets.insert(floor, v);
}
self.total_cache += total_added;
Ok(())
}
pub fn remove(&mut self, slot: S, k: &K) {
let slot = self.to_storage_slot(slot);
let ekey = entry_key(&slot, k);
if !self.store.contains_key(&ekey) {
return;
}
let mut staged = StagedRows::new();
staged.del(ekey);
let mut kept = self.levels.len();
while kept > 0 && self.levels[kept - 1].buckets.len() < 2 {
for floor in self.levels[kept - 1].buckets.keys() {
staged.del(level_key(kept as u8, floor));
}
kept -= 1;
}
let slot_cnt = self.slot_entry_cnt(&slot);
if slot_cnt <= 1 {
staged.del(level_key(0, &slot));
} else {
staged.put(level_key(0, &slot), encode_cnt(slot_cnt - 1).to_vec());
}
let mut decs: Vec<(usize, S, Option<EntryCnt>)> = Vec::with_capacity(kept);
for (i, lv) in self.levels.iter().take(kept).enumerate() {
let floor = slot.floor_align(&lv.floor_base);
let cnt = match lv.buckets.get(&floor).copied() {
Some(n) => n,
None => continue,
};
if cnt <= 1 {
staged.del(level_key(i as u8 + 1, &floor));
decs.push((i, floor, None));
} else {
staged.put(level_key(i as u8 + 1, &floor), encode_cnt(cnt - 1).to_vec());
decs.push((i, floor, Some(cnt - 1)));
}
}
staged.put(
TOTAL_KEY.to_vec(),
encode_cnt(self.total_cache.saturating_sub(1)).to_vec(),
);
staged
.commit(&mut self.store)
.expect("vsdb: SlotDex remove batch commit failed");
self.levels.truncate(kept);
for (i, floor, v) in decs {
match v {
Some(v) => {
self.levels[i].buckets.insert(floor, v);
}
None => {
self.levels[i].buckets.remove(&floor);
}
}
}
self.total_cache = self.total_cache.saturating_sub(1);
}
pub fn clear(&mut self) {
self.store.clear();
self.levels.clear();
self.total_cache = 0;
}
fn stage_level_growth(&self, staged: &mut StagedRows) -> Option<Level<S>> {
self.stage_level_growth_over(staged, &[], &BTreeMap::new())
}
fn stage_level_growth_over(
&self,
staged: &mut StagedRows,
grown: &[Level<S>],
bumps: &BTreeMap<(usize, S), EntryCnt>,
) -> Option<Level<S>> {
let n = self.levels.len() + grown.len();
let new_level_no = n as u8 + 1;
let mut newtop = Level::new(new_level_no, &self.tier_capacity);
if let Some(top) = grown.last().or_else(|| self.levels.last()) {
let top_idx = n - 1;
let mut view = top.buckets.clone();
for ((i, floor), v) in bumps {
if *i == top_idx {
view.insert(floor.clone(), *v);
}
}
if view.len() as i128 <= self.tier_capacity.as_i128() {
return None;
}
for (slot, cnt) in view {
let floor = slot.floor_align(&newtop.floor_base);
*newtop.buckets.entry(floor).or_insert(0) += cnt;
}
} else {
for (rk, rv) in self.level0_range(Bound::Unbounded, Bound::Unbounded) {
let (slot, cnt) = decode_level_row::<S>(&rk, &rv);
let floor = slot.floor_align(&newtop.floor_base);
*newtop.buckets.entry(floor).or_insert(0) += cnt;
}
}
for (floor, cnt) in &newtop.buckets {
staged.put(level_key(new_level_no, floor), encode_cnt(*cnt).to_vec());
}
Some(newtop)
}
fn level0_range(
&self,
lo: Bound<S>,
hi: Bound<S>,
) -> impl DoubleEndedIterator<
Item = (vsdb_core::common::RawKey, vsdb_core::common::RawValue),
> + '_ {
let lo = match &lo {
Bound::Unbounded => Bound::Included(level_key(0, &S::MIN)),
Bound::Included(s) => Bound::Included(level_key(0, s)),
Bound::Excluded(s) => Bound::Excluded(level_key(0, s)),
};
let hi = match &hi {
Bound::Unbounded => Bound::Excluded(level_prefix(1)),
Bound::Included(s) => Bound::Included(level_key(0, s)),
Bound::Excluded(s) => Bound::Excluded(level_key(0, s)),
};
self.store.range((bound_to_raw(lo), bound_to_raw(hi)))
}
fn entry_range(
&self,
lo: Bound<S>,
hi: Bound<S>,
) -> Option<
impl DoubleEndedIterator<
Item = (vsdb_core::common::RawKey, vsdb_core::common::RawValue),
> + '_,
> {
let lo = entry_lower_bound(&lo)?;
let hi = entry_upper_bound(&hi);
Some(self.store.range((bound_to_raw(lo), bound_to_raw(hi))))
}
fn slot_entry_cnt(&self, slot: &S) -> EntryCnt {
self.store
.get(level_key(0, slot))
.map(|v| decode_cnt(&v))
.unwrap_or(0)
}
fn level0_walk_desc(
&self,
lo: Bound<S>,
hi: Bound<S>,
mut visit: impl FnMut(S, EntryCnt) -> bool,
) {
let splits: Vec<S> = match self.levels.first() {
Some(l1) => {
let above_lo = match &lo {
Bound::Included(s) | Bound::Excluded(s) => {
Bound::Excluded(s.clone())
}
Bound::Unbounded => Bound::Unbounded,
};
let nonempty = match (&above_lo, &hi) {
(Bound::Unbounded, _) | (_, Bound::Unbounded) => true,
(
Bound::Included(a) | Bound::Excluded(a),
Bound::Included(b) | Bound::Excluded(b),
) => a < b,
};
if nonempty {
let stride =
(32 / self.tier_capacity.as_i128()).clamp(1, 32) as usize;
l1.buckets
.range((above_lo, hi.clone()))
.rev()
.map(|(f, _)| f.clone())
.skip(stride - 1)
.step_by(stride)
.collect()
} else {
vec![]
}
}
None => vec![],
};
let mut buf: Vec<(S, EntryCnt)> = vec![];
let mut upper = hi;
for split in splits {
buf.clear();
for (rk, rv) in self.level0_range(Bound::Included(split.clone()), upper) {
buf.push(decode_level_row::<S>(&rk, &rv));
}
for (slot, cnt) in buf.drain(..).rev() {
if !visit(slot, cnt) {
return;
}
}
upper = Bound::Excluded(split);
}
buf.clear();
for (rk, rv) in self.level0_range(lo, upper) {
buf.push(decode_level_row::<S>(&rk, &rv));
}
for (slot, cnt) in buf.drain(..).rev() {
if !visit(slot, cnt) {
return;
}
}
}
pub fn get_entries_by_page(
&self,
page_size: PageSize,
page_index: PageIndex, reverse_order: bool,
) -> Vec<K> {
self.get_entries_by_page_slot(None, None, page_size, page_index, reverse_order)
}
pub fn get_entries_by_page_slot(
&self,
slot_left_bound: Option<S>, slot_right_bound: Option<S>, page_size: PageSize,
page_index: PageIndex, reverse_order: bool,
) -> Vec<K> {
let (slot_min, slot_max, storage_is_reversed) =
self.transform_range(slot_left_bound, slot_right_bound);
if slot_max < slot_min {
return vec![];
}
if 0 == page_size || 0 == self.total() {
return vec![];
}
self.get_entries(
slot_min,
slot_max,
page_size,
page_index,
reverse_order ^ storage_is_reversed,
)
}
fn distance_to_the_rightmost_slot(&self, slot: &S) -> Distance {
if *slot == S::MAX {
return 0;
}
self.total() as Distance
- self.distance_to_the_leftmost_slot(slot)
- self.slot_entry_cnt(slot) as Distance
}
fn distance_to_the_leftmost_slot(&self, slot: &S) -> Distance {
if *slot == S::MIN {
return 0;
}
let mut left_bound = S::MIN;
let mut ret = 0;
for lv in self.levels.iter().rev() {
let right_bound = slot.floor_align(&lv.floor_base);
ret += lv
.buckets
.range(left_bound.clone()..right_bound.clone())
.map(|(_, cnt)| *cnt as Distance)
.sum::<Distance>();
left_bound = right_bound;
}
ret += self
.level0_range(Bound::Included(left_bound), Bound::Excluded(slot.clone()))
.map(|(_, v)| decode_cnt(&v) as Distance)
.sum::<Distance>();
ret
}
fn offsets_from_the_leftmost_slot(
&self,
slot_start: &S, page_size: PageSize,
page_index: PageIndex,
) -> (SkipNum, TakeNum) {
let skip_n = self.distance_to_the_leftmost_slot(slot_start)
+ (page_size as Distance) * (page_index as Distance);
(skip_n as SkipNum, page_size as TakeNum)
}
fn locate_page_start(&self, global_skip_n: EntryCnt) -> (Bound<S>, SkipNum) {
let mut slot_start = Bound::Included(S::MIN);
let mut remaining: u64 = global_skip_n;
for lv in self.levels.iter().rev() {
let mut hdr = lv
.buckets
.range((slot_start.clone(), Bound::Unbounded))
.peekable();
while let Some(entry_cnt) = hdr.next().map(|(_, cnt)| *cnt) {
if entry_cnt > remaining {
break;
} else {
slot_start = hdr
.peek()
.map(|(s, _)| Bound::Included((*s).clone()))
.unwrap_or(Bound::Excluded(S::MAX));
remaining -= entry_cnt;
}
}
}
let mut hdr = self
.level0_range(slot_start.clone(), Bound::Unbounded)
.peekable();
while let Some(entry_cnt) = hdr.next().map(|(_, v)| decode_cnt(&v)) {
if entry_cnt > remaining {
break;
} else {
slot_start = hdr
.peek()
.map(|(rk, _)| {
Bound::Included(
S::from_bytes(rk[2..].to_vec())
.expect("SlotDex: corrupt level-row floor bytes"),
)
})
.unwrap_or(Bound::Excluded(S::MAX));
remaining -= entry_cnt;
}
}
(slot_start, remaining)
}
fn locate_page_rstart(&self, global_skip_n: EntryCnt) -> (Bound<S>, SkipNum) {
let mut slot_end: Bound<S> = Bound::Unbounded;
let mut remaining: u64 = global_skip_n;
for lv in self.levels.iter().rev() {
for (floor, entry_cnt) in
lv.buckets.range((Bound::Unbounded, slot_end.clone())).rev()
{
if *entry_cnt > remaining {
break;
}
slot_end = Bound::Excluded(floor.clone());
remaining -= *entry_cnt;
}
}
let mut slot_end_cell = slot_end;
let mut remaining_cell = remaining;
self.level0_walk_desc(
Bound::Unbounded,
slot_end_cell.clone(),
|slot, entry_cnt| {
if entry_cnt > remaining_cell {
return false;
}
slot_end_cell = Bound::Excluded(slot);
remaining_cell -= entry_cnt;
true
},
);
(slot_end_cell, remaining_cell)
}
fn get_entries(
&self,
slot_start: S, slot_end: S, page_size: PageSize,
page_index: PageIndex,
reverse: bool,
) -> Vec<K> {
if slot_end < slot_start {
return vec![];
}
if reverse {
return self
.get_entries_reverse(slot_start, slot_end, page_size, page_index);
}
let (global_skip_n, take_n) =
self.offsets_from_the_leftmost_slot(&slot_start, page_size, page_index);
let (slot_start_actual, local_skip_n) = self.locate_page_start(global_skip_n);
let iter = match self.entry_range(slot_start_actual, Bound::Included(slot_end)) {
Some(iter) => iter,
None => return vec![],
};
iter.skip(local_skip_n as usize)
.take(take_n as usize)
.map(|(rk, _)| decode_entry_row::<S, K>(&rk).1)
.collect()
}
fn get_entries_reverse(
&self,
slot_start: S, slot_end: S, page_size: PageSize,
page_index: PageIndex,
) -> Vec<K> {
let global_skip_n = self.distance_to_the_rightmost_slot(&slot_end)
+ (page_size as Distance) * (page_index as Distance);
let global_skip_n = u64::try_from(global_skip_n).unwrap_or(u64::MAX);
let (slot_end_actual, local_skip_n) = self.locate_page_rstart(global_skip_n);
let mut to_skip = local_skip_n as usize;
let mut remaining = page_size as usize;
let mut plan: Vec<(S, usize, usize)> = vec![];
self.level0_walk_desc(
Bound::Included(slot_start.clone()),
slot_end_actual,
|slot, n| {
let n = n as usize;
if to_skip >= n {
to_skip -= n;
return true;
}
let take = (n - to_skip).min(remaining);
plan.push((slot, to_skip, take));
remaining -= take;
to_skip = 0;
remaining > 0
},
);
let Some((last, _, _)) = plan.last() else {
return vec![];
};
let lo = last.clone();
let hi = plan[0].0.clone();
let iter = match self.entry_range(Bound::Included(lo), Bound::Included(hi)) {
Some(iter) => iter,
None => return vec![],
};
let mut segments: BTreeMap<S, Vec<K>> =
plan.iter().map(|(s, ..)| (s.clone(), vec![])).collect();
let quota: BTreeMap<S, (usize, usize)> = plan
.iter()
.map(|(s, skip, take)| (s.clone(), (*skip, *take)))
.collect();
let mut cur: Option<(S, usize, usize, usize)> = None; for (rk, _) in iter {
let (slot, k) = decode_entry_row::<S, K>(&rk);
match &mut cur {
Some((s, skip, take, seen)) if *s == slot => {
*seen += 1;
if *seen > *skip && segments[s].len() < *take {
segments.get_mut(s).expect("planned").push(k);
}
}
_ => {
let Some(&(skip, take)) = quota.get(&slot) else {
cur = Some((slot, usize::MAX, 0, 0));
continue;
};
if skip == 0 && take > 0 {
segments.get_mut(&slot).expect("planned").push(k);
}
cur = Some((slot, skip, take, 1));
}
}
}
let mut ret = Vec::with_capacity(page_size as usize);
for (slot, ..) in &plan {
ret.extend(segments.remove(slot).expect("planned"));
}
ret
}
pub fn entry_cnt_within_two_slots(&self, slot_start: S, slot_end: S) -> EntryCnt {
let (slot_min, slot_max, _) =
self.transform_range(Some(slot_start), Some(slot_end));
if slot_min > slot_max {
0
} else {
let cnt = self.distance_to_the_leftmost_slot(&slot_max)
- self.distance_to_the_leftmost_slot(&slot_min)
+ self.slot_entry_cnt(&slot_max) as Distance;
cnt as EntryCnt
}
}
pub fn total_by_slot(&self, slot_start: Option<S>, slot_end: Option<S>) -> EntryCnt {
let slot_start = slot_start.unwrap_or(S::MIN);
let slot_end = slot_end.unwrap_or(S::MAX);
if S::MIN == slot_start && S::MAX == slot_end {
self.total_cache
} else {
self.entry_cnt_within_two_slots(slot_start, slot_end)
}
}
pub fn total(&self) -> EntryCnt {
self.total_by_slot(None, None)
}
#[inline(always)]
fn to_storage_slot(&self, logical_slot: S) -> S {
if self.swap_order {
!logical_slot
} else {
logical_slot
}
}
fn transform_range(
&self,
logical_min: Option<S>,
logical_max: Option<S>,
) -> (S, S, bool) {
let min = logical_min.unwrap_or(S::MIN);
let max = logical_max.unwrap_or(S::MAX);
if self.swap_order {
(self.to_storage_slot(max), self.to_storage_slot(min), true)
} else {
(min, max, false)
}
}
}
pub type SlotDex32<K> = SlotDex<u32, K>;
pub type SlotDex64<K> = SlotDex<u64, K>;
pub type SlotDex128<K> = SlotDex<u128, K>;
fn _assert_send_sync() {
fn require<T: Send + Sync>() {}
require::<SlotDex<u64, u64>>();
}
#[cfg(test)]
mod test;