use std::borrow::Cow;
use std::ops::{Bound, Range};
use std::sync::Arc;
use std::vec;
use anyhow::{Result, bail};
use reblessive::tree::Stk;
use crate::catalog::providers::TableProvider;
use crate::catalog::{DatabaseId, NamespaceId, Record};
use crate::ctx::{Context, FrozenContext};
use crate::dbs::distinct::SyncDistinct;
use crate::dbs::{Iterable, Iterator, Operable, Options, Processable, Statement};
use crate::doc::{DocumentContext, NsDbCtx};
use crate::err::Error;
use crate::expr::dir::Dir;
use crate::expr::lookup::{ComputedLookupSubject, LookupKind};
use crate::expr::statements::relate::RelateThrough;
use crate::idx::planner::iterators::{IndexItemRecord, IteratorRef, RecordIterator};
use crate::idx::planner::{IterationStage, RecordStrategy, ScanDirection};
use crate::key::{graph, record, r#ref};
use crate::kvs::{KVKey, KVValue, Key, NORMAL_BATCH_SIZE, Transaction, Val};
use crate::val::{RecordId, RecordIdKey, RecordIdKeyRange, TableName, Value};
impl Iterable {
#[instrument(level = "trace", name = "Iterable::iterate", skip_all)]
pub(super) async fn iterate(
self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
stm: &Statement<'_>,
ite: &mut Iterator,
dis: Option<&mut SyncDistinct>,
) -> Result<()> {
if !self.iteration_stage_check(ctx) {
return Ok(());
}
let txn = ctx.tx();
let mut concurrent_collector = ConcurrentCollector {
stk,
ctx,
opt,
txn: &txn,
stm,
ite,
};
if let Some(dis) = dis {
let mut distinct_collector = ConcurrentDistinctCollector {
coll: concurrent_collector,
dis,
};
distinct_collector.collect_iterable(ctx, opt, self).await?;
} else {
concurrent_collector.collect_iterable(ctx, opt, self).await?;
}
Ok(())
}
fn iteration_stage_check(&self, ctx: &FrozenContext) -> bool {
match self {
Iterable::Table(_doc_ctx, tb, _, _) | Iterable::Index(_doc_ctx, tb, _, _) => {
if let Some(IterationStage::BuildKnn) = ctx.get_iteration_stage()
&& let Some(qp) = ctx.get_query_planner()
&& let Some(exe) = qp.get_query_executor(tb)
{
return exe.has_bruteforce_knn();
}
}
_ => {}
}
true
}
}
pub(super) enum Collectable {
Lookup(DocumentContext, LookupKind, Key),
RangeKey(DocumentContext, Key),
TableKey(DocumentContext, Key),
Relatable {
doc_ctx: DocumentContext,
f: RecordId,
v: RelateThrough,
w: RecordId,
o: Option<Value>,
},
RecordId(DocumentContext, RecordId),
GenerateRecordId(DocumentContext, TableName),
Value(NsDbCtx, Value),
Defer(DocumentContext, RecordId),
Mergeable(DocumentContext, TableName, Option<RecordIdKey>, Value),
KeyVal(DocumentContext, Key, Val),
Count(DocumentContext, usize),
IndexItem(DocumentContext, IndexItemRecord),
IndexItemKey(DocumentContext, IndexItemRecord),
}
impl Collectable {
#[instrument(level = "trace", name = "Collectable::prepare", skip_all)]
pub(super) async fn prepare(
self,
ctx: &FrozenContext,
opt: &Options,
txn: &Transaction,
rid_only: bool,
) -> Result<Processable> {
match self {
Self::Lookup(doc_ctx, kind, key) => {
Self::process_lookup(doc_ctx, ctx, opt, txn, kind, key, rid_only).await
}
Self::RangeKey(doc_ctx, key) => Self::process_range_key(doc_ctx, key).await,
Self::TableKey(doc_ctx, key) => Self::process_table_key(doc_ctx, key).await,
Self::Relatable {
doc_ctx,
f,
v,
w,
o,
} => Self::process_relatable(doc_ctx, txn, f, v, w, o, rid_only).await,
Self::RecordId(doc_ctx, record_id) => {
Self::process_record(opt, doc_ctx, txn, record_id, rid_only).await
}
Self::GenerateRecordId(doc_ctx, table) => Self::process_yield(doc_ctx, table).await,
Self::Value(doc_ctx, value) => Ok(Self::process_value(doc_ctx, value)),
Self::Defer(doc_ctx, key) => Self::process_defer(doc_ctx, key).await,
Self::Mergeable(doc_ctx, tb, id, o) => {
Self::process_mergeable(doc_ctx, tb, id, o).await
}
Self::KeyVal(doc_ctx, key, val) => Ok(Self::process_key_val(doc_ctx, &key, &val)?),
Self::Count(doc_ctx, c) => Ok(Self::process_count(doc_ctx, c)),
Self::IndexItem(doc_ctx, i) => {
Self::process_index_item(doc_ctx, txn, i, rid_only).await
}
Self::IndexItemKey(doc_ctx, i) => Ok(Self::process_index_item_key(doc_ctx, i)),
}
}
#[instrument(level = "trace", skip_all)]
async fn process_lookup(
mut doc_ctx: DocumentContext,
ctx: &FrozenContext,
opt: &Options,
txn: &Transaction,
kind: LookupKind,
key: Key,
rid_only: bool,
) -> Result<Processable> {
let (ft, fk) = match kind {
LookupKind::Graph(_) => {
let gra = graph::Graph::decode_key(&key)?;
(gra.edge.table, gra.edge.key)
}
LookupKind::Reference => {
let refe = r#ref::Ref::decode_key(&key)?;
(refe.ft.into_owned(), refe.fk.into_owned())
}
};
if ft != doc_ctx.tb()?.name {
let tb = txn
.get_or_add_tb(None, &doc_ctx.ns().name, &doc_ctx.db().name, &ft, opt.version)
.await?;
let parent = NsDbCtx {
ns: Arc::clone(doc_ctx.ns()),
db: Arc::clone(doc_ctx.db()),
};
let mutating = matches!(doc_ctx, DocumentContext::NsDbTbMutCtx(_));
doc_ctx =
DocumentContext::initialise(ctx, &parent, tb, &ft, opt.version, mutating).await?;
}
let record = if rid_only {
Arc::new(Default::default())
} else {
txn.get_record(
doc_ctx.ns().namespace_id,
doc_ctx.db().database_id,
&ft,
&fk,
opt.version,
)
.await?
};
let rid = RecordId {
table: ft,
key: fk,
};
let val = Operable::Value(record);
Ok(Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(rid.into()),
ir: None,
val,
})
}
#[instrument(level = "trace", skip_all)]
async fn process_range_key(doc_ctx: DocumentContext, key: Key) -> Result<Processable> {
let key = record::RecordKey::decode_key(&key)?;
let val = Record::new(Value::Null);
let rid = RecordId {
table: key.tb.into_owned(),
key: key.id,
};
let val = Operable::Value(val.into());
let pro = Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysOnly,
generate: None,
rid: Some(rid.into()),
ir: None,
val,
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
async fn process_table_key(doc_ctx: DocumentContext, key: Key) -> Result<Processable> {
let key = record::RecordKey::decode_key(&key)?;
let rid = RecordId {
table: key.tb.into_owned(),
key: key.id,
};
let pro = Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysOnly,
generate: None,
rid: Some(rid.into()),
ir: None,
val: Operable::Value(Record::new(Value::Null).into_read_only()),
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
async fn process_relatable(
doc_ctx: DocumentContext,
txn: &Transaction,
f: RecordId,
through: RelateThrough,
w: RecordId,
o: Option<Value>,
rid_only: bool,
) -> Result<Processable> {
let pro = match (rid_only, through) {
(true, RelateThrough::Table(v)) => Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysOnly,
generate: Some(v),
rid: None,
ir: None,
val: Operable::Value(Default::default()),
},
(false, RelateThrough::Table(v)) => Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: Some(v),
rid: None,
ir: None,
val: Operable::Relate(Default::default(), f, w, o.map(|v| v.into())),
},
(true, RelateThrough::RecordId(v)) => Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysOnly,
generate: None,
rid: Some(v.into()),
ir: None,
val: Operable::Value(Default::default()),
},
(false, RelateThrough::RecordId(v)) if o.is_some() => Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(v.into()),
ir: None,
val: Operable::Relate(Default::default(), f, w, o.map(|v| v.into())),
},
(false, RelateThrough::RecordId(v)) => {
let val = txn
.get_record(
doc_ctx.ns().namespace_id,
doc_ctx.db().database_id,
&v.table,
&v.key,
None,
)
.await?;
let val = Operable::Relate(val, f, w, o.map(|v| v.into()));
Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(v.into()),
ir: None,
val,
}
}
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
async fn process_record(
opt: &Options,
doc_ctx: DocumentContext,
txn: &Transaction,
record_id: RecordId,
rid_only: bool,
) -> Result<Processable> {
let val = if rid_only {
Record::new(Value::Null).into_read_only()
} else {
txn.get_record(
doc_ctx.ns().namespace_id,
doc_ctx.db().database_id,
&record_id.table,
&record_id.key,
opt.version,
)
.await?
};
let val = Operable::Value(val);
let pro = Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(record_id.into()),
ir: None,
val,
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
async fn process_yield(doc_ctx: DocumentContext, table_name: TableName) -> Result<Processable> {
let pro = Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: Some(table_name),
rid: None,
ir: None,
val: Operable::Value(Default::default()),
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
fn process_value(doc_ctx: NsDbCtx, v: Value) -> Processable {
let rid = match &v {
Value::RecordId(rid) => Some(Arc::new(rid.clone())),
_ => None,
};
Processable {
doc_ctx: DocumentContext::NsDbCtx(doc_ctx),
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid,
ir: None,
val: Operable::Value(Record::new(v).into_read_only()),
}
}
#[instrument(level = "trace", skip_all)]
async fn process_defer(doc_ctx: DocumentContext, v: RecordId) -> Result<Processable> {
let pro = Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(v.into()),
ir: None,
val: Operable::Value(Default::default()),
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
async fn process_mergeable(
doc_ctx: DocumentContext,
tb: TableName,
id: Option<RecordIdKey>,
o: Value,
) -> Result<Processable> {
let pro = if let Some(id) = id {
Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(RecordId::new(tb, id).into()),
ir: None,
val: Operable::Insert(Default::default(), o.into()),
}
} else {
Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysOnly,
generate: Some(tb),
rid: None,
ir: None,
val: Operable::Insert(Default::default(), o.into()),
}
};
Ok(pro)
}
#[instrument(level = "trace", skip_all)]
fn process_key_val(doc_ctx: DocumentContext, key: &Key, val: &[u8]) -> Result<Processable> {
let key = record::RecordKey::decode_key(key)?;
let rid = RecordId {
table: key.tb.into_owned(),
key: key.id,
};
let val = Record::kv_decode_value(val, rid.clone())?;
let val = Operable::Value(val.into());
Ok(Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(rid.into()),
ir: None,
val,
})
}
#[instrument(level = "trace", skip_all)]
fn process_count(doc_ctx: DocumentContext, count: usize) -> Processable {
Processable {
record_strategy: RecordStrategy::Count,
generate: None,
doc_ctx,
rid: None,
ir: None,
val: Operable::Count(count),
}
}
#[instrument(level = "trace", skip_all)]
fn process_index_item_key(doc_ctx: DocumentContext, i: IndexItemRecord) -> Processable {
let (t, v, ir) = i.consume();
Processable {
record_strategy: RecordStrategy::KeysOnly,
generate: None,
doc_ctx,
rid: Some(t),
ir: Some(Arc::new(ir)),
val: Operable::Value(v.unwrap_or_else(|| Record::new(Value::Null).into_read_only())),
}
}
#[instrument(level = "trace", skip_all)]
async fn process_index_item(
doc_ctx: DocumentContext,
txn: &Transaction,
i: IndexItemRecord,
rid_only: bool,
) -> Result<Processable> {
let (t, v, ir) = i.consume();
let v = if let Some(v) = v {
v
} else if rid_only {
Record::new(Value::Null).into_read_only()
} else {
txn.get_record(
doc_ctx.ns().namespace_id,
doc_ctx.db().database_id,
&t.table,
&t.key,
None,
)
.await?
};
let pro = Processable {
doc_ctx,
record_strategy: RecordStrategy::KeysAndValues,
generate: None,
rid: Some(t),
ir: Some(ir.into()),
val: Operable::Value(v),
};
Ok(pro)
}
}
pub(super) struct ConcurrentCollector<'a> {
stk: &'a mut Stk,
ctx: &'a FrozenContext,
opt: &'a Options,
txn: &'a Transaction,
stm: &'a Statement<'a>,
ite: &'a mut Iterator,
}
impl Collector for ConcurrentCollector<'_> {
#[instrument(level = "trace", skip_all)]
async fn collect(&mut self, collectable: Collectable) -> Result<()> {
if self.ite.skippable() > 0 {
self.ite.skipped(1);
return Ok(());
}
let pro = collectable.prepare(self.ctx, self.opt, self.txn, false).await?;
self.ite.process(self.stk, self.ctx, self.opt, self.stm, pro).await?;
Ok(())
}
fn iterator(&mut self) -> &mut Iterator {
self.ite
}
}
pub(super) struct ConcurrentDistinctCollector<'a> {
coll: ConcurrentCollector<'a>,
dis: &'a mut SyncDistinct,
}
impl Collector for ConcurrentDistinctCollector<'_> {
#[instrument(level = "trace", skip_all)]
async fn collect(&mut self, collectable: Collectable) -> Result<()> {
let skippable = self.coll.ite.skippable() > 0;
let pro =
collectable.prepare(self.coll.ctx, self.coll.opt, self.coll.txn, skippable).await?;
if self.dis.check_already_processed(&pro) {
return Ok(());
}
if skippable {
self.coll.ite.skipped(1);
return Ok(());
}
self.coll
.ite
.process(self.coll.stk, self.coll.ctx, self.coll.opt, self.coll.stm, pro)
.await?;
Ok(())
}
fn iterator(&mut self) -> &mut Iterator {
self.coll.ite
}
}
pub(super) trait Collector {
async fn collect(&mut self, collected: Collectable) -> Result<()>;
fn max_fetch_size(&mut self) -> u32 {
if let Some(l) = self.iterator().start_limit() {
*l
} else {
NORMAL_BATCH_SIZE
}
}
fn iterator(&mut self) -> &mut Iterator;
fn check_query_planner_context<'b>(
ctx: &'b FrozenContext,
table: &'b TableName,
) -> Cow<'b, FrozenContext> {
if let Some(qp) = ctx.get_query_planner()
&& let Some(exe) = qp.get_query_executor(table)
{
let mut ctx = Context::new_child(ctx);
ctx.set_query_executor(exe.clone());
return Cow::Owned(ctx.freeze());
}
Cow::Borrowed(ctx)
}
#[instrument(level = "trace", name = "Collector::collect_iterable", skip_all)]
async fn collect_iterable(
&mut self,
ctx: &FrozenContext,
opt: &Options,
iterable: Iterable,
) -> Result<()> {
if ctx.is_done(None).await? {
return Ok(());
}
match iterable {
Iterable::Value(doc_ctx, v) => {
if v.is_nullish() {
return Ok(());
}
return self.collect(Collectable::Value(doc_ctx, v)).await;
}
Iterable::GenerateRecordId(doc_ctx, v) => {
self.collect(Collectable::GenerateRecordId(doc_ctx, v)).await?
}
Iterable::RecordId(doc_ctx, v) => {
self.collect(Collectable::RecordId(doc_ctx, v)).await?
}
Iterable::Defer(doc_ctx, v) => self.collect(Collectable::Defer(doc_ctx, v)).await?,
Iterable::Lookup {
doc_ctx,
kind,
from,
what,
} => self.collect_lookup(ctx, opt, doc_ctx, from, kind, what).await?,
Iterable::Range(doc_ctx, tb, v, rs, sc) => match rs {
RecordStrategy::Count => {
self.collect_range_count(ctx, opt, doc_ctx, &tb, v).await?
}
RecordStrategy::KeysOnly => {
self.collect_range_keys(ctx, opt, doc_ctx, &tb, v, sc).await?
}
RecordStrategy::KeysAndValues => {
self.collect_range(ctx, opt, doc_ctx, &tb, v, sc).await?
}
},
Iterable::Table(doc_ctx, table, rs, sc) => {
let ctx = Self::check_query_planner_context(ctx, &table);
match rs {
RecordStrategy::Count => {
self.collect_table_count(&ctx, opt, doc_ctx, &table).await?
}
RecordStrategy::KeysOnly => {
self.collect_table_keys(&ctx, opt, doc_ctx, &table, sc).await?
}
RecordStrategy::KeysAndValues => {
self.collect_table(&ctx, opt, doc_ctx, &table, sc).await?
}
}
}
Iterable::Index(doc_ctx, v, irf, rs) => {
if let Some(qp) = ctx.get_query_planner()
&& let Some(exe) = qp.get_query_executor(&v)
{
let mut ctx = Context::new_child(ctx);
ctx.set_query_executor(exe.clone());
let ctx = ctx.freeze();
return self.collect_index_items(&ctx, doc_ctx, irf, rs).await;
}
self.collect_index_items(ctx, doc_ctx, irf, rs).await?
}
Iterable::Mergeable(doc_ctx, tb, id, o) => {
self.collect(Collectable::Mergeable(doc_ctx, tb, id, o)).await?
}
Iterable::Relatable(doc_ctx, f, v, w, o) => {
self.collect(Collectable::Relatable {
doc_ctx,
f,
v,
w,
o,
})
.await?
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn start_skip(
&mut self,
ctx: &FrozenContext,
opt: &Options,
mut rng: Range<Key>,
sc: ScanDirection,
) -> Result<Option<Range<Key>>> {
let ite = self.iterator();
let skippable = ite.skippable();
if skippable == 0 {
return Ok(Some(rng));
}
let txn = ctx.tx();
let mut cursor = txn.open_keys_cursor(rng.clone(), sc, 0, opt.version).await?;
let mut skipped = 0;
let mut last_key: Vec<u8> = vec![];
'outer: loop {
let remaining = skippable.saturating_sub(skipped).min(NORMAL_BATCH_SIZE as usize);
if remaining == 0 {
break;
}
let batch = cursor.next_batch(remaining as u32).await?;
if batch.is_empty() {
break;
}
for key in &batch {
if ctx.is_done(Some(skipped)).await? {
break 'outer;
}
last_key.clear();
last_key.extend_from_slice(key);
skipped += 1;
}
}
if last_key.is_empty() {
return Ok(None);
}
ite.skipped(skipped);
match sc {
ScanDirection::Forward => {
last_key.push(0xFF);
rng.start = last_key;
}
ScanDirection::Backward => {
rng.end = last_key;
}
}
Ok(Some(rng))
}
#[instrument(level = "trace", skip_all)]
async fn collect_table(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
table: &TableName,
sc: ScanDirection,
) -> Result<()> {
let ns = doc_ctx.ns().namespace_id;
let db = doc_ctx.db().database_id;
let beg = record::prefix(ns, db, table)?;
let end = record::suffix(ns, db, table)?;
let Some(rng) = self.start_skip(ctx, opt, beg..end, sc).await? else {
return Ok(());
};
let txn = ctx.tx();
let mut cursor = txn.open_vals_cursor(rng, sc, 0, opt.version).await?;
let mut count = 0;
'outer: loop {
let batch = cursor.next_batch(NORMAL_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
let owned: Vec<(Key, Val)> =
batch.iter().map(|(k, v)| (k.to_vec(), v.to_vec())).collect();
for (k, v) in owned {
if ctx.is_done(Some(count)).await? {
break 'outer;
}
self.collect(Collectable::KeyVal(doc_ctx.clone(), k, v)).await?;
count += 1;
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_table_keys(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
table: &TableName,
sc: ScanDirection,
) -> Result<()> {
let ns = doc_ctx.ns().namespace_id;
let db = doc_ctx.db().database_id;
let beg = record::prefix(ns, db, table)?;
let end = record::suffix(ns, db, table)?;
let rng = if let Some(rng) = self.start_skip(ctx, opt, beg..end, sc).await? {
rng
} else {
return Ok(());
};
let txn = ctx.tx();
let mut cursor = txn.open_keys_cursor(rng, sc, 0, opt.version).await?;
let mut count = 0;
'outer: loop {
let batch = cursor.next_batch(NORMAL_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
let owned: Vec<Key> = batch.iter().map(|k| k.to_vec()).collect();
for k in owned {
if ctx.is_done(Some(count)).await? {
break 'outer;
}
self.collect(Collectable::TableKey(doc_ctx.clone(), k)).await?;
count += 1;
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_table_count(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
v: &TableName,
) -> Result<()> {
let ns = doc_ctx.ns().namespace_id;
let db = doc_ctx.db().database_id;
let beg = record::prefix(ns, db, v)?;
let end = record::suffix(ns, db, v)?;
let count = ctx.tx().count(beg..end, opt.version).await?;
self.collect(Collectable::Count(doc_ctx, count)).await?;
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn range_prepare(
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
r: RecordIdKeyRange,
) -> Result<(Vec<u8>, Vec<u8>)> {
let beg = match &r.start {
Bound::Unbounded => record::prefix(ns, db, tb)?,
Bound::Included(v) => record::new(ns, db, tb, v).encode_key()?,
Bound::Excluded(v) => {
let mut key = record::new(ns, db, tb, v).encode_key()?;
key.push(0x00);
key
}
};
let end = match &r.end {
Bound::Unbounded => record::suffix(ns, db, tb)?,
Bound::Excluded(v) => record::new(ns, db, tb, v).encode_key()?,
Bound::Included(v) => {
let mut key = record::new(ns, db, tb, v).encode_key()?;
key.push(0x00);
key
}
};
Ok((beg, end))
}
#[instrument(level = "trace", skip_all)]
async fn collect_range(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
table_name: &TableName,
r: RecordIdKeyRange,
sc: ScanDirection,
) -> Result<()> {
let ns = doc_ctx.ns().namespace_id;
let db = doc_ctx.db().database_id;
let (beg, end) = Self::range_prepare(ns, db, table_name, r).await?;
let rng = if let Some(rng) = self.start_skip(ctx, opt, beg..end, sc).await? {
rng
} else {
return Ok(());
};
let txn = ctx.tx();
let mut cursor = txn.open_vals_cursor(rng, sc, 0, None).await?;
let mut count = 0;
'outer: loop {
let batch = cursor.next_batch(NORMAL_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
let owned: Vec<(Key, Val)> =
batch.iter().map(|(k, v)| (k.to_vec(), v.to_vec())).collect();
for (k, v) in owned {
if ctx.is_done(Some(count)).await? {
break 'outer;
}
self.collect(Collectable::KeyVal(doc_ctx.clone(), k, v)).await?;
count += 1;
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_range_keys(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
tb: &TableName,
r: RecordIdKeyRange,
sc: ScanDirection,
) -> Result<()> {
let ns = doc_ctx.ns().namespace_id;
let db = doc_ctx.db().database_id;
let txn = ctx.tx();
let (beg, end) = Self::range_prepare(ns, db, tb, r).await?;
let rng = if let Some(rng) = self.start_skip(ctx, opt, beg..end, sc).await? {
rng
} else {
return Ok(());
};
let mut cursor = txn.open_keys_cursor(rng, sc, 0, opt.version).await?;
let mut count = 0;
'outer: loop {
let batch = cursor.next_batch(NORMAL_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
let owned: Vec<Key> = batch.iter().map(|k| k.to_vec()).collect();
for k in owned {
if ctx.is_done(Some(count)).await? {
break 'outer;
}
self.collect(Collectable::RangeKey(doc_ctx.clone(), k)).await?;
count += 1;
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_range_count(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
tb: &TableName,
r: RecordIdKeyRange,
) -> Result<()> {
let txn = ctx.tx();
let (beg, end) =
Self::range_prepare(doc_ctx.ns().namespace_id, doc_ctx.db().database_id, tb, r).await?;
let count = txn.count(beg..end, opt.version).await?;
self.collect(Collectable::Count(doc_ctx, count)).await?;
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_lookup(
&mut self,
ctx: &FrozenContext,
opt: &Options,
doc_ctx: DocumentContext,
from: RecordId,
kind: LookupKind,
what: Vec<ComputedLookupSubject>,
) -> Result<()> {
let ns = doc_ctx.ns().namespace_id;
let db = doc_ctx.db().database_id;
let tb = &from.table;
let id = &from.key;
let keys = match (what.is_empty(), &kind) {
(true, LookupKind::Reference) => {
vec![(r#ref::prefix(ns, db, tb, id)?, r#ref::suffix(ns, db, tb, id)?)]
}
(true, LookupKind::Graph(dir)) => match dir {
Dir::Both => {
vec![(graph::prefix(ns, db, tb, id)?, graph::suffix(ns, db, tb, id)?)]
}
Dir::In => vec![(
graph::egprefix(ns, db, tb, id, *dir)?,
graph::egsuffix(ns, db, tb, id, *dir)?,
)],
Dir::Out => vec![(
graph::egprefix(ns, db, tb, id, *dir)?,
graph::egsuffix(ns, db, tb, id, *dir)?,
)],
},
(false, LookupKind::Graph(Dir::Both)) => what
.iter()
.flat_map(|v| {
[
v.presuf(ns, db, tb, id, &LookupKind::Graph(Dir::In)),
v.presuf(ns, db, tb, id, &LookupKind::Graph(Dir::Out)),
]
})
.collect::<Result<Vec<_>>>()?,
(false, kind) => {
what.iter().map(|v| v.presuf(ns, db, tb, id, kind)).collect::<Result<Vec<_>>>()?
}
};
let txn = ctx.tx();
'keys: for (beg, end) in keys {
let mut cursor =
txn.open_keys_cursor(beg..end, ScanDirection::Forward, 0, opt.version).await?;
let mut count = 0;
loop {
let batch = cursor.next_batch(NORMAL_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
let owned: Vec<Key> = batch.iter().map(|k| k.to_vec()).collect();
for key in owned {
if ctx.is_done(Some(count)).await? {
break 'keys;
}
self.collect(Collectable::Lookup(doc_ctx.clone(), kind.clone(), key)).await?;
count += 1;
}
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_index_items(
&mut self,
ctx: &FrozenContext,
doc_ctx: DocumentContext,
irf: IteratorRef,
rs: RecordStrategy,
) -> Result<()> {
let Some(exe) = ctx.get_query_executor() else {
bail!(Error::QueryNotExecuted {
message: "No QueryExecutor has been found.".to_string(),
})
};
let Some(iterator) =
exe.new_iterator(doc_ctx.ns().namespace_id, doc_ctx.db().database_id, irf).await?
else {
bail!(Error::QueryNotExecuted {
message: "No iterator has been found.".to_string(),
})
};
let txn = ctx.tx();
match rs {
RecordStrategy::Count => {
self.collect_index_item_count(ctx, &txn, doc_ctx, iterator).await?
}
RecordStrategy::KeysOnly => {
self.collect_index_item_key(ctx, &txn, doc_ctx, iterator).await?
}
RecordStrategy::KeysAndValues => {
self.collect_index_item_key_value(ctx, &txn, doc_ctx, iterator).await?
}
}
return Ok(());
}
#[instrument(level = "trace", skip_all)]
async fn collect_index_item_key(
&mut self,
ctx: &FrozenContext,
txn: &Transaction,
doc_ctx: DocumentContext,
mut iterator: RecordIterator,
) -> Result<()> {
let fetch_size = self.max_fetch_size();
while !ctx.is_done(None).await? {
let records: Vec<IndexItemRecord> = iterator.next_batch(ctx, txn, fetch_size).await?;
if records.is_empty() {
break;
}
for (count, record) in records.into_iter().enumerate() {
if ctx.is_done(Some(count)).await? {
break;
}
self.collect(Collectable::IndexItemKey(doc_ctx.clone(), record)).await?;
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_index_item_key_value(
&mut self,
ctx: &FrozenContext,
txn: &Transaction,
doc_ctx: DocumentContext,
mut iterator: RecordIterator,
) -> Result<()> {
let fetch_size = self.max_fetch_size();
while !ctx.is_done(None).await? {
let records: Vec<IndexItemRecord> = iterator.next_batch(ctx, txn, fetch_size).await?;
if records.is_empty() {
break;
}
for (count, record) in records.into_iter().enumerate() {
if ctx.is_done(Some(count)).await? {
break;
}
self.collect(Collectable::IndexItem(doc_ctx.clone(), record)).await?;
}
}
Ok(())
}
#[instrument(level = "trace", skip_all)]
async fn collect_index_item_count(
&mut self,
ctx: &FrozenContext,
txn: &Transaction,
doc_ctx: DocumentContext,
mut iterator: RecordIterator,
) -> Result<()> {
let mut total_count = 0;
let fetch_size = self.max_fetch_size();
while !ctx.is_done(None).await? {
let count = iterator.next_count(ctx, txn, fetch_size).await?;
if count == 0 {
break;
}
total_count += count;
}
self.collect(Collectable::Count(doc_ctx, total_count)).await
}
}