use std::collections::HashSet;
use std::sync::Arc;
use common::future::stream::{self, Yielder};
use futures::StreamExt;
use reblessive::TreeStack;
use roaring::RoaringTreemap;
use super::common::{fetch_and_filter_records_measured, next_batch_len, probe_batch_len};
use super::pipeline::{ScanPipeline, build_field_state};
use super::resolved::ResolvedTableContext;
use crate::catalog::{DatabaseId, Error, Index, NamespaceId};
use crate::err::EngineError;
use crate::exec::index::access_path::{BTreeAccess, IndexRef};
use crate::exec::index::iterator::btree::{
INDEX_BATCH_SIZE, bitmap_scan_range, decode_entry_doc_ids,
};
use crate::exec::permission::{
PhysicalPermission, convert_permission_to_physical_runtime, should_check_perms,
tables_have_full_select_permission, validate_record_user_access,
};
use crate::exec::{
AccessMode, BoxFut, ContextLevel, ExecOperator, ExecutionContext, FlowResult, OperatorMetrics,
ValueBatch, ValueBatchStream, monitor_stream,
};
use crate::expr::operator::MatchesOperator;
use crate::expr::{ControlFlow, ControlFlowExt};
use crate::iam::Action;
use crate::idx::IndexKeyBase;
use crate::idx::docids::TableDocIds;
use crate::idx::ft::fulltext::FullTextIndex;
use crate::kvs::util::scan;
use crate::kvs::{CachePolicy, Transaction};
use crate::legacy::analyzer_function::LegacyAnalyzerFunction;
use crate::val::{RecordId, TableName};
const RESOLVE_BATCH_SIZE: usize = 1000;
const RESOLVE_PROBE_ROWS: usize = 32;
const AND_CHILD_BUDGET_FACTOR: u64 = 64;
fn reject_versioned_execution(ctx: &ExecutionContext) -> std::result::Result<(), ControlFlow> {
if ctx.version_stamp().is_some() {
return Err(ControlFlow::Err(anyhow::anyhow!(
"Bitmap index plans do not support VERSION queries"
)));
}
Ok(())
}
#[derive(Debug)]
pub(crate) struct BitmapNode {
kind: BitmapNodeKind,
children_dyn: Vec<Arc<dyn ExecOperator>>,
metrics: Arc<OperatorMetrics>,
}
#[derive(Debug)]
pub(crate) enum BitmapNodeKind {
BTree {
index_ref: IndexRef,
access: BTreeAccess,
},
FullText {
index_ref: IndexRef,
query: String,
operator: MatchesOperator,
},
Graph {
source: crate::val::RecordId,
direction: crate::expr::Dir,
edge_tables: Vec<TableName>,
},
And {
children: Vec<Arc<BitmapNode>>,
},
Or {
children: Vec<Arc<BitmapNode>>,
},
AndNot {
base: Arc<BitmapNode>,
subtract: Arc<BitmapNode>,
},
}
enum BranchBitmap {
Ready(RoaringTreemap),
Overflow,
}
impl BitmapNode {
pub(crate) fn btree(index_ref: IndexRef, access: BTreeAccess) -> Arc<Self> {
Arc::new(Self {
kind: BitmapNodeKind::BTree {
index_ref,
access,
},
children_dyn: Vec::new(),
metrics: Arc::new(OperatorMetrics::new()),
})
}
pub(crate) fn fulltext(
index_ref: IndexRef,
query: String,
operator: MatchesOperator,
) -> Arc<Self> {
Arc::new(Self {
kind: BitmapNodeKind::FullText {
index_ref,
query,
operator,
},
children_dyn: Vec::new(),
metrics: Arc::new(OperatorMetrics::new()),
})
}
pub(crate) fn graph(
source: crate::val::RecordId,
direction: crate::expr::Dir,
edge_tables: Vec<TableName>,
) -> Arc<Self> {
Arc::new(Self {
kind: BitmapNodeKind::Graph {
source,
direction,
edge_tables,
},
children_dyn: Vec::new(),
metrics: Arc::new(OperatorMetrics::new()),
})
}
pub(crate) fn and(children: Vec<Arc<BitmapNode>>) -> Arc<Self> {
let children_dyn =
children.iter().map(|c| Arc::clone(c) as Arc<dyn ExecOperator>).collect();
Arc::new(Self {
kind: BitmapNodeKind::And {
children,
},
children_dyn,
metrics: Arc::new(OperatorMetrics::new()),
})
}
pub(crate) fn or(children: Vec<Arc<BitmapNode>>) -> Arc<Self> {
let children_dyn =
children.iter().map(|c| Arc::clone(c) as Arc<dyn ExecOperator>).collect();
Arc::new(Self {
kind: BitmapNodeKind::Or {
children,
},
children_dyn,
metrics: Arc::new(OperatorMetrics::new()),
})
}
pub(crate) fn and_not(base: Arc<BitmapNode>, subtract: Arc<BitmapNode>) -> Arc<Self> {
let children_dyn = vec![
Arc::clone(&base) as Arc<dyn ExecOperator>,
Arc::clone(&subtract) as Arc<dyn ExecOperator>,
];
Arc::new(Self {
kind: BitmapNodeKind::AndNot {
base,
subtract,
},
children_dyn,
metrics: Arc::new(OperatorMetrics::new()),
})
}
fn build_bitmap<'a>(
&'a self,
bctx: &'a BitmapBuildContext<'a>,
anchored: bool,
) -> BoxFut<'a, std::result::Result<BranchBitmap, ControlFlow>> {
Box::pin(async move {
let result = match &self.kind {
BitmapNodeKind::BTree {
index_ref,
access,
} => self.build_btree_bitmap(bctx, index_ref, access, anchored).await?,
BitmapNodeKind::FullText {
index_ref,
query,
operator,
} => {
BranchBitmap::Ready(
self.build_fulltext_bitmap(bctx, index_ref, query, operator).await?,
)
}
BitmapNodeKind::Graph {
source,
direction,
edge_tables,
} => self.build_graph_bitmap(bctx, source, *direction, edge_tables).await?,
BitmapNodeKind::And {
children,
} => {
let mut acc: Option<RoaringTreemap> = None;
for (i, child) in children.iter().enumerate() {
let child_budget = match &acc {
Some(acc) if !bctx.strict && bctx.budget > 0 => bctx.budget.min(
usize::try_from(acc.len().saturating_mul(AND_CHILD_BUDGET_FACTOR))
.unwrap_or(usize::MAX),
),
_ => bctx.budget,
};
let child_bctx = BitmapBuildContext {
budget: child_budget,
..*bctx
};
match child.build_bitmap(&child_bctx, anchored && i == 0).await? {
BranchBitmap::Ready(bitmap) => {
acc = Some(match acc {
None => bitmap,
Some(mut acc) => {
if bitmap.len() < acc.len() {
let mut bitmap = bitmap;
bitmap &= &acc;
bitmap
} else {
acc &= &bitmap;
acc
}
}
});
if acc.as_ref().is_some_and(|a| a.is_empty()) {
break;
}
}
BranchBitmap::Overflow if bctx.strict => {
return Ok(BranchBitmap::Overflow);
}
BranchBitmap::Overflow => continue,
}
}
match acc {
Some(acc) => BranchBitmap::Ready(acc),
None => BranchBitmap::Overflow,
}
}
BitmapNodeKind::Or {
children,
} => {
let mut acc = RoaringTreemap::new();
let mut overflow = false;
for child in children {
match child.build_bitmap(bctx, anchored).await? {
BranchBitmap::Ready(bitmap) => acc |= bitmap,
BranchBitmap::Overflow => {
overflow = true;
break;
}
}
}
if overflow {
BranchBitmap::Overflow
} else {
BranchBitmap::Ready(acc)
}
}
BitmapNodeKind::AndNot {
base,
subtract,
} => {
match base.build_bitmap(bctx, anchored).await? {
BranchBitmap::Ready(mut acc) => {
if !acc.is_empty() {
match subtract.build_bitmap(bctx, false).await? {
BranchBitmap::Ready(sub) => acc -= sub,
BranchBitmap::Overflow if bctx.strict => {
return Ok(BranchBitmap::Overflow);
}
BranchBitmap::Overflow => {}
}
}
BranchBitmap::Ready(acc)
}
BranchBitmap::Overflow => BranchBitmap::Overflow,
}
}
};
match &result {
BranchBitmap::Ready(bitmap) => self.metrics.add_output_rows(bitmap.len()),
BranchBitmap::Overflow => self.metrics.set_branch_dropped(),
}
Ok(result)
})
}
pub(crate) async fn build_exact_cardinality(
&self,
ctx: &ExecutionContext,
table: &TableName,
) -> std::result::Result<u64, ControlFlow> {
reject_versioned_execution(ctx)?;
let db_ctx = ctx.database().context("Bitmap cardinality requires database context")?;
let ns = db_ctx.ns_ctx.ns.namespace_id;
let db = db_ctx.db.database_id;
let txn = ctx.txn();
let doc_ids = TableDocIds::new(ns, db, table.clone());
let bctx = BitmapBuildContext {
ctx,
txn: txn.as_ref(),
ns,
db,
table,
doc_ids: &doc_ids,
budget: 0,
strict: false,
};
match self.build_bitmap(&bctx, true).await? {
BranchBitmap::Ready(bitmap) => Ok(bitmap.len()),
BranchBitmap::Overflow => Err(ControlFlow::Err(anyhow::anyhow!(
"An exact bitmap plan overflowed its branch budget"
))),
}
}
pub(crate) async fn build_allowlist(
&self,
ctx: &ExecutionContext,
table: &TableName,
budget: usize,
) -> std::result::Result<Option<RoaringTreemap>, ControlFlow> {
reject_versioned_execution(ctx)?;
let db_ctx = ctx.database().context("KNN prefilter requires database context")?;
let ns = db_ctx.ns_ctx.ns.namespace_id;
let db = db_ctx.db.database_id;
let txn = ctx.txn();
let doc_ids = TableDocIds::new(ns, db, table.clone());
let bctx = BitmapBuildContext {
ctx,
txn: txn.as_ref(),
ns,
db,
table,
doc_ids: &doc_ids,
budget,
strict: true,
};
match self.build_bitmap(&bctx, false).await? {
BranchBitmap::Ready(bitmap) => Ok(Some(bitmap)),
BranchBitmap::Overflow => Ok(None),
}
}
async fn build_btree_bitmap(
&self,
bctx: &BitmapBuildContext<'_>,
index_ref: &IndexRef,
access: &BTreeAccess,
anchored: bool,
) -> std::result::Result<BranchBitmap, ControlFlow> {
let ix = index_ref.definition();
let mut range = bitmap_scan_range(bctx.ns, bctx.db, ix, access)
.context("Failed to compute bitmap scan range")?;
let mut docs = RoaringTreemap::new();
let mut missing: Vec<RecordId> = Vec::new();
let mut drained = 0usize;
loop {
if bctx.ctx.cancellation().is_cancelled() {
return Err(ControlFlow::Err(anyhow::anyhow!(EngineError::QueryCancelled)));
}
let res = scan(&mut range, bctx.txn, INDEX_BATCH_SIZE)
.await
.context("Failed to scan index entries")?;
if res.is_empty() {
break;
}
drained += decode_entry_doc_ids(res, &mut docs, &mut missing)
.context("Failed to decode index entry doc-IDs")?;
if !anchored && bctx.budget > 0 && drained > bctx.budget {
return Ok(BranchBitmap::Overflow);
}
}
for rid in missing {
match bctx
.doc_ids
.get_doc_id(bctx.txn, &rid.key)
.await
.context("Failed to resolve a doc-ID for an index entry")?
{
Some(doc_id) => {
docs.insert(doc_id);
}
None => {
return Err(ControlFlow::Err(anyhow::anyhow!(
"Index '{}' on table '{}' contains an entry without a doc-ID mapping; \
run `REBUILD INDEX {} ON {}` to repair it",
ix.name,
ix.table_name,
ix.name,
ix.table_name,
)));
}
}
}
Ok(BranchBitmap::Ready(docs))
}
async fn build_graph_bitmap(
&self,
bctx: &BitmapBuildContext<'_>,
source: &crate::val::RecordId,
direction: crate::expr::Dir,
edge_tables: &[TableName],
) -> std::result::Result<BranchBitmap, ControlFlow> {
use std::ops::Bound;
use super::graph_keys::{EdgeTableSpec, compute_graph_ranges};
if bctx.strict {
return Ok(BranchBitmap::Overflow);
}
let db_ctx = bctx.ctx.database().context("Bitmap graph leaf requires database context")?;
if should_check_perms(db_ctx, Action::View)?
&& !tables_have_full_select_permission(
bctx.txn,
db_ctx.ns_name(),
db_ctx.db_name(),
edge_tables,
)
.await
{
return Ok(BranchBitmap::Overflow);
}
let specs: Vec<EdgeTableSpec> = edge_tables
.iter()
.map(|table| EdgeTableSpec {
table: table.clone(),
range_start: Bound::Unbounded,
range_end: Bound::Unbounded,
})
.collect();
let ranges =
compute_graph_ranges(bctx.ns, bctx.db, source, direction, &specs, bctx.ctx).await?;
let mut docs = RoaringTreemap::new();
let mut batch_targets: Vec<crate::val::RecordIdKey> = Vec::new();
let mut drained = 0usize;
for mut range in ranges {
loop {
if bctx.ctx.cancellation().is_cancelled() {
return Err(ControlFlow::Err(anyhow::anyhow!(EngineError::QueryCancelled)));
}
let keys =
crate::kvs::util::scan_keys(&mut range.range, bctx.txn, INDEX_BATCH_SIZE)
.await
.context("Failed to scan adjacency keys")?;
if keys.is_empty() {
break;
}
drained += keys.len();
if bctx.budget > 0 && drained > bctx.budget {
return Ok(BranchBitmap::Overflow);
}
batch_targets.clear();
for key in &keys {
let decoded = range.decoder.decode(key)?;
match decoded.target {
Some(target) if target.table == *bctx.table => {
batch_targets.push(target.key);
}
Some(_) => {}
None => return Ok(BranchBitmap::Overflow),
}
}
if batch_targets.is_empty() {
continue;
}
batch_targets.sort_unstable();
batch_targets.dedup_by(|a, b| a == b && a.addresses_same_record(b));
let ids = bctx
.doc_ids
.get_doc_ids_batch(bctx.txn, &batch_targets)
.await
.context("Failed to resolve reachable records to doc-IDs")?;
docs.extend(ids.into_iter().flatten());
}
}
Ok(BranchBitmap::Ready(docs))
}
async fn build_fulltext_bitmap(
&self,
bctx: &BitmapBuildContext<'_>,
index_ref: &IndexRef,
query: &str,
operator: &MatchesOperator,
) -> std::result::Result<RoaringTreemap, ControlFlow> {
let index_def = index_ref.definition();
let ft_params = match &index_def.index {
Index::FullText(params) => params,
_ => {
return Err(ControlFlow::Err(anyhow::anyhow!(
"Index '{}' is not a full-text index",
index_def.name
)));
}
};
let root = bctx.ctx.root();
let frozen_ctx = &root.ctx;
let opt =
root.options.as_ref().context("Bitmap full-text scan requires Options context")?;
let ikb = IndexKeyBase::new(bctx.ns, bctx.db, bctx.table.clone(), index_def.index_id);
let fti = FullTextIndex::new(
frozen_ctx.get_index_stores(),
bctx.txn,
ikb,
ft_params,
&frozen_ctx.config.idx.file_allowlist,
index_def.format_version,
)
.await
.context("Failed to open full-text index")?;
let query_terms = {
let az_fn = LegacyAnalyzerFunction::new(frozen_ctx, opt);
let mut stack = TreeStack::new();
stack
.enter(|stk| fti.extract_querying_terms(stk, frozen_ctx, &az_fn, query.to_owned()))
.finish()
.await
.context("Failed to extract query terms")?
};
if query_terms.is_empty() {
return Ok(RoaringTreemap::new());
}
Ok(FullTextIndex::merged_postings(&query_terms, operator.operator).unwrap_or_default())
}
}
impl ExecOperator for BitmapNode {
fn name(&self) -> &'static str {
match &self.kind {
BitmapNodeKind::BTree {
..
} => "BitmapIndexScan",
BitmapNodeKind::FullText {
..
} => "BitmapFullTextScan",
BitmapNodeKind::Graph {
..
} => "BitmapGraphScan",
BitmapNodeKind::And {
..
} => "BitmapAnd",
BitmapNodeKind::Or {
..
} => "BitmapOr",
BitmapNodeKind::AndNot {
..
} => "BitmapAndNot",
}
}
fn attrs(&self) -> Vec<(String, String)> {
match &self.kind {
BitmapNodeKind::BTree {
index_ref,
access,
} => vec![
("index".to_string(), index_ref.name.to_string()),
("access".to_string(), access.describe()),
],
BitmapNodeKind::FullText {
index_ref,
query,
..
} => vec![
("index".to_string(), index_ref.name.to_string()),
("query".to_string(), query.clone()),
],
BitmapNodeKind::Graph {
source,
direction,
edge_tables,
} => vec![
("source".to_string(), surrealdb_types::ToSql::to_sql(source)),
(
"direction".to_string(),
match direction {
crate::expr::Dir::In => "<-",
crate::expr::Dir::Out => "->",
crate::expr::Dir::Both => "<->",
}
.to_string(),
),
(
"edges".to_string(),
edge_tables.iter().map(|t| t.to_string()).collect::<Vec<_>>().join(", "),
),
],
_ => vec![],
}
}
fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
self.children_dyn.iter().collect()
}
fn required_context(&self) -> ContextLevel {
ContextLevel::Database
}
fn access_mode(&self) -> AccessMode {
AccessMode::ReadOnly
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn execute(&self, _ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
Err(ControlFlow::Err(anyhow::anyhow!(
"{} is not directly executable; it is evaluated by its BitmapResolve parent",
self.name()
)))
}
}
struct BitmapBuildContext<'a> {
ctx: &'a ExecutionContext,
txn: &'a Transaction,
ns: NamespaceId,
db: DatabaseId,
table: &'a TableName,
doc_ids: &'a TableDocIds,
budget: usize,
strict: bool,
}
#[derive(Debug)]
pub struct BitmapResolve {
pub table_name: TableName,
root: Arc<BitmapNode>,
root_dyn: Arc<dyn ExecOperator>,
needed_fields: Option<HashSet<String>>,
resolved: Option<ResolvedTableContext>,
overflow_fallback: Option<Arc<dyn ExecOperator>>,
fallback_taken: Arc<std::sync::atomic::AtomicBool>,
metrics: Arc<OperatorMetrics>,
}
impl BitmapResolve {
pub(crate) fn new(
table_name: TableName,
root: Arc<BitmapNode>,
needed_fields: Option<HashSet<String>>,
) -> Self {
let root_dyn = Arc::clone(&root) as Arc<dyn ExecOperator>;
Self {
table_name,
root,
root_dyn,
needed_fields,
resolved: None,
overflow_fallback: None,
fallback_taken: Arc::new(std::sync::atomic::AtomicBool::new(false)),
metrics: Arc::new(OperatorMetrics::new()),
}
}
pub(crate) fn with_resolved(mut self, resolved: ResolvedTableContext) -> Self {
self.resolved = Some(resolved);
self
}
pub(crate) fn with_overflow_fallback(mut self, fallback: Arc<dyn ExecOperator>) -> Self {
self.overflow_fallback = Some(fallback);
self
}
}
impl ExecOperator for BitmapResolve {
fn name(&self) -> &'static str {
"BitmapResolve"
}
fn attrs(&self) -> Vec<(String, String)> {
let mut attrs = vec![("table".to_string(), self.table_name.to_string())];
if self.fallback_taken.load(std::sync::atomic::Ordering::Relaxed) {
attrs.push(("union_tier".to_string(), "fallback".to_string()));
}
attrs
}
fn children(&self) -> Vec<&Arc<dyn ExecOperator>> {
vec![&self.root_dyn]
}
fn required_context(&self) -> ContextLevel {
ContextLevel::Database
}
fn access_mode(&self) -> AccessMode {
AccessMode::ReadOnly
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn execute(&self, ctx: &ExecutionContext) -> FlowResult<ValueBatchStream> {
let db_ctx = ctx.database()?.clone();
validate_record_user_access(&db_ctx)?;
let check_perms = should_check_perms(&db_ctx, Action::View)?;
let table_name = self.table_name.clone();
let root = Arc::clone(&self.root);
let resolved = self.resolved.clone();
let needed_fields = self.needed_fields.clone();
let overflow_fallback = self.overflow_fallback.clone();
let fallback_taken = Arc::clone(&self.fallback_taken);
let ctx = ctx.clone();
let stream = stream::try_async_stream(async move |mut yielder: Yielder<_>| {
reject_versioned_execution(&ctx)?;
let db_ctx = ctx.database().context("BitmapResolve requires database context")?;
let ns = Arc::clone(&db_ctx.ns_ctx.ns);
let db = Arc::clone(&db_ctx.db);
let txn = ctx.txn();
let select_permission = if let Some(ref res) = resolved {
res.select_permission(check_perms)
} else if check_perms {
let table_def =
db_ctx.get_table_def(&table_name, None).await.context("Failed to get table")?;
if let Some(def) = &table_def {
convert_permission_to_physical_runtime(&def.permissions.select, &ctx)
.await
.context("Failed to convert permission")?
} else {
Err(ControlFlow::Err(anyhow::Error::new(Error::TbNotFound {
name: table_name.clone(),
})))?
}
} else {
PhysicalPermission::Allow
};
if matches!(select_permission, PhysicalPermission::Deny) {
return Ok(());
}
let field_state = if let Some(ref res) = resolved {
res.field_state_for_projection(needed_fields.as_ref())
} else {
build_field_state(&ctx, &table_name, check_perms, needed_fields.as_ref()).await?
};
let mut pipeline = ScanPipeline::new(
PhysicalPermission::Allow,
None,
field_state,
check_perms,
None,
0,
);
let doc_ids = TableDocIds::new(ns.namespace_id, db.database_id, table_name.clone());
let bctx = BitmapBuildContext {
ctx: &ctx,
txn: txn.as_ref(),
ns: ns.namespace_id,
db: db.database_id,
table: &table_name,
doc_ids: &doc_ids,
budget: *surrealdb_cnf::BITMAP_BRANCH_BUDGET,
strict: false,
};
let anchored = overflow_fallback.is_none();
let bitmap = match root.build_bitmap(&bctx, anchored).await? {
BranchBitmap::Ready(bitmap) => bitmap,
BranchBitmap::Overflow => {
let Some(fallback) = overflow_fallback else {
Err(ControlFlow::Err(anyhow::anyhow!(
"The anchored bitmap plan root overflowed its branch budget"
)))?;
unreachable!()
};
fallback_taken.store(true, std::sync::atomic::Ordering::Relaxed);
let mut stream = fallback.execute(&ctx)?;
while let Some(batch) = stream.next().await {
yielder.emit(batch?).await;
}
return Ok(());
}
};
let mut iter = bitmap.into_iter();
let batch_bytes = ctx.root().ctx.config.exec.scan_batch_bytes;
let mut chunk_len = probe_batch_len(batch_bytes, RESOLVE_PROBE_ROWS);
loop {
if ctx.cancellation().is_cancelled() {
Err(ControlFlow::Err(anyhow::anyhow!(EngineError::QueryCancelled)))?;
}
let chunk: Vec<u64> = iter.by_ref().take(chunk_len).collect();
if chunk.is_empty() {
break;
}
let keys = doc_ids
.get_record_ids_batch(txn.as_ref(), &chunk)
.await
.context("Failed to resolve doc-IDs to record IDs")?;
let mut rids = Vec::with_capacity(keys.len());
for key in keys.into_iter().flatten() {
rids.push(RecordId {
table: table_name.clone(),
key,
});
}
if rids.is_empty() {
continue;
}
let fetched = fetch_and_filter_records_measured(
&ctx,
&txn,
ns.namespace_id,
db.database_id,
&rids,
&select_permission,
check_perms,
None,
CachePolicy::ReadOnly,
)
.await?;
let mut values = fetched.values;
pipeline.process_batch(&mut values, &ctx).await?;
let held_bytes = fetched.fetched_bytes.max(pipeline.materialised_bytes());
chunk_len = next_batch_len(
held_bytes,
fetched.fetched_rows,
chunk_len,
batch_bytes,
RESOLVE_BATCH_SIZE,
);
if !values.is_empty() {
yielder.emit(ValueBatch::new(values)).await;
}
}
Ok(())
});
Ok(monitor_stream(Box::pin(stream), "BitmapResolve", &self.metrics))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::exec::operators::test_util::TestDb;
use crate::kvs::TransactionType;
use crate::val::Value;
fn anonymous() -> crate::dbs::Session {
crate::dbs::Session::default().with_ns("test").with_db("test")
}
async fn index_ref(ctx: &ExecutionContext, table: &str, name: &str) -> IndexRef {
let indexes = ctx
.database()
.expect("a database context")
.get_table_indexes(&TableName::from(table), None)
.await
.expect("catalog read");
let idx =
indexes.iter().position(|ix| ix.name.as_str() == name).expect("the index is defined");
IndexRef::new(indexes, idx)
}
fn field(row: &Value, name: &str) -> Value {
match row {
Value::Object(obj) => obj.get(name).cloned().unwrap_or(Value::None),
_ => Value::None,
}
}
#[tokio::test]
async fn wide_records_resolve_in_batches_sized_to_the_byte_budget() {
const ROWS: usize = 1400;
let db = TestDb::new(&format!(
"DEFINE TABLE t SCHEMALESS;
DEFINE INDEX n_idx ON t FIELDS n;
FOR $i IN 0..{ROWS} {{ CREATE t SET n = 1, emb = array::repeat(0.5f, 1024) }};"
))
.await;
let ctx = db.exec_ctx().await;
let index_ref = index_ref(&ctx, "t", "n_idx").await;
let root = BitmapNode::btree(index_ref, BTreeAccess::Equality(Value::from(1)));
let scan: Arc<dyn ExecOperator> = Arc::new(BitmapResolve::new("t".into(), root, None));
let mut stream = scan.execute(&ctx).expect("execute should succeed");
let mut batches = Vec::new();
while let Some(batch) = stream.next().await {
batches.push(batch.expect("batch should be Ok").into_values());
}
let ids: HashSet<Value> = batches.iter().flatten().map(|row| field(row, "id")).collect();
assert_eq!(ids.len(), ROWS, "every record once");
assert_eq!(batches.iter().map(Vec::len).sum::<usize>(), ROWS, "no record twice");
assert_eq!(batches[0].len(), RESOLVE_PROBE_ROWS, "the first batch is the probe");
let first_bytes = batches[0].iter().map(super::super::common::approx_value_size).sum();
let sized = next_batch_len(
first_bytes,
batches[0].len(),
RESOLVE_PROBE_ROWS,
crate::exec::config::DEFAULT_SCAN_BATCH_BYTES,
RESOLVE_BATCH_SIZE,
);
assert!(sized < RESOLVE_BATCH_SIZE, "wide rows shrink the batch, got {sized}");
let (last, middle) = batches[1..].split_last().expect("more than one batch");
assert!(middle.iter().all(|b| b.len() == sized), "later batches take the sized length");
assert!(last.len() <= sized, "the final batch holds the remainder");
}
#[cfg(all(feature = "kv-mem", not(target_family = "wasm")))]
#[tokio::test]
async fn the_configured_batch_bytes_size_the_batches() {
const ROWS: usize = 400;
let config = surrealdb_cnf::ConfigMap::empty().with_key_value("scan_batch_bytes", "256KiB");
let db = TestDb::new_with_config(
&format!(
"DEFINE TABLE t SCHEMALESS;
DEFINE INDEX n_idx ON t FIELDS n;
FOR $i IN 0..{ROWS} {{ CREATE t SET n = 1, emb = array::repeat(0.5f, 1024) }};"
),
config,
)
.await;
let ctx = db.exec_ctx().await;
let index_ref = index_ref(&ctx, "t", "n_idx").await;
let root = BitmapNode::btree(index_ref, BTreeAccess::Equality(Value::from(1)));
let scan: Arc<dyn ExecOperator> = Arc::new(BitmapResolve::new("t".into(), root, None));
let mut stream = scan.execute(&ctx).expect("execute should succeed");
let mut lens = Vec::new();
while let Some(batch) = stream.next().await {
lens.push(batch.expect("batch should be Ok").into_values().len());
}
assert_eq!(lens.iter().sum::<usize>(), ROWS, "every record once: {lens:?}");
assert!(lens[0] <= 4, "the probe follows the configured budget: {lens:?}");
assert!(lens[1] <= 8, "later batches follow the configured budget: {lens:?}");
}
#[tokio::test]
async fn wide_computed_fields_size_the_batches() {
const ROWS: usize = 2500;
let db = TestDb::new(&format!(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD wide ON t COMPUTED array::repeat(0.5f, 1024);
DEFINE INDEX n_idx ON t FIELDS n;
FOR $i IN 0..{ROWS} {{ CREATE t SET n = 1 }};"
))
.await;
let ctx = db.exec_ctx().await;
let index_ref = index_ref(&ctx, "t", "n_idx").await;
let root = BitmapNode::btree(index_ref, BTreeAccess::Equality(Value::from(1)));
let scan: Arc<dyn ExecOperator> = Arc::new(BitmapResolve::new("t".into(), root, None));
let mut stream = scan.execute(&ctx).expect("execute should succeed");
let mut lens = Vec::new();
while let Some(batch) = stream.next().await {
let rows = batch.expect("batch should be Ok").into_values();
assert!(rows.iter().all(|r| field(r, "wide") != Value::None), "wide materialised");
lens.push(rows.len());
}
assert_eq!(lens.iter().sum::<usize>(), ROWS, "every record once: {lens:?}");
let largest = lens.iter().copied().max().unwrap_or(0);
assert!(largest <= 2 * RESOLVE_PROBE_ROWS, "largest batch {largest}: {lens:?}");
}
#[tokio::test]
async fn wide_records_after_narrow_ones_stay_in_small_batches() {
const NARROW: usize = 100;
const WIDE: usize = 1400;
let db = TestDb::new(&format!(
"DEFINE TABLE t SCHEMALESS;
DEFINE INDEX n_idx ON t FIELDS n;
FOR $i IN 0..{NARROW} {{ CREATE t SET n = 1 }};
FOR $i IN 0..{WIDE} {{ CREATE t SET n = 1, emb = array::repeat(0.5f, 1024) }};"
))
.await;
let ctx = db.exec_ctx().await;
let index_ref = index_ref(&ctx, "t", "n_idx").await;
let root = BitmapNode::btree(index_ref, BTreeAccess::Equality(Value::from(1)));
let scan: Arc<dyn ExecOperator> = Arc::new(BitmapResolve::new("t".into(), root, None));
let mut stream = scan.execute(&ctx).expect("execute should succeed");
let mut lens = Vec::new();
while let Some(batch) = stream.next().await {
lens.push(batch.expect("batch should be Ok").into_values().len());
}
assert_eq!(lens.iter().sum::<usize>(), NARROW + WIDE, "every record once: {lens:?}");
let largest = lens.iter().copied().max().unwrap_or(0);
assert!(largest <= 4 * RESOLVE_PROBE_ROWS, "largest batch {largest}: {lens:?}");
}
#[tokio::test]
async fn rows_a_permission_denies_still_size_the_batch() {
const ROWS: usize = 1500;
let db = TestDb::new_with_auth(&format!(
"DEFINE TABLE t SCHEMALESS PERMISSIONS FOR select WHERE allow = true;
DEFINE INDEX n_idx ON t FIELDS n;
FOR $i IN 0..{ROWS} {{
IF $i % 10 == 0 {{ CREATE t SET n = 1, allow = true }}
ELSE {{ CREATE t SET n = 1, allow = false, emb = array::repeat(0.5f, 1024) }}
}};"
))
.await;
let ctx = db.exec_ctx_as(&anonymous(), TransactionType::Read).await;
let index_ref = index_ref(&ctx, "t", "n_idx").await;
let root = BitmapNode::btree(index_ref, BTreeAccess::Equality(Value::from(1)));
let scan: Arc<dyn ExecOperator> = Arc::new(BitmapResolve::new("t".into(), root, None));
let mut stream = scan.execute(&ctx).expect("execute should succeed");
let mut batches = Vec::new();
while let Some(batch) = stream.next().await {
batches.push(batch.expect("batch should be Ok").into_values());
}
let rows: usize = batches.iter().map(Vec::len).sum();
assert_eq!(rows, ROWS / 10, "only the permitted records are returned");
assert!(batches.len() > 10, "requests stay small: {} batches", batches.len());
}
}