use std::collections::VecDeque;
use std::ops::Range;
use std::sync::Arc;
use ahash::HashSet;
use anyhow::{Result, bail};
use surrealdb_types::ToSql;
use crate::catalog::{DatabaseId, IndexDefinition, IndexId, NamespaceId, Record};
use crate::ctx::FrozenContext;
use crate::err::Error;
use crate::expr::BinaryOperator;
use crate::idx::ft::fulltext::FullTextHitsIterator;
use crate::idx::planner::plan::RangeValue;
use crate::idx::planner::tree::IndexReference;
use crate::idx::seqdocids::DocId;
use crate::idx::{IndexKeyBase, bump_compaction_generation, read_compaction_generation};
use crate::key::index::Index;
use crate::key::index::iu::IndexCountKey;
use crate::kvs::{COUNT_BATCH_SIZE, KVKey, Key, Transaction, Val};
use crate::val::{Array, RecordId, TableName, Value};
pub(crate) type IteratorRef = usize;
#[derive(Debug)]
pub(crate) struct IteratorRecord {
irf: IteratorRef,
doc_id: Option<DocId>,
dist: Option<f64>,
}
impl IteratorRecord {
pub(crate) fn irf(&self) -> IteratorRef {
self.irf
}
pub(crate) fn doc_id(&self) -> Option<DocId> {
self.doc_id
}
pub(crate) fn dist(&self) -> Option<f64> {
self.dist
}
}
impl From<IteratorRef> for IteratorRecord {
fn from(irf: IteratorRef) -> Self {
IteratorRecord {
irf,
doc_id: None,
dist: None,
}
}
}
pub(crate) trait IteratorBatch {
fn empty() -> Self;
fn with_capacity(capacity: usize) -> Self;
fn from_one(record: IndexItemRecord) -> Self;
fn add(&mut self, record: IndexItemRecord);
fn len(&self) -> usize;
fn is_empty(&self) -> bool;
}
impl IteratorBatch for Vec<IndexItemRecord> {
fn empty() -> Self {
Vec::from([])
}
fn with_capacity(capacity: usize) -> Self {
Vec::with_capacity(capacity)
}
fn from_one(record: IndexItemRecord) -> Self {
Vec::from([record])
}
fn add(&mut self, record: IndexItemRecord) {
self.push(record)
}
fn len(&self) -> usize {
Vec::len(self)
}
fn is_empty(&self) -> bool {
Vec::is_empty(self)
}
}
impl IteratorBatch for VecDeque<IndexItemRecord> {
fn empty() -> Self {
VecDeque::from([])
}
fn with_capacity(capacity: usize) -> Self {
VecDeque::with_capacity(capacity)
}
fn from_one(record: IndexItemRecord) -> Self {
VecDeque::from([record])
}
fn add(&mut self, record: IndexItemRecord) {
self.push_back(record)
}
fn len(&self) -> usize {
VecDeque::len(self)
}
fn is_empty(&self) -> bool {
VecDeque::is_empty(self)
}
}
pub(crate) enum RecordIterator {
IndexEqual(IndexEqualThingIterator),
IndexRange(IndexRangeThingIterator),
IndexRangeReverse(IndexRangeReverseThingIterator),
IndexUnion(IndexUnionThingIterator),
IndexJoin(Box<IndexJoinThingIterator>),
IndexCount(IndexCountThingIterator),
UniqueEqual(UniqueEqualThingIterator),
UniqueRange(UniqueRangeThingIterator),
UniqueRangeReverse(UniqueRangeReverseThingIterator),
UniqueUnion(UniqueUnionThingIterator),
UniqueJoin(Box<UniqueJoinThingIterator>),
FullTextMatches(Box<MatchesThingIterator<FullTextHitsIterator>>),
Knn(KnnIterator),
}
impl RecordIterator {
pub(crate) async fn next_batch<B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
txn: &Transaction,
size: u32,
) -> Result<B> {
match self {
Self::IndexEqual(i) => i.next_batch(txn, size).await,
Self::UniqueEqual(i) => i.next_batch(txn, size).await,
Self::IndexRange(i) => i.next_batch(txn, size).await,
Self::IndexRangeReverse(i) => i.next_batch(txn, size).await,
Self::UniqueRange(i) => i.next_batch(txn, size).await,
Self::UniqueRangeReverse(i) => i.next_batch(txn, size).await,
Self::IndexUnion(i) => i.next_batch(ctx, txn, size).await,
Self::UniqueUnion(i) => i.next_batch(ctx, txn, size).await,
Self::FullTextMatches(i) => i.next_batch(ctx, txn, size).await,
Self::Knn(i) => i.next_batch(ctx, size).await,
Self::IndexJoin(i) => Box::pin(i.next_batch(ctx, txn, size)).await,
Self::UniqueJoin(i) => Box::pin(i.next_batch(ctx, txn, size)).await,
Self::IndexCount(_) => {
bail!(Error::unreachable("IndexCount should not be used with next_batch"))
}
}
}
pub(crate) async fn next_count(
&mut self,
ctx: &FrozenContext,
txn: &Transaction,
size: u32,
) -> Result<usize> {
match self {
Self::IndexEqual(i) => i.next_count(txn, size).await,
Self::UniqueEqual(i) => i.next_count(txn, size).await,
Self::IndexRange(i) => i.next_count(txn, size).await,
Self::IndexRangeReverse(i) => i.next_count(txn, size).await,
Self::UniqueRange(i) => i.next_count(txn, size).await,
Self::UniqueRangeReverse(i) => i.next_count(txn, size).await,
Self::IndexUnion(i) => i.next_count(ctx, txn, size).await,
Self::UniqueUnion(i) => i.next_count(ctx, txn, size).await,
Self::FullTextMatches(i) => i.next_count(ctx, txn, size).await,
Self::Knn(i) => i.next_count(ctx, size).await,
Self::IndexJoin(i) => Box::pin(i.next_count(ctx, txn, size)).await,
Self::UniqueJoin(i) => Box::pin(i.next_count(ctx, txn, size)).await,
Self::IndexCount(i) => i.next_count(ctx, txn, size).await,
}
}
}
pub(crate) enum IndexItemRecord {
Key(Arc<RecordId>, IteratorRecord),
KeyValue(Arc<RecordId>, Arc<Record>, IteratorRecord),
}
impl IndexItemRecord {
fn new(t: Arc<RecordId>, ir: IteratorRecord, val: Option<Arc<Record>>) -> Self {
if let Some(val) = val {
Self::KeyValue(t, val, ir)
} else {
Self::Key(t, ir)
}
}
fn new_key(t: RecordId, ir: IteratorRecord) -> Self {
Self::Key(Arc::new(t), ir)
}
fn record_id(&self) -> &RecordId {
match self {
Self::Key(t, _) => t,
Self::KeyValue(t, _, _) => t,
}
}
pub(crate) fn consume(self) -> (Arc<RecordId>, Option<Arc<Record>>, IteratorRecord) {
match self {
Self::Key(t, ir) => (t, None, ir),
Self::KeyValue(t, v, ir) => (t, Some(v), ir),
}
}
}
pub(crate) struct IndexEqualThingIterator {
irf: IteratorRef,
beg: Vec<u8>,
end: Vec<u8>,
}
impl IndexEqualThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
fd: &Array,
) -> Result<Self> {
let (beg, end) = Self::get_beg_end(ns, db, ix, fd)?;
Ok(Self {
irf,
beg,
end,
})
}
fn get_beg_end(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
fd: &Array,
) -> Result<(Vec<u8>, Vec<u8>)> {
Ok(if ix.cols.len() == 1 {
(
Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, fd)?,
Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, fd)?,
)
} else {
(
Index::prefix_ids_composite_beg(ns, db, &ix.table_name, ix.index_id, fd)?,
Index::prefix_ids_composite_end(ns, db, &ix.table_name, ix.index_id, fd)?,
)
})
}
async fn next_scan(
tx: &Transaction,
beg: &mut Vec<u8>,
end: &[u8],
limit: u32,
) -> Result<Vec<(Key, Val)>> {
let min = beg.clone();
let max = end.to_owned();
let res = tx.scan(min..max, limit, 0, None).await?;
if let Some((key, _)) = res.last() {
let mut key = key.clone();
key.push(0x00); *beg = key;
}
Ok(res)
}
async fn next_scan_batch<B: IteratorBatch>(
tx: &Transaction,
irf: IteratorRef,
beg: &mut Vec<u8>,
end: &[u8],
limit: u32,
) -> Result<B> {
let res = Self::next_scan(tx, beg, end, limit).await?;
let mut records = B::with_capacity(res.len());
res.into_iter().try_for_each(|(_, val)| -> Result<()> {
records.add(IndexItemRecord::new_key(revision::from_slice(&val)?, irf.into()));
Ok(())
})?;
Ok(records)
}
async fn next_batch<B: IteratorBatch>(&mut self, tx: &Transaction, limit: u32) -> Result<B> {
Self::next_scan_batch(tx, self.irf, &mut self.beg, &self.end, limit).await
}
async fn next_count(&mut self, tx: &Transaction, limit: u32) -> Result<usize> {
Ok(Self::next_scan(tx, &mut self.beg, &self.end, limit).await?.len())
}
}
struct RangeScan {
beg: Key,
end: Key,
beg_excl_match_checked: bool,
end_excl_match_checked: bool,
}
impl RangeScan {
fn new(beg_key: Key, beg_incl: bool, end_key: Key, end_incl: bool) -> Self {
Self {
beg: beg_key,
end: end_key,
beg_excl_match_checked: beg_incl,
end_excl_match_checked: end_incl,
}
}
fn range(&self) -> Range<Key> {
self.beg.clone()..self.end.clone()
}
fn matches(&mut self, k: &Key) -> bool {
if !self.beg_excl_match_checked && self.beg.eq(k) {
self.beg_excl_match_checked = true;
return false; }
if !self.end_excl_match_checked && self.end.eq(k) {
self.end_excl_match_checked = true;
return false; }
true }
fn matches_end(&mut self) -> bool {
if !self.end_excl_match_checked && self.end.eq(&self.end) {
self.end_excl_match_checked = true;
return false;
}
true
}
}
struct ReverseRangeScan {
r: RangeScan,
beg_incl: bool,
end_incl: bool,
}
impl ReverseRangeScan {
fn new(r: RangeScan) -> Self {
Self {
beg_incl: r.beg_excl_match_checked,
end_incl: r.end_excl_match_checked,
r,
}
}
fn matches_check(&self, k: &Key) -> bool {
if !self.r.beg_excl_match_checked && self.r.beg.eq(k) {
return false;
}
if !self.r.end_excl_match_checked && self.r.end.eq(k) {
return false;
}
true
}
}
pub(crate) struct IndexRangeThingIterator {
irf: IteratorRef,
r: RangeScan,
}
impl IndexRangeThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: RangeValue,
to: RangeValue,
) -> Result<Self> {
Ok(Self {
irf,
r: Self::range_scan(ns, db, ix, from, to)?,
})
}
pub(super) fn full_range(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
) -> Result<Self> {
Self::new(irf, ns, db, ix, RangeValue::default(), RangeValue::default())
}
pub(super) fn compound_range(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexReference,
prefix: &[Value],
ranges: &[(BinaryOperator, Arc<Value>)],
) -> Result<Self> {
let (from, to) = Self::reduce_range(ranges)?;
Ok(Self {
irf,
r: Self::range_scan_prefix(ns, db, ix, prefix, from, to)?,
})
}
fn reduce_range(ranges: &[(BinaryOperator, Arc<Value>)]) -> Result<(RangeValue, RangeValue)> {
let mut from = vec![];
let mut to = vec![];
for (op, v) in ranges {
let key = storekey::encode_vec(v.as_ref()).map_err(|_| Error::Unencodable)?;
match op {
BinaryOperator::LessThan => to.push((key, false, Arc::clone(v))),
BinaryOperator::LessThanEqual => to.push((key, true, Arc::clone(v))),
BinaryOperator::MoreThan => from.push((key, true, Arc::clone(v))),
BinaryOperator::MoreThanEqual => from.push((key, false, Arc::clone(v))),
_ => {
bail!(Error::Unreachable(format!(
"Invalid operator for range extraction {}",
op.to_sql()
)))
}
}
}
let cmp =
|(a1, a2, _): &(Vec<u8>, bool, Arc<Value>),
(b1, b2, _): &(Vec<u8>, bool, Arc<Value>)| { b1.cmp(a1).then_with(|| b2.cmp(a2)) };
from.sort_unstable_by(cmp);
to.sort_unstable_by(cmp);
let from = if let Some((_, inclusivity, val)) = from.into_iter().next() {
RangeValue {
value: Some(val),
inclusive: !inclusivity,
}
} else {
RangeValue::default()
};
let to = if let Some((_, inclusivity, val)) = to.into_iter().next_back() {
RangeValue {
value: Some(val),
inclusive: inclusivity,
}
} else {
RangeValue::default()
};
Ok((from, to))
}
fn range_scan(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: RangeValue,
to: RangeValue,
) -> Result<RangeScan> {
let (from_inclusive, to_inclusive) = (from.inclusive, to.inclusive);
let beg = Self::compute_beg(ns, db, &ix.table_name, ix.index_id, from)?;
let end = Self::compute_end(ns, db, &ix.table_name, ix.index_id, to)?;
Ok(RangeScan::new(beg, from_inclusive, end, to_inclusive))
}
fn compute_beg(
ns: NamespaceId,
db: DatabaseId,
ix_what: &TableName,
index_id: IndexId,
from: RangeValue,
) -> Result<Vec<u8>> {
let Some(value) = from.value else {
return Index::prefix_beg(ns, db, ix_what, index_id);
};
let array = Array::from(vec![value.as_ref().clone()]);
if from.inclusive {
Index::prefix_ids_beg(ns, db, ix_what, index_id, &array)
} else {
Index::prefix_ids_end(ns, db, ix_what, index_id, &array)
}
}
fn compute_end(
ns: NamespaceId,
db: DatabaseId,
ix_what: &TableName,
index_id: IndexId,
to: RangeValue,
) -> Result<Vec<u8>> {
let Some(value) = to.value else {
return Index::prefix_end(ns, db, ix_what, index_id);
};
let array = Array::from(vec![value.as_ref().clone()]);
if to.inclusive {
Index::prefix_ids_end(ns, db, ix_what, index_id, &array)
} else {
Index::prefix_ids_beg(ns, db, ix_what, index_id, &array)
}
}
fn range_scan_prefix(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
prefix: &[Value],
from: RangeValue,
to: RangeValue,
) -> Result<RangeScan> {
let prefix_array: Array = if prefix.is_empty() {
Array(Vec::with_capacity(1))
} else {
Array::from(prefix.to_vec())
};
let (from_inclusive, to_inclusive) = (from.inclusive, to.inclusive);
let beg = match from.value {
None => {
Index::prefix_ids_composite_beg(ns, db, &ix.table_name, ix.index_id, &prefix_array)?
}
Some(v) => Self::compute_beg_with_prefix(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
&v,
from_inclusive,
)?,
};
let end = match to.value {
None => {
Index::prefix_ids_composite_end(ns, db, &ix.table_name, ix.index_id, &prefix_array)?
}
Some(v) => Self::compute_end_with_prefix(
ns,
db,
&ix.table_name,
ix.index_id,
&prefix_array,
&v,
to_inclusive,
)?,
};
Ok(RangeScan::new(beg, from_inclusive, end, to_inclusive))
}
fn compute_beg_with_prefix(
ns: NamespaceId,
db: DatabaseId,
ix_what: &TableName,
index_id: IndexId,
prefix: &Array,
value: &Arc<Value>,
inclusive: bool,
) -> Result<Vec<u8>> {
let mut fd = prefix.clone();
fd.0.push(value.as_ref().clone());
if inclusive {
Index::prefix_ids_beg(ns, db, ix_what, index_id, &fd)
} else {
Index::prefix_ids_end(ns, db, ix_what, index_id, &fd)
}
}
fn compute_end_with_prefix(
ns: NamespaceId,
db: DatabaseId,
ix_what: &TableName,
index_id: IndexId,
prefix: &Array,
value: &Arc<Value>,
inclusive: bool,
) -> Result<Vec<u8>> {
let mut fd = prefix.clone();
fd.0.push(value.as_ref().clone());
if inclusive {
Index::prefix_ids_end(ns, db, ix_what, index_id, &fd)
} else {
Index::prefix_ids_beg(ns, db, ix_what, index_id, &fd)
}
}
async fn next_scan(&mut self, tx: &Transaction, limit: u32) -> Result<Vec<(Key, Val)>> {
let res = tx.scan(self.r.range(), limit, 0, None).await?;
if let Some((key, _)) = res.last() {
self.r.beg.clone_from(key);
self.r.beg.push(0x00);
}
Ok(res)
}
async fn next_keys(&mut self, tx: &Transaction, limit: u32) -> Result<Vec<Key>> {
let res = tx.keys(self.r.range(), limit, 0, None).await?;
if let Some(key) = res.last() {
self.r.beg.clone_from(key);
self.r.beg.push(0x00);
}
Ok(res)
}
async fn next_batch<B: IteratorBatch>(&mut self, tx: &Transaction, limit: u32) -> Result<B> {
let res = self.next_scan(tx, limit).await?;
let mut records = B::with_capacity(res.len());
res.into_iter().filter(|(k, _)| self.r.matches(k)).try_for_each(
|(_, v)| -> Result<()> {
records.add(IndexItemRecord::new_key(revision::from_slice(&v)?, self.irf.into()));
Ok(())
},
)?;
Ok(records)
}
async fn next_count(&mut self, tx: &Transaction, limit: u32) -> Result<usize> {
let res = self.next_keys(tx, limit).await?;
let count = res.into_iter().filter(|k| self.r.matches(k)).count();
Ok(count)
}
}
pub(crate) struct IndexRangeReverseThingIterator {
irf: IteratorRef,
r: ReverseRangeScan,
}
impl IndexRangeReverseThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: RangeValue,
to: RangeValue,
) -> Result<Self> {
Ok(Self {
irf,
r: ReverseRangeScan::new(IndexRangeThingIterator::range_scan(ns, db, ix, from, to)?),
})
}
pub(super) fn full_range(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
) -> Result<Self> {
Self::new(irf, ns, db, ix, RangeValue::default(), RangeValue::default())
}
async fn check_batch_ending(
&mut self,
tx: &Transaction,
limit: &mut u32,
) -> Result<Option<IndexItemRecord>> {
if !self.r.end_incl || !self.r.matches_check(&self.r.r.end) {
return Ok(None);
}
self.r.r.end_excl_match_checked = true;
if let Some(v) = tx.get(&self.r.r.end, None).await? {
*limit -= 1;
Ok(Some(IndexItemRecord::new_key(revision::from_slice(&v)?, self.irf.into())))
} else {
Ok(None)
}
}
async fn check_keys_ending(&mut self, tx: &Transaction, limit: &mut u32) -> Result<bool> {
if !self.r.end_incl || !self.r.matches_check(&self.r.r.end) {
return Ok(false);
}
self.r.r.end_excl_match_checked = true;
if tx.exists(&self.r.r.end, None).await? {
*limit -= 1;
Ok(true)
} else {
Ok(false)
}
}
async fn next_batch<B: IteratorBatch>(
&mut self,
tx: &Transaction,
mut limit: u32,
) -> Result<B> {
let ending = self.check_batch_ending(tx, &mut limit).await?;
let res = if limit > 0 {
tx.scanr(self.r.r.range(), limit, 0, None).await?
} else {
vec![]
};
let mut records = B::with_capacity(res.len() + (ending.is_some() as usize));
if let Some(r) = ending {
records.add(r);
}
let last_key = res.last().map(|(k, _)| k.clone());
res.into_iter().filter(|(k, _)| self.r.r.matches(k)).try_for_each(
|(_, v)| -> Result<()> {
records.add(IndexItemRecord::new_key(revision::from_slice(&v)?, self.irf.into()));
Ok(())
},
)?;
if let Some(key) = last_key {
self.r.r.end = key;
}
if self.r.end_incl {
self.r.end_incl = false;
}
Ok(records)
}
async fn next_count(&mut self, tx: &Transaction, mut limit: u32) -> Result<usize> {
let mut count = self.check_keys_ending(tx, &mut limit).await? as usize;
let res = if limit > 0 {
tx.keysr(self.r.r.range(), limit, 0, None).await?
} else {
vec![]
};
count += res.iter().filter(|k| self.r.r.matches(k)).count();
if let Some(key) = res.last() {
self.r.r.end.clone_from(key);
}
if self.r.end_incl {
self.r.end_incl = false;
}
Ok(count)
}
}
pub(crate) struct IndexUnionThingIterator {
irf: IteratorRef,
values: VecDeque<(Vec<u8>, Vec<u8>)>,
current: Option<(Vec<u8>, Vec<u8>)>,
}
impl IndexUnionThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
fds: &[Array],
) -> Result<Self> {
let mut values: VecDeque<(Vec<u8>, Vec<u8>)> = VecDeque::with_capacity(fds.len());
for fd in fds {
let (beg, end) = IndexEqualThingIterator::get_beg_end(ns, db, ix, fd)?;
values.push_back((beg, end));
}
let current = values.pop_front();
Ok(Self {
irf,
values,
current,
})
}
async fn next_batch<B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<B> {
while let Some(r) = &mut self.current {
if ctx.is_done(None).await? {
break;
}
let records: B =
IndexEqualThingIterator::next_scan_batch(tx, self.irf, &mut r.0, &r.1, limit)
.await?;
if !records.is_empty() {
return Ok(records);
}
self.current = self.values.pop_front();
}
Ok(B::empty())
}
async fn next_count(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<usize> {
while let Some(r) = &mut self.current {
if ctx.is_done(None).await? {
break;
}
let res = IndexEqualThingIterator::next_scan(tx, &mut r.0, &r.1, limit).await?;
if !res.is_empty() {
return Ok(res.len());
}
self.current = self.values.pop_front();
}
Ok(0)
}
}
struct JoinThingIterator {
ns: NamespaceId,
db: DatabaseId,
ix: IndexReference,
remote_iterators: VecDeque<RecordIterator>,
current_remote: Option<RecordIterator>,
current_remote_batch: VecDeque<IndexItemRecord>,
current_local: Option<RecordIterator>,
distinct: HashSet<Key>,
}
impl JoinThingIterator {
pub(super) fn new(
ns: NamespaceId,
db: DatabaseId,
ix: IndexReference,
remote_iterators: VecDeque<RecordIterator>,
) -> Result<Self> {
Ok(Self {
ns,
db,
ix,
current_remote: None,
current_remote_batch: VecDeque::with_capacity(1),
remote_iterators,
current_local: None,
distinct: Default::default(),
})
}
}
impl JoinThingIterator {
async fn next_current_remote_batch(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<bool> {
while !ctx.is_done(None).await? {
if let Some(it) = &mut self.current_remote {
self.current_remote_batch = it.next_batch(ctx, tx, limit).await?;
if !self.current_remote_batch.is_empty() {
return Ok(true);
}
}
self.current_remote = self.remote_iterators.pop_front();
if self.current_remote.is_none() {
break;
}
}
Ok(false)
}
async fn next_current_local<F>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
new_iter: F,
) -> Result<bool>
where
F: Fn(NamespaceId, DatabaseId, &IndexDefinition, Value) -> Result<RecordIterator>,
{
let mut count = 0;
while !ctx.is_done(None).await? {
while let Some(r) = self.current_remote_batch.pop_front() {
if ctx.is_done(Some(count)).await? {
break;
}
let record = r.record_id();
let value: Value = Value::from(record.clone());
let k: Key = revision::to_vec(record)?;
if self.distinct.insert(k) {
self.current_local = Some(new_iter(self.ns, self.db, &self.ix, value)?);
return Ok(true);
}
count += 1;
}
if !self.next_current_remote_batch(ctx, tx, limit).await? {
break;
}
}
Ok(false)
}
async fn next_batch<F, B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
new_iter: F,
) -> Result<B>
where
F: Fn(NamespaceId, DatabaseId, &IndexDefinition, Value) -> Result<RecordIterator> + Copy,
{
while !ctx.is_done(None).await? {
if let Some(current_local) = &mut self.current_local {
let records: B = current_local.next_batch(ctx, tx, limit).await?;
if !records.is_empty() {
return Ok(records);
}
}
if !self.next_current_local(ctx, tx, limit, new_iter).await? {
break;
}
}
Ok(B::empty())
}
async fn next_count<F>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
new_iter: F,
) -> Result<usize>
where
F: Fn(NamespaceId, DatabaseId, &IndexDefinition, Value) -> Result<RecordIterator> + Copy,
{
while !ctx.is_done(None).await? {
if let Some(current_local) = &mut self.current_local {
let count = current_local.next_count(ctx, tx, limit).await?;
if count > 0 {
return Ok(count);
}
}
if !self.next_current_local(ctx, tx, limit, new_iter).await? {
break;
}
}
Ok(0)
}
}
pub(crate) struct IndexJoinThingIterator(IteratorRef, JoinThingIterator);
impl IndexJoinThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: IndexReference,
remote_iterators: VecDeque<RecordIterator>,
) -> Result<Self> {
Ok(Self(irf, JoinThingIterator::new(ns, db, ix, remote_iterators)?))
}
async fn next_batch<B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<B> {
let new_iter = |ns: NamespaceId, db: DatabaseId, ix: &IndexDefinition, value: Value| {
let fd = Array::from(vec![value]);
let it = IndexEqualThingIterator::new(self.0, ns, db, ix, &fd)?;
Ok(RecordIterator::IndexEqual(it))
};
self.1.next_batch(ctx, tx, limit, new_iter).await
}
async fn next_count(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<usize> {
let new_iter = |ns: NamespaceId, db: DatabaseId, ix: &IndexDefinition, value: Value| {
let fd = Array::from(vec![value]);
let it = IndexEqualThingIterator::new(self.0, ns, db, ix, &fd)?;
Ok(RecordIterator::IndexEqual(it))
};
self.1.next_count(ctx, tx, limit, new_iter).await
}
}
pub(crate) struct UniqueEqualThingIterator {
irf: IteratorRef,
inner: UniqueEqualThingInner,
}
enum UniqueEqualThingInner {
PointGet(Option<Key>),
PrefixScan {
beg: Key,
end: Key,
done: bool,
},
}
impl UniqueEqualThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
a: &Array,
) -> Result<Self> {
let inner = if a.is_any_none_or_null() {
let beg = Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, a)?;
let end = Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, a)?;
UniqueEqualThingInner::PrefixScan {
beg,
end,
done: false,
}
} else {
let key = Index::new(ns, db, &ix.table_name, ix.index_id, a, None).encode_key()?;
UniqueEqualThingInner::PointGet(Some(key))
};
Ok(Self {
irf,
inner,
})
}
async fn next_batch<B: IteratorBatch>(&mut self, tx: &Transaction, limit: u32) -> Result<B> {
match &mut self.inner {
UniqueEqualThingInner::PointGet(key) => {
if let Some(key) = key.take()
&& let Some(val) = tx.get(&key, None).await?
{
let rid: RecordId = revision::from_slice(&val)?;
let record = IndexItemRecord::new_key(rid, self.irf.into());
return Ok(B::from_one(record));
}
Ok(B::empty())
}
UniqueEqualThingInner::PrefixScan {
beg,
end,
done,
} => {
if *done {
return Ok(B::empty());
}
let res = tx.scan(beg.clone()..end.clone(), limit, 0, None).await?;
if res.is_empty() {
*done = true;
return Ok(B::empty());
}
if let Some((key, _)) = res.last() {
beg.clone_from(key);
beg.push(0x00);
}
let mut records = B::with_capacity(res.len());
for (_key, val) in res {
let rid: RecordId = revision::from_slice(&val)?;
records.add(IndexItemRecord::new_key(rid, self.irf.into()));
}
Ok(records)
}
}
}
async fn next_count(&mut self, tx: &Transaction, limit: u32) -> Result<usize> {
match &mut self.inner {
UniqueEqualThingInner::PointGet(key) => {
if let Some(key) = key.take()
&& tx.exists(&key, None).await?
{
return Ok(1);
}
Ok(0)
}
UniqueEqualThingInner::PrefixScan {
beg,
end,
done,
} => {
if *done {
return Ok(0);
}
let res = tx.keys(beg.clone()..end.clone(), limit, 0, None).await?;
if res.is_empty() {
*done = true;
return Ok(0);
}
if let Some(key) = res.last() {
beg.clone_from(key);
beg.push(0x00);
}
Ok(res.len())
}
}
}
}
pub(crate) struct UniqueRangeThingIterator {
irf: IteratorRef,
r: RangeScan,
done: bool,
}
impl UniqueRangeThingIterator {
fn range_scan(
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: &RangeValue,
to: &RangeValue,
) -> Result<RangeScan> {
let from_bound = from.value.as_ref().map(|v| (v.as_ref(), from.inclusive));
let to_bound = to.value.as_ref().map(|v| (v.as_ref(), to.inclusive));
let (beg, beg_incl) = Self::compute_beg(ns, db, &ix.table_name, ix.index_id, from_bound)?;
let (end, end_incl) = Self::compute_end(ns, db, &ix.table_name, ix.index_id, to_bound)?;
Ok(RangeScan::new(beg, beg_incl, end, end_incl))
}
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: &RangeValue,
to: &RangeValue,
) -> Result<Self> {
let r = Self::range_scan(ns, db, ix, from, to)?;
Ok(Self {
irf,
r,
done: false,
})
}
pub(super) fn full_range(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
) -> Result<Self> {
let from = RangeValue::default();
let to = RangeValue::default();
Self::new(irf, ns, db, ix, &from, &to)
}
pub(super) fn compound_range(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexReference,
prefix: &[Value],
ranges: &[(BinaryOperator, Arc<Value>)],
) -> Result<Self> {
let (from, to) = IndexRangeThingIterator::reduce_range(ranges)?;
let r = IndexRangeThingIterator::range_scan_prefix(ns, db, ix, prefix, from, to)?;
Ok(Self {
irf,
r,
done: false,
})
}
fn compute_beg(
ns: NamespaceId,
db: DatabaseId,
ix_what: &TableName,
index_id: IndexId,
from: Option<(&Value, bool)>,
) -> Result<(Vec<u8>, bool)> {
let Some((from, inclusive)) = from else {
return Ok((Index::prefix_beg(ns, db, ix_what, index_id)?, true));
};
let array = Array::from(vec![from.clone()]);
if array.is_any_none_or_null() {
let key = if inclusive {
Index::prefix_ids_beg(ns, db, ix_what, index_id, &array)?
} else {
Index::prefix_ids_end(ns, db, ix_what, index_id, &array)?
};
return Ok((key, true));
}
Ok((Index::new(ns, db, ix_what, index_id, &array, None).encode_key()?, inclusive))
}
fn compute_end(
ns: NamespaceId,
db: DatabaseId,
ix_what: &TableName,
index_id: IndexId,
to: Option<(&Value, bool)>,
) -> Result<(Vec<u8>, bool)> {
let Some((to, inclusive)) = to else {
return Ok((Index::prefix_end(ns, db, ix_what, index_id)?, false));
};
let array = Array::from(vec![to.clone()]);
if array.is_any_none_or_null() {
let key = if inclusive {
Index::prefix_ids_end(ns, db, ix_what, index_id, &array)?
} else {
Index::prefix_ids_beg(ns, db, ix_what, index_id, &array)?
};
return Ok((key, false));
}
Ok((Index::new(ns, db, ix_what, index_id, &array, None).encode_key()?, inclusive))
}
async fn next_batch<B: IteratorBatch>(
&mut self,
tx: &Transaction,
mut limit: u32,
) -> Result<B> {
if self.done {
return Ok(B::empty());
}
limit += 1;
let res = tx.scan(self.r.range(), limit, 0, None).await?;
let mut records = B::with_capacity(res.len());
for (k, v) in res {
limit -= 1;
if limit == 0 {
self.r.beg = k;
return Ok(records);
}
if self.r.matches(&k) {
let rid: RecordId = revision::from_slice(&v)?;
records.add(IndexItemRecord::new_key(rid, self.irf.into()));
}
}
if self.r.matches_end()
&& let Some(v) = tx.get(&self.r.end, None).await?
{
let rid: RecordId = revision::from_slice(&v)?;
records.add(IndexItemRecord::new_key(rid, self.irf.into()));
}
self.done = true;
Ok(records)
}
async fn next_count(&mut self, tx: &Transaction, mut limit: u32) -> Result<usize> {
if self.done {
return Ok(0);
}
limit += 1;
let res = tx.keys(self.r.range(), limit, 0, None).await?;
let mut count = 0;
for k in res {
limit -= 1;
if limit == 0 {
self.r.beg = k;
return Ok(count);
}
if self.r.matches(&k) {
count += 1;
}
}
if self.r.matches_end() && tx.exists(&self.r.end, None).await? {
count += 1;
}
self.done = true;
Ok(count)
}
}
pub(crate) struct UniqueRangeReverseThingIterator {
irf: IteratorRef,
r: ReverseRangeScan,
done: bool,
}
impl UniqueRangeReverseThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
from: &RangeValue,
to: &RangeValue,
) -> Result<Self> {
let r = ReverseRangeScan::new(UniqueRangeThingIterator::range_scan(ns, db, ix, from, to)?);
Ok(Self {
irf,
r,
done: false,
})
}
pub(super) fn full_range(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
) -> Result<Self> {
let from = RangeValue::default();
let to = RangeValue::default();
Self::new(irf, ns, db, ix, &from, &to)
}
async fn next_batch<B: IteratorBatch>(
&mut self,
tx: &Transaction,
mut limit: u32,
) -> Result<B> {
if self.done {
return Ok(B::empty());
}
let ending_record = if self.r.end_incl {
self.r.end_incl = false;
if let Some(v) = tx.get(&self.r.r.end, None).await? {
let rid: RecordId = revision::from_slice(&v)?;
let record = IndexItemRecord::new_key(rid, self.irf.into());
limit -= 1;
if limit == 0 {
return Ok(B::from_one(record));
}
Some(record)
} else {
None
}
} else {
None
};
let mut res = tx.scanr(self.r.r.range(), limit, 0, None).await?;
if let Some((k, _)) = res.last() {
self.r.r.end.clone_from(k);
if self.r.r.beg.eq(k) {
self.done = true;
if !self.r.beg_incl {
res.remove(res.len() - 1);
}
}
}
let mut records = B::with_capacity(res.len() + ending_record.is_some() as usize);
if let Some(record) = ending_record {
records.add(record);
}
for (_, v) in res {
let rid: RecordId = revision::from_slice(&v)?;
records.add(IndexItemRecord::new_key(rid, self.irf.into()));
}
Ok(records)
}
async fn next_count(&mut self, tx: &Transaction, mut limit: u32) -> Result<usize> {
if self.done {
return Ok(0);
}
let mut count = 0;
if self.r.end_incl {
self.r.end_incl = false;
if tx.exists(&self.r.r.end, None).await? {
count += 1;
limit -= 1;
if limit == 0 {
return Ok(count);
}
}
}
let mut res = tx.keysr(self.r.r.range(), limit, 0, None).await?;
if let Some(k) = res.last() {
self.r.r.end.clone_from(k);
if self.r.r.beg.eq(k) {
self.done = true;
if !self.r.beg_incl {
res.remove(res.len() - 1);
}
}
}
count += res.len();
Ok(count)
}
}
pub(crate) struct UniqueUnionThingIterator {
irf: IteratorRef,
entries: VecDeque<UniqueUnionEntry>,
}
enum UniqueUnionEntry {
PointGet(Key),
PrefixScan {
beg: Key,
end: Key,
},
}
impl UniqueUnionThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
fds: &[Array],
) -> Result<Self> {
let mut entries = VecDeque::with_capacity(fds.len());
for fd in fds {
if fd.is_any_none_or_null() {
let beg = Index::prefix_ids_beg(ns, db, &ix.table_name, ix.index_id, fd)?;
let end = Index::prefix_ids_end(ns, db, &ix.table_name, ix.index_id, fd)?;
entries.push_back(UniqueUnionEntry::PrefixScan {
beg,
end,
});
} else {
let key = Index::new(ns, db, &ix.table_name, ix.index_id, fd, None).encode_key()?;
entries.push_back(UniqueUnionEntry::PointGet(key));
}
}
Ok(Self {
irf,
entries,
})
}
async fn next_batch<B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<B> {
let limit = limit as usize;
let mut results = B::with_capacity(limit.min(self.entries.len()));
let mut count = 0;
while let Some(entry) = self.entries.pop_front() {
if ctx.is_done(Some(count)).await? {
break;
}
let remaining = (limit - results.len()) as u32;
match entry {
UniqueUnionEntry::PointGet(key) => {
if let Some(val) = tx.get(&key, None).await? {
count += 1;
let rid: RecordId = revision::from_slice(&val)?;
results.add(IndexItemRecord::new_key(rid, self.irf.into()));
}
}
UniqueUnionEntry::PrefixScan {
mut beg,
end,
} => {
let res = tx.scan(beg.clone()..end.clone(), remaining, 0, None).await?;
for (key, val) in &res {
count += 1;
let rid: RecordId = revision::from_slice(val)?;
results.add(IndexItemRecord::new_key(rid, self.irf.into()));
beg.clone_from(key);
}
if !res.is_empty() {
beg.push(0x00);
self.entries.push_front(UniqueUnionEntry::PrefixScan {
beg,
end,
});
}
}
}
if results.len() >= limit {
break;
}
}
Ok(results)
}
async fn next_count(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<usize> {
let limit = limit as usize;
let mut count = 0;
while let Some(entry) = self.entries.pop_front() {
if ctx.is_done(Some(count)).await? {
break;
}
let remaining = (limit - count) as u32;
match entry {
UniqueUnionEntry::PointGet(key) => {
if tx.exists(&key, None).await? {
count += 1;
}
}
UniqueUnionEntry::PrefixScan {
mut beg,
end,
} => {
let res = tx.keys(beg.clone()..end.clone(), remaining, 0, None).await?;
count += res.len();
if let Some(key) = res.last() {
beg.clone_from(key);
beg.push(0x00);
self.entries.push_front(UniqueUnionEntry::PrefixScan {
beg,
end,
});
}
}
}
if count >= limit {
break;
}
}
Ok(count)
}
}
pub(crate) struct UniqueJoinThingIterator(IteratorRef, JoinThingIterator);
impl UniqueJoinThingIterator {
pub(super) fn new(
irf: IteratorRef,
ns: NamespaceId,
db: DatabaseId,
ix: IndexReference,
remote_iterators: VecDeque<RecordIterator>,
) -> Result<Self> {
Ok(Self(irf, JoinThingIterator::new(ns, db, ix, remote_iterators)?))
}
async fn next_batch<B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<B> {
let new_iter = |ns: NamespaceId, db: DatabaseId, ix: &IndexDefinition, value: Value| {
let array = Array::from(vec![value]);
let it = UniqueEqualThingIterator::new(self.0, ns, db, ix, &array)?;
Ok(RecordIterator::UniqueEqual(it))
};
self.1.next_batch(ctx, tx, limit, new_iter).await
}
async fn next_count(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<usize> {
let new_iter = |ns: NamespaceId, db: DatabaseId, ix: &IndexDefinition, value: Value| {
let array = Array::from(vec![value]);
let it = UniqueEqualThingIterator::new(self.0, ns, db, ix, &array)?;
Ok(RecordIterator::UniqueEqual(it))
};
self.1.next_count(ctx, tx, limit, new_iter).await
}
}
pub(crate) trait MatchesHitsIterator {
fn len(&self) -> usize;
async fn next(&mut self, tx: &Transaction) -> Result<Option<(RecordId, DocId)>>;
}
pub(crate) struct MatchesThingIterator<T>
where
T: MatchesHitsIterator,
{
irf: IteratorRef,
hits_left: usize,
hits: Option<T>,
}
impl<T> MatchesThingIterator<T>
where
T: MatchesHitsIterator,
{
pub(super) fn new(irf: IteratorRef, hits: Option<T>) -> Self {
let hits_left = hits.as_ref().map(|h| h.len()).unwrap_or(0);
Self {
irf,
hits,
hits_left,
}
}
async fn next_batch<B: IteratorBatch>(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<B> {
if let Some(hits) = &mut self.hits {
let limit = limit as usize;
let mut count = 0;
let mut records = B::with_capacity(limit.min(self.hits_left));
while limit > records.len() {
if ctx.is_done(Some(count)).await? {
break;
}
if let Some((thg, doc_id)) = hits.next(tx).await? {
let ir = IteratorRecord {
irf: self.irf,
doc_id: Some(doc_id),
dist: None,
};
records.add(IndexItemRecord::new_key(thg, ir));
self.hits_left -= 1;
} else {
break;
}
count += 1;
}
Ok(records)
} else {
Ok(B::empty())
}
}
async fn next_count(
&mut self,
ctx: &FrozenContext,
tx: &Transaction,
limit: u32,
) -> Result<usize> {
if let Some(hits) = &mut self.hits {
let limit = limit as usize;
let mut count = 0;
while limit > count {
if ctx.is_done(Some(count)).await? {
break;
}
if let Some((_, _)) = hits.next(tx).await? {
count += 1;
self.hits_left -= 1;
} else {
break;
}
}
Ok(count)
} else {
Ok(0)
}
}
}
pub(crate) type KnnIteratorResult = (Arc<RecordId>, f64, Option<Arc<Record>>);
pub(crate) struct KnnIterator {
irf: IteratorRef,
res: VecDeque<KnnIteratorResult>,
}
impl KnnIterator {
pub(super) fn new(irf: IteratorRef, res: VecDeque<KnnIteratorResult>) -> Self {
Self {
irf,
res,
}
}
async fn next_batch<B: IteratorBatch>(&mut self, ctx: &FrozenContext, limit: u32) -> Result<B> {
let limit = limit as usize;
let mut count = 0;
let mut records = B::with_capacity(limit.min(self.res.len()));
while limit > records.len() {
if ctx.is_done(Some(count)).await? {
break;
}
if let Some((thing, dist, val)) = self.res.pop_front() {
let ir = IteratorRecord {
irf: self.irf,
doc_id: None,
dist: Some(dist),
};
records.add(IndexItemRecord::new(thing, ir, val));
} else {
break;
}
count += 1;
}
Ok(records)
}
async fn next_count(&mut self, ctx: &FrozenContext, limit: u32) -> Result<usize> {
let limit = limit as usize;
let mut count = 0;
while limit > count {
if ctx.is_done(Some(count)).await? {
break;
}
if self.res.pop_front().is_some() {
count += 1;
} else {
break;
}
}
Ok(count)
}
}
pub(crate) struct IndexCountThingIterator(Option<Range<Key>>);
pub(crate) struct IndexCountCompactionPlan {
generation: Option<u64>,
count: i64,
has_delta: bool,
has_more: bool,
keys: Vec<Key>,
}
impl IndexCountCompactionPlan {
pub(crate) fn has_work(&self) -> bool {
self.has_delta
}
pub(crate) fn has_more(&self) -> bool {
self.has_more
}
}
impl IndexCountThingIterator {
pub(in crate::idx) fn new(
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
ix: IndexId,
) -> Result<Self> {
Ok(Self(Some(IndexCountKey::range(ns, db, tb, ix)?)))
}
async fn next_count(
&mut self,
ctx: &FrozenContext,
txn: &Transaction,
_limit: u32,
) -> Result<usize> {
if let Some(range) = self.0.take() {
let mut count: i64 = 0;
let mut loops = 0;
let mut current_range = Some(range);
while let Some(range) = current_range {
let batch = txn.batch_keys(range, COUNT_BATCH_SIZE, None).await?;
for key in batch.result.iter() {
loops += 1;
ctx.is_done(Some(loops)).await?;
let iu = IndexCountKey::decode_key(key)?;
if iu.pos {
count += iu.count as i64;
} else {
count -= iu.count as i64;
}
}
current_range = batch.next;
ctx.is_done(None).await?;
}
Ok(count as usize)
} else {
Ok(0)
}
}
pub(in crate::idx) async fn prepare_compaction(
&mut self,
ikb: &IndexKeyBase,
txn: &Transaction,
) -> Result<IndexCountCompactionPlan> {
self.prepare_compaction_with_limit(ikb, txn, COUNT_BATCH_SIZE).await
}
async fn prepare_compaction_with_limit(
&mut self,
ikb: &IndexKeyBase,
txn: &Transaction,
limit: u32,
) -> Result<IndexCountCompactionPlan> {
let generation = read_compaction_generation(txn, &ikb.new_iv_key()).await?;
let Some(range) = self.0.take() else {
return Ok(IndexCountCompactionPlan {
generation,
count: 0,
has_delta: false,
has_more: false,
keys: Vec::new(),
});
};
let mut count: i64 = 0;
let mut has_delta = false;
let mut has_more = false;
let mut keys = Vec::new();
let mut delta_count = 0;
let mut loops = 0;
let mut current_range = Some(range.clone());
let limit = limit.max(1);
while let Some(r) = current_range.take() {
let batch = txn.batch_keys(r, limit.saturating_add(1), None).await?;
for key in batch.result.iter() {
loops += 1;
if loops % 1000 == 0 {
yield_now!()
}
let iu = IndexCountKey::decode_key(key)?;
if iu.uid.is_some() && delta_count >= limit {
has_more = true;
current_range = None;
break;
}
if iu.pos {
count += iu.count as i64;
} else {
count -= iu.count as i64;
}
if iu.uid.is_some() {
has_delta = true;
delta_count += 1;
}
keys.push(key.clone());
}
if has_more {
break;
}
current_range = batch.next;
}
has_more |= current_range.is_some();
Ok(IndexCountCompactionPlan {
generation,
count,
has_delta,
has_more,
keys,
})
}
pub(in crate::idx) async fn apply_compaction(
ikb: &IndexKeyBase,
txn: &Transaction,
plan: IndexCountCompactionPlan,
) -> Result<bool> {
if !plan.has_work() {
return Ok(false);
}
if !bump_compaction_generation(txn, &ikb.new_iv_key(), plan.generation).await? {
return Ok(false);
}
for key in plan.keys {
txn.del(&key).await?;
}
let count = plan.count;
let pos = count.is_positive();
let count = count.unsigned_abs();
let compact_key =
IndexCountKey::new(ikb.ns(), ikb.db(), ikb.table(), ikb.index(), None, pos, count);
txn.set(&compact_key, &()).await?;
Ok(true)
}
#[cfg(test)]
pub(in crate::idx) async fn compaction(
&mut self,
ikb: &IndexKeyBase,
txn: &Transaction,
) -> Result<()> {
let plan = self.prepare_compaction(ikb, txn).await?;
Self::apply_compaction(ikb, txn, plan).await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use uuid::Uuid;
use super::*;
use crate::catalog::{DatabaseId, IndexId, NamespaceId};
use crate::idx::IndexKeyBase;
use crate::key::index::iu::IndexCountKey;
use crate::kvs::Datastore;
use crate::kvs::LockType::Optimistic;
use crate::kvs::TransactionType::{Read, Write};
async fn count_value(ds: &Datastore, ikb: &IndexKeyBase) -> usize {
let mut count_iter =
IndexCountThingIterator::new(ikb.ns(), ikb.db(), ikb.table(), ikb.index()).unwrap();
let tx = Arc::new(ds.transaction(Read, Optimistic).await.unwrap());
let mut ctx = ds.setup_ctx().unwrap();
ctx.set_transaction(Arc::clone(&tx));
let ctx = ctx.freeze();
let count = count_iter.next_count(&ctx, &ctx.tx(), u32::MAX).await.unwrap();
tx.cancel().await.unwrap();
count
}
#[tokio::test]
async fn test_consecutive_compactions_do_not_fail() {
let ns = NamespaceId(1);
let db = DatabaseId(2);
let tb: TableName = "test_tb".into();
let ix = IndexId(3);
let ikb = IndexKeyBase::new(ns, db, tb.clone(), ix);
let ds = Datastore::new("memory").await.unwrap();
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let uid1 = (Uuid::new_v4(), Uuid::new_v4());
let uid2 = (Uuid::new_v4(), Uuid::new_v4());
let k1 = IndexCountKey::new(ns, db, &tb, ix, Some(uid1), true, 10);
let k2 = IndexCountKey::new(ns, db, &tb, ix, Some(uid2), true, 5);
tx.set(&k1, &()).await.unwrap();
tx.set(&k2, &()).await.unwrap();
tx.commit().await.unwrap();
}
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.compaction(&ikb, &tx)
.await
.unwrap();
tx.commit().await.unwrap();
}
{
let count = count_value(&ds, &ikb).await;
assert_eq!(count, 15, "first compaction should yield count 15");
}
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let uid3 = (Uuid::new_v4(), Uuid::new_v4());
let k3 = IndexCountKey::new(ns, db, &tb, ix, Some(uid3), true, 7);
tx.set(&k3, &()).await.unwrap();
tx.commit().await.unwrap();
}
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.compaction(&ikb, &tx)
.await
.unwrap();
tx.commit().await.unwrap();
}
{
let count = count_value(&ds, &ikb).await;
assert_eq!(count, 22, "second compaction should yield count 22 (15 + 7)");
}
}
#[tokio::test]
async fn count_compaction_preserves_post_snapshot_deltas() {
let ns = NamespaceId(1);
let db = DatabaseId(2);
let tb: TableName = "test_tb".into();
let ix = IndexId(3);
let ikb = IndexKeyBase::new(ns, db, tb.clone(), ix);
let ds = Datastore::new("memory").await.unwrap();
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let uid1 = (Uuid::new_v4(), Uuid::new_v4());
let k1 = IndexCountKey::new(ns, db, &tb, ix, Some(uid1), true, 10);
tx.set(&k1, &()).await.unwrap();
tx.commit().await.unwrap();
}
let plan = {
let tx = ds.transaction(Read, Optimistic).await.unwrap();
let plan = IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.prepare_compaction(&ikb, &tx)
.await
.unwrap();
tx.cancel().await.unwrap();
plan
};
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let uid2 = (Uuid::new_v4(), Uuid::new_v4());
let k2 = IndexCountKey::new(ns, db, &tb, ix, Some(uid2), true, 7);
tx.set(&k2, &()).await.unwrap();
tx.commit().await.unwrap();
}
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
assert!(IndexCountThingIterator::apply_compaction(&ikb, &tx, plan).await.unwrap());
tx.commit().await.unwrap();
}
let tx = ds.transaction(Read, Optimistic).await.unwrap();
assert_eq!(tx.get(&ikb.new_iv_key(), None).await.unwrap(), Some(1));
let range = IndexCountKey::range(ns, db, &tb, ix).unwrap();
assert_eq!(
tx.count(range, None).await.unwrap(),
2,
"compacted root and post-snapshot delta should remain"
);
tx.cancel().await.unwrap();
assert_eq!(count_value(&ds, &ikb).await, 17);
}
#[tokio::test]
async fn count_compaction_batches_visible_deltas() {
let ns = NamespaceId(1);
let db = DatabaseId(2);
let tb: TableName = "test_tb".into();
let ix = IndexId(3);
let ikb = IndexKeyBase::new(ns, db, tb.clone(), ix);
let ds = Datastore::new("memory").await.unwrap();
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
for count in [10, 5, 7] {
let uid = (Uuid::new_v4(), Uuid::new_v4());
let key = IndexCountKey::new(ns, db, &tb, ix, Some(uid), true, count);
tx.set(&key, &()).await.unwrap();
}
tx.commit().await.unwrap();
}
let plan = {
let tx = ds.transaction(Read, Optimistic).await.unwrap();
let plan = IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.prepare_compaction_with_limit(&ikb, &tx, 2)
.await
.unwrap();
tx.cancel().await.unwrap();
plan
};
assert!(plan.has_work());
assert!(plan.has_more());
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
assert!(IndexCountThingIterator::apply_compaction(&ikb, &tx, plan).await.unwrap());
tx.commit().await.unwrap();
}
let tx = ds.transaction(Read, Optimistic).await.unwrap();
assert_eq!(
tx.count(IndexCountKey::range(ns, db, &tb, ix).unwrap(), None).await.unwrap(),
2,
"first batch should leave compacted root plus one residual delta"
);
tx.cancel().await.unwrap();
assert_eq!(count_value(&ds, &ikb).await, 22);
let plan = {
let tx = ds.transaction(Read, Optimistic).await.unwrap();
let plan = IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.prepare_compaction_with_limit(&ikb, &tx, 2)
.await
.unwrap();
tx.cancel().await.unwrap();
plan
};
assert!(plan.has_work());
assert!(!plan.has_more());
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
assert!(IndexCountThingIterator::apply_compaction(&ikb, &tx, plan).await.unwrap());
tx.commit().await.unwrap();
}
let tx = ds.transaction(Read, Optimistic).await.unwrap();
assert_eq!(
tx.count(IndexCountKey::range(ns, db, &tb, ix).unwrap(), None).await.unwrap(),
1,
"second batch should collapse all deltas into one compacted root"
);
tx.cancel().await.unwrap();
assert_eq!(count_value(&ds, &ikb).await, 22);
}
#[tokio::test]
async fn count_compaction_generation_allows_only_one_winner() {
let ns = NamespaceId(1);
let db = DatabaseId(2);
let tb: TableName = "test_tb".into();
let ix = IndexId(3);
let ikb = IndexKeyBase::new(ns, db, tb.clone(), ix);
let ds = Datastore::new("memory").await.unwrap();
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let uid1 = (Uuid::new_v4(), Uuid::new_v4());
let k1 = IndexCountKey::new(ns, db, &tb, ix, Some(uid1), true, 10);
tx.set(&k1, &()).await.unwrap();
tx.commit().await.unwrap();
}
let plan1 = {
let tx = ds.transaction(Read, Optimistic).await.unwrap();
let plan = IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.prepare_compaction(&ikb, &tx)
.await
.unwrap();
tx.cancel().await.unwrap();
plan
};
let plan2 = {
let tx = ds.transaction(Read, Optimistic).await.unwrap();
let plan = IndexCountThingIterator::new(ns, db, &tb, ix)
.unwrap()
.prepare_compaction(&ikb, &tx)
.await
.unwrap();
tx.cancel().await.unwrap();
plan
};
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
assert!(IndexCountThingIterator::apply_compaction(&ikb, &tx, plan1).await.unwrap());
tx.commit().await.unwrap();
}
{
let tx = ds.transaction(Write, Optimistic).await.unwrap();
assert!(!IndexCountThingIterator::apply_compaction(&ikb, &tx, plan2).await.unwrap());
tx.cancel().await.unwrap();
}
let tx = ds.transaction(Read, Optimistic).await.unwrap();
assert_eq!(tx.get(&ikb.new_iv_key(), None).await.unwrap(), Some(1));
tx.cancel().await.unwrap();
assert_eq!(count_value(&ds, &ikb).await, 10);
}
}