use anyhow::Result;
use crate::catalog::{DatabaseId, IndexDefinition, NamespaceId};
use crate::exec::index::access_path::RangeBound;
use crate::expr::BinaryOperator;
use crate::idx::planner::ScanDirection;
use crate::key::index::Index;
use crate::kvs::{KVKey, Key, Transaction, Val};
use crate::val::{Array, RecordId, Value};
const INDEX_BATCH_SIZE: u32 = 1000;
fn decode_record_ids(res: Vec<(Key, Val)>) -> Result<Vec<RecordId>> {
let mut records = Vec::with_capacity(res.len());
for (_, val) in res {
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
Ok(records)
}
pub(crate) struct IndexEqualIterator {
beg: Vec<u8>,
end: Vec<u8>,
reverse: bool,
done: bool,
}
impl IndexEqualIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
value: &Value,
) -> Result<Self> {
Self::with_direction(ns, db, ix, value, false)
}
pub(crate) fn with_direction(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
value: &Value,
reverse: bool,
) -> Result<Self> {
let array = Array::from(vec![value.clone()]);
let beg = Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &array)?;
let end = Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &array)?;
Ok(Self {
beg,
end,
reverse,
done: false,
})
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
if self.reverse {
self.next_batch_reverse(tx).await
} else {
self.next_batch_forward(tx).await
}
}
async fn next_batch_forward(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
let res = tx.scan(self.beg.clone()..self.end.clone(), INDEX_BATCH_SIZE, 0, None).await?;
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.beg.clone_from(key);
self.beg.push(0x00);
}
decode_record_ids(res)
}
async fn next_batch_reverse(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
let res = tx.scanr(self.beg.clone()..self.end.clone(), INDEX_BATCH_SIZE, 0, None).await?;
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.end.clone_from(key);
}
let mut records = Vec::with_capacity(res.len());
for (_, val) in res {
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
Ok(records)
}
}
pub(crate) struct UniqueEqualIterator {
inner: UniqueEqualInner,
}
enum UniqueEqualInner {
PointGet(Option<Key>),
PrefixScan {
beg: Key,
end: Key,
done: bool,
},
}
impl UniqueEqualIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
value: &Value,
) -> Result<Self> {
let array = Array::from(vec![value.clone()]);
let inner = if array.is_any_none_or_null() {
let beg = Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &array)?;
let end = Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &array)?;
UniqueEqualInner::PrefixScan {
beg,
end,
done: false,
}
} else {
let key = Index::new(ns, db, &ix.table_name, ix.index_id, &array, None).encode_key()?;
UniqueEqualInner::PointGet(Some(key))
};
Ok(Self {
inner,
})
}
pub async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
match &mut self.inner {
UniqueEqualInner::PointGet(key) => {
let Some(key) = key.take() else {
return Ok(Vec::new());
};
if let Some(val) = tx.get(&key, None).await? {
let rid: RecordId = revision::from_slice(&val)?;
Ok(vec![rid])
} else {
Ok(Vec::new())
}
}
UniqueEqualInner::PrefixScan {
beg,
end,
done,
} => {
if *done {
return Ok(Vec::new());
}
let res = tx.scan(beg.clone()..end.clone(), INDEX_BATCH_SIZE, 0, None).await?;
if res.is_empty() {
*done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
beg.clone_from(key);
beg.push(0x00);
}
decode_record_ids(res)
}
}
}
}
fn compute_index_range_beg_key(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
) -> Result<(Key, bool)> {
if let Some(from) = from {
let array = Array::from(vec![from.value.clone()]);
if from.inclusive {
Ok((Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &array)?, true))
} else {
Ok((Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &array)?, false))
}
} else {
Ok((Index::prefix_beg(ns, db, &ix.table_name, ix.index_id)?, true))
}
}
fn compute_index_range_end_key(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
to: Option<&RangeBound>,
) -> Result<(Key, bool)> {
if let Some(to) = to {
let array = Array::from(vec![to.value.clone()]);
if to.inclusive {
Ok((Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &array)?, true))
} else {
Ok((Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &array)?, false))
}
} else {
Ok((Index::prefix_end(ns, db, &ix.table_name, ix.index_id)?, true))
}
}
pub(crate) struct IndexRangeForwardIterator {
beg: Key,
end: Key,
beg_checked: bool,
done: bool,
}
impl IndexRangeForwardIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
to: Option<&RangeBound>,
) -> Result<Self> {
let (beg, beg_inclusive) = compute_index_range_beg_key(ns, db, ix, from)?;
let (end, _end_inclusive) = compute_index_range_end_key(ns, db, ix, to)?;
Ok(Self {
beg,
end,
beg_checked: beg_inclusive,
done: false,
})
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let check_exclusive_beg = if self.beg_checked {
None
} else {
Some(self.beg.clone())
};
let res = tx.scan(self.beg.clone()..self.end.clone(), INDEX_BATCH_SIZE, 0, None).await?;
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.beg.clone_from(key);
self.beg.push(0x00);
}
self.beg_checked = true;
let mut records = Vec::with_capacity(res.len());
for (key, val) in res {
if let Some(ref exclusive_key) = check_exclusive_beg
&& key == *exclusive_key
{
continue;
}
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
Ok(records)
}
}
pub(crate) struct IndexRangeBackwardIterator {
beg: Key,
end: Key,
end_checked: bool,
exclude_beg_key: Option<Key>,
done: bool,
}
impl IndexRangeBackwardIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
to: Option<&RangeBound>,
) -> Result<Self> {
let (beg, beg_inclusive) = compute_index_range_beg_key(ns, db, ix, from)?;
let (end, end_inclusive) = compute_index_range_end_key(ns, db, ix, to)?;
let exclude_beg_key = if beg_inclusive {
None
} else {
Some(beg.clone())
};
Ok(Self {
beg,
end,
end_checked: end_inclusive,
exclude_beg_key,
done: false,
})
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let check_exclusive_end = if self.end_checked {
None
} else {
Some(self.end.clone())
};
let res = tx.scanr(self.beg.clone()..self.end.clone(), INDEX_BATCH_SIZE, 0, None).await?;
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.end.clone_from(key);
}
self.end_checked = true;
let mut records = Vec::with_capacity(res.len());
for (key, val) in res {
if let Some(ref exclusive_key) = check_exclusive_end
&& key == *exclusive_key
{
continue;
}
if let Some(ref beg_key) = self.exclude_beg_key
&& key == *beg_key
{
continue;
}
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
Ok(records)
}
}
pub(crate) enum IndexRangeIterator {
Forward(IndexRangeForwardIterator),
Backward(IndexRangeBackwardIterator),
}
impl IndexRangeIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
to: Option<&RangeBound>,
direction: ScanDirection,
) -> Result<Self> {
match direction {
ScanDirection::Forward => {
Ok(Self::Forward(IndexRangeForwardIterator::new(ns, db, ix, from, to)?))
}
ScanDirection::Backward => {
Ok(Self::Backward(IndexRangeBackwardIterator::new(ns, db, ix, from, to)?))
}
}
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
match self {
Self::Forward(iter) => iter.next_batch(tx).await,
Self::Backward(iter) => iter.next_batch(tx).await,
}
}
}
fn compute_unique_range_beg_key(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
) -> Result<(Key, bool)> {
if let Some(from) = from {
let array = Array::from(vec![from.value.clone()]);
if array.is_any_none_or_null() {
let key = if from.inclusive {
Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &array)?
} else {
Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &array)?
};
Ok((key, true))
} else {
let key = Index::new(ns, db, &ix.table_name, ix.index_id, &array, None).encode_key()?;
Ok((key, from.inclusive))
}
} else {
Ok((Index::prefix_beg(ns, db, &ix.table_name, ix.index_id)?, true))
}
}
fn compute_unique_range_end_key(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
to: Option<&RangeBound>,
) -> Result<(Key, bool)> {
if let Some(to) = to {
let array = Array::from(vec![to.value.clone()]);
if array.is_any_none_or_null() {
let key = if to.inclusive {
Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &array)?
} else {
Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &array)?
};
Ok((key, false))
} else {
let key = Index::new(ns, db, &ix.table_name, ix.index_id, &array, None).encode_key()?;
Ok((key, to.inclusive))
}
} else {
Ok((Index::prefix_end(ns, db, &ix.table_name, ix.index_id)?, false))
}
}
pub(crate) struct UniqueRangeForwardIterator {
beg: Key,
end: Key,
beg_checked: bool,
end_inclusive: bool,
done: bool,
}
impl UniqueRangeForwardIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
to: Option<&RangeBound>,
) -> Result<Self> {
let (beg, beg_inclusive) = compute_unique_range_beg_key(ns, db, ix, from)?;
let (end, end_inclusive) = compute_unique_range_end_key(ns, db, ix, to)?;
Ok(Self {
beg,
end,
beg_checked: beg_inclusive,
end_inclusive,
done: false,
})
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let check_exclusive_beg = if self.beg_checked {
None
} else {
Some(self.beg.clone())
};
let limit = INDEX_BATCH_SIZE + 1;
let res = tx.scan(self.beg.clone()..self.end.clone(), limit, 0, None).await?;
if res.is_empty() {
self.done = true;
if self.end_inclusive
&& let Some(val) = tx.get(&self.end, None).await?
{
let rid: RecordId = revision::from_slice(&val)?;
return Ok(vec![rid]);
}
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.beg.clone_from(key);
self.beg.push(0x00);
}
self.beg_checked = true;
let mut records = Vec::with_capacity(res.len());
for (key, val) in res {
if let Some(ref exclusive_key) = check_exclusive_beg
&& key == *exclusive_key
{
continue;
}
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
Ok(records)
}
}
pub(crate) struct UniqueRangeBackwardIterator {
beg: Key,
end: Key,
end_checked: bool,
end_inclusive: bool,
original_end: Key,
exclude_beg_key: Option<Key>,
done: bool,
}
impl UniqueRangeBackwardIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
to: Option<&RangeBound>,
) -> Result<Self> {
let (beg, beg_inclusive) = compute_unique_range_beg_key(ns, db, ix, from)?;
let (end, end_inclusive) = compute_unique_range_end_key(ns, db, ix, to)?;
let exclude_beg_key = if beg_inclusive {
None
} else {
Some(beg.clone())
};
Ok(Self {
beg,
original_end: end.clone(),
end,
end_checked: end_inclusive,
end_inclusive,
exclude_beg_key,
done: false,
})
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let mut records = Vec::new();
if self.end_inclusive {
self.end_inclusive = false;
if let Some(val) = tx.get(&self.original_end, None).await? {
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
}
let check_exclusive_end = if self.end_checked {
None
} else {
Some(self.end.clone())
};
let limit = INDEX_BATCH_SIZE + 1;
let res = tx.scanr(self.beg.clone()..self.end.clone(), limit, 0, None).await?;
if res.is_empty() {
self.done = true;
if records.is_empty() {
return Ok(Vec::new());
}
return Ok(records);
}
if let Some((key, _)) = res.last() {
self.end.clone_from(key);
}
self.end_checked = true;
records.reserve(res.len());
for (key, val) in res {
if let Some(ref exclusive_key) = check_exclusive_end
&& key == *exclusive_key
{
continue;
}
if let Some(ref beg_key) = self.exclude_beg_key
&& key == *beg_key
{
continue;
}
let rid: RecordId = revision::from_slice(&val)?;
records.push(rid);
}
Ok(records)
}
}
pub(crate) enum UniqueRangeIterator {
Forward(UniqueRangeForwardIterator),
Backward(UniqueRangeBackwardIterator),
}
impl UniqueRangeIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: Option<&RangeBound>,
to: Option<&RangeBound>,
direction: ScanDirection,
) -> Result<Self> {
match direction {
ScanDirection::Forward => {
Ok(Self::Forward(UniqueRangeForwardIterator::new(ns, db, ix, from, to)?))
}
ScanDirection::Backward => {
Ok(Self::Backward(UniqueRangeBackwardIterator::new(ns, db, ix, from, to)?))
}
}
}
pub(crate) async fn next_batch(&mut self, tx: &Transaction) -> Result<Vec<RecordId>> {
match self {
Self::Forward(iter) => iter.next_batch(tx).await,
Self::Backward(iter) => iter.next_batch(tx).await,
}
}
}
pub(crate) struct CompoundEqualIterator {
beg: Vec<u8>,
end: Vec<u8>,
done: bool,
direction: ScanDirection,
}
impl CompoundEqualIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
prefix: &[Value],
range: Option<&(BinaryOperator, Value)>,
direction: ScanDirection,
) -> Result<Self> {
let (beg, end) = compute_compound_key_range(ns, db, ix, prefix, range)?;
Ok(Self {
beg,
end,
done: false,
direction,
})
}
pub(crate) async fn next_batch(
&mut self,
tx: &Transaction,
limit: u32,
) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let scan_limit = limit.min(INDEX_BATCH_SIZE);
let res = match self.direction {
ScanDirection::Forward => {
tx.scan(self.beg.clone()..self.end.clone(), scan_limit, 0, None).await?
}
ScanDirection::Backward => {
tx.scanr(self.beg.clone()..self.end.clone(), scan_limit, 0, None).await?
}
};
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
match self.direction {
ScanDirection::Forward => {
self.beg.clone_from(key);
self.beg.push(0x00);
}
ScanDirection::Backward => {
self.end.clone_from(key);
}
}
}
decode_record_ids(res)
}
}
pub(crate) struct CompoundRangeForwardIterator {
beg: Vec<u8>,
end: Vec<u8>,
done: bool,
}
impl CompoundRangeForwardIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
prefix: &[Value],
range: &(BinaryOperator, Value),
) -> Result<Self> {
let (beg, end) = compute_compound_key_range(ns, db, ix, prefix, Some(range))?;
Ok(Self {
beg,
end,
done: false,
})
}
pub(crate) async fn next_batch(
&mut self,
tx: &Transaction,
limit: u32,
) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let scan_limit = limit.min(INDEX_BATCH_SIZE);
let res = tx.scan(self.beg.clone()..self.end.clone(), scan_limit, 0, None).await?;
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.beg.clone_from(key);
self.beg.push(0x00);
}
decode_record_ids(res)
}
}
pub(crate) enum CompoundRangeIterator {
Forward(CompoundRangeForwardIterator),
Backward(CompoundRangeBackwardIterator),
}
impl CompoundRangeIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
prefix: &[Value],
range: &(BinaryOperator, Value),
direction: ScanDirection,
) -> Result<Self> {
match direction {
ScanDirection::Forward => {
Ok(Self::Forward(CompoundRangeForwardIterator::new(ns, db, ix, prefix, range)?))
}
ScanDirection::Backward => {
Ok(Self::Backward(CompoundRangeBackwardIterator::new(ns, db, ix, prefix, range)?))
}
}
}
pub(crate) async fn next_batch(
&mut self,
tx: &Transaction,
limit: u32,
) -> Result<Vec<RecordId>> {
match self {
Self::Forward(iter) => iter.next_batch(tx, limit).await,
Self::Backward(iter) => iter.next_batch(tx, limit).await,
}
}
}
pub(crate) struct CompoundRangeBackwardIterator {
beg: Vec<u8>,
end: Vec<u8>,
done: bool,
}
impl CompoundRangeBackwardIterator {
pub(crate) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
prefix: &[Value],
range: &(BinaryOperator, Value),
) -> Result<Self> {
let (beg, end) = compute_compound_key_range(ns, db, ix, prefix, Some(range))?;
Ok(Self {
beg,
end,
done: false,
})
}
pub(crate) async fn next_batch(
&mut self,
tx: &Transaction,
limit: u32,
) -> Result<Vec<RecordId>> {
if self.done {
return Ok(Vec::new());
}
let scan_limit = limit.min(INDEX_BATCH_SIZE);
let res = tx.scanr(self.beg.clone()..self.end.clone(), scan_limit, 0, None).await?;
if res.is_empty() {
self.done = true;
return Ok(Vec::new());
}
if let Some((key, _)) = res.last() {
self.end.clone_from(key);
}
decode_record_ids(res)
}
}
fn compute_compound_key_range(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
prefix: &[Value],
range: Option<&(BinaryOperator, Value)>,
) -> Result<(Vec<u8>, Vec<u8>)> {
let prefix_array = Array::from(prefix.to_vec());
if let Some((op, val)) = range {
let mut key_values: Vec<Value> = prefix.to_vec();
key_values.push(val.clone());
let key_array = Array::from(key_values);
match op {
BinaryOperator::Equal | BinaryOperator::ExactEqual => {
let beg = Index::prefix_ids_composite_beg(
ns,
db,
&ix.table_name,
ix.index_id,
&key_array,
)?;
let end = Index::prefix_ids_composite_end(
ns,
db,
&ix.table_name,
ix.index_id,
&key_array,
)?;
Ok((beg, end))
}
BinaryOperator::MoreThan => {
let beg = Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &key_array)?;
let end = Index::prefix_ids_composite_end(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
)?;
Ok((beg, end))
}
BinaryOperator::MoreThanEqual => {
let beg = Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &key_array)?;
let end = Index::prefix_ids_composite_end(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
)?;
Ok((beg, end))
}
BinaryOperator::LessThan => {
let beg = Index::prefix_ids_composite_beg(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
)?;
let end = Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, &key_array)?;
Ok((beg, end))
}
BinaryOperator::LessThanEqual => {
let beg = Index::prefix_ids_composite_beg(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
)?;
let end = Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, &key_array)?;
Ok((beg, end))
}
_ => {
let beg = Index::prefix_ids_composite_beg(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
)?;
let end = Index::prefix_ids_composite_end(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
)?;
Ok((beg, end))
}
}
} else {
let beg =
Index::prefix_ids_composite_beg(ns, db, &ix.table_name, ix.index_id, &prefix_array)?;
let end =
Index::prefix_ids_composite_end(ns, db, &ix.table_name, ix.index_id, &prefix_array)?;
Ok((beg, end))
}
}