use std::borrow::Cow;
use std::sync::Arc;
use common::future::stream::{self, Yielder};
use futures::StreamExt;
use tracing::instrument;
use super::pipeline::{
build_field_state, determine_scan_direction, eval_limit_expr, kv_scan_stream,
table_read_restricted_fields,
};
use super::{FullTextScan, IndexScan, KnnScan};
use crate::catalog::providers::TableProvider;
use crate::catalog::{DatabaseId, Error, NamespaceId, table_select_permission};
use crate::err::EngineError;
use crate::exec::index::access_path::{
AccessPath, ElementColumns, may_have_array_columns, select_access_path,
with_array_columns_as_elements, without_btree_indexes,
};
use crate::exec::operators::scan::pipeline::ScanPipeline;
use crate::exec::permission::{
PhysicalPermission, convert_permission_to_physical_runtime, should_check_perms,
validate_record_user_access,
};
use crate::exec::planner::util::{
SELECT_ITERATION_PARAMS, fold_condition_expressions, index_covers_ordering,
resolve_condition_params, resolve_projection_field_idioms,
strip_knn_and_matches_from_condition, strip_knn_from_condition,
};
use crate::exec::pre_decode_filter::{PreDecodeFilterStatus, pre_decode_filter_for_execute};
use crate::exec::{
AccessMode, ContextLevel, EvalContext, ExecOperator, ExecutionContext, FlowResult,
OperatorMetrics, PhysicalExpr, ValueBatch, ValueBatchStream, monitor_stream,
};
use crate::expr::order::Ordering;
use crate::expr::with::With;
use crate::expr::{Cond, ControlFlow, ControlFlowExt};
use crate::iam::Action;
use crate::key::schema::RecordPrefix;
use crate::kvs::Direction;
use crate::val::{TableName, Value};
#[derive(Debug, Clone)]
pub struct DynamicScan {
pub(crate) source: Arc<dyn PhysicalExpr>,
pub(crate) version: Option<Arc<dyn PhysicalExpr>>,
pub(crate) cond: Option<Cond>,
pub(crate) order: Option<Ordering>,
pub(crate) with: Option<With>,
pub(crate) needed_fields: Option<std::collections::HashSet<String>>,
pub(crate) predicate: Option<Arc<dyn PhysicalExpr>>,
pub(crate) limit: Option<Arc<dyn PhysicalExpr>>,
pub(crate) start: Option<Arc<dyn PhysicalExpr>>,
pub(crate) metrics: Arc<OperatorMetrics>,
pub(crate) knn_context: Option<Arc<crate::exec::function::KnnContext>>,
pub(crate) pre_decode_filter_status: PreDecodeFilterStatus,
pub(crate) for_update: bool,
}
impl DynamicScan {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
source: Arc<dyn PhysicalExpr>,
version: Option<Arc<dyn PhysicalExpr>>,
cond: Option<Cond>,
order: Option<Ordering>,
with: Option<With>,
needed_fields: Option<std::collections::HashSet<String>>,
predicate: Option<Arc<dyn PhysicalExpr>>,
limit: Option<Arc<dyn PhysicalExpr>>,
start: Option<Arc<dyn PhysicalExpr>>,
) -> Self {
Self {
source,
version,
cond,
order,
with,
needed_fields,
predicate,
limit,
start,
metrics: Arc::new(OperatorMetrics::new()),
knn_context: None,
pre_decode_filter_status: PreDecodeFilterStatus::NotApplicable,
for_update: false,
}
}
pub(crate) fn with_pre_decode_filter(mut self, status: PreDecodeFilterStatus) -> Self {
self.pre_decode_filter_status = status;
self
}
pub(crate) fn with_for_update(mut self, for_update: bool) -> Self {
self.for_update = for_update;
self
}
pub(crate) fn with_knn_context(
mut self,
knn_context: Option<Arc<crate::exec::function::KnnContext>>,
) -> Self {
self.knn_context = knn_context;
self
}
}
impl ExecOperator for DynamicScan {
fn name(&self) -> &'static str {
"DynamicScan"
}
fn attrs(&self) -> Vec<(String, String)> {
let mut attrs = vec![("source".to_string(), self.source.to_sql())];
if let Some(ref pred) = self.predicate {
attrs.push(("predicate".to_string(), pred.to_sql()));
}
if let Some(ref limit) = self.limit {
attrs.push(("limit".to_string(), limit.to_sql()));
}
if let Some(ref start) = self.start {
attrs.push(("offset".to_string(), start.to_sql()));
}
if let Some(s) = self.pre_decode_filter_status.explain_text() {
attrs.push(("pre_decode_filter".to_string(), s.to_string()));
}
if self.for_update {
attrs.push(("for_update".to_string(), "true".to_string()));
}
attrs
}
fn required_context(&self) -> ContextLevel {
let exprs_ctx = [
Some(&self.source),
self.version.as_ref(),
self.predicate.as_ref(),
self.limit.as_ref(),
self.start.as_ref(),
]
.into_iter()
.flatten()
.map(|e| e.required_context())
.max()
.unwrap_or(ContextLevel::Root);
exprs_ctx.max(ContextLevel::Database)
}
fn metrics(&self) -> Option<&OperatorMetrics> {
Some(&self.metrics)
}
fn expressions(&self) -> Vec<(&str, &Arc<dyn PhysicalExpr>)> {
let mut exprs = vec![("source", &self.source)];
if let Some(ref version) = self.version {
exprs.push(("version", version));
}
if let Some(ref pred) = self.predicate {
exprs.push(("predicate", pred));
}
if let Some(ref limit) = self.limit {
exprs.push(("limit", limit));
}
if let Some(ref start) = self.start {
exprs.push(("start", start));
}
exprs
}
fn access_mode(&self) -> AccessMode {
let mut mode = if self.for_update {
AccessMode::ReadWrite
} else {
AccessMode::ReadOnly
};
mode = mode.combine(self.source.access_mode());
if let Some(ref version) = self.version {
mode = mode.combine(version.access_mode());
}
if let Some(ref pred) = self.predicate {
mode = mode.combine(pred.access_mode());
}
if let Some(ref limit) = self.limit {
mode = mode.combine(limit.access_mode());
}
if let Some(ref start) = self.start {
mode = mode.combine(start.access_mode());
}
mode
}
#[instrument(name = "Scan::execute", level = "trace", skip_all)]
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 source_expr = Arc::clone(&self.source);
let version = self.version.clone();
let cond = self.cond.clone();
let order = self.order.clone();
let with = self.with.clone();
let needed_fields = self.needed_fields.clone();
let predicate = self.predicate.clone();
let limit_expr = self.limit.clone();
let start_expr = self.start.clone();
let knn_context = self.knn_context.clone();
let pre_decode_filter_status = self.pre_decode_filter_status.clone();
let for_update = self.for_update;
let ctx = ctx.clone();
let stream = stream::try_async_stream(async move |mut yielder: Yielder<_>| {
let db_ctx = ctx.database().context("Scan requires database context")?;
let ns = Arc::clone(&db_ctx.ns_ctx.ns);
let db = Arc::clone(&db_ctx.db);
let eval_ctx = EvalContext::from_exec_ctx(&ctx);
let table_value = source_expr.evaluate(eval_ctx).await?;
if for_update && !matches!(table_value, Value::RecordId(_)) {
return Err(ControlFlow::Err(anyhow::Error::new(Error::Query {
message: crate::dbs::FOR_UPDATE_TARGETS_ERROR.to_string(),
})));
}
let table_name = match table_value {
Value::Table(t) => t,
Value::RecordId(rid) => {
let version: Option<u64> = match &version {
Some(expr) => {
let eval_ctx = EvalContext::from_exec_ctx(&ctx);
let v = expr.evaluate(eval_ctx).await?;
Some(
v.cast_to::<crate::val::Datetime>()
.map_err(|e| anyhow::anyhow!("{e}"))?
.to_version_stamp(ctx.txn().timestamp_impl().as_ref())?,
)
}
None => ctx.version_stamp(),
};
let limit_val: Option<usize> = match &limit_expr {
Some(expr) => Some(eval_limit_expr(&**expr, &ctx).await?),
None => None,
};
let start_val: usize = match &start_expr {
Some(expr) => eval_limit_expr(&**expr, &ctx).await?,
None => 0,
};
if limit_val == Some(0) && !for_update {
return Ok(());
}
let results = super::record_id::execute_record_lookup(
&rid,
version,
for_update,
check_perms,
needed_fields.as_ref(),
&ctx,
predicate.as_ref(),
limit_val,
start_val,
None,
&pre_decode_filter_status,
)
.await?;
if !results.is_empty() {
yielder.emit(ValueBatch::new(results)).await;
}
return Ok(());
}
Value::Array(arr) => {
let limit_val: Option<usize> = match &limit_expr {
Some(expr) => Some(eval_limit_expr(&**expr, &ctx).await?),
None => None,
};
let start_val: usize = match &start_expr {
Some(expr) => eval_limit_expr(&**expr, &ctx).await?,
None => 0,
};
if limit_val == Some(0) {
return Ok(());
}
let mut values = arr.0;
if let Some(ref pred) = predicate {
let mut write_idx = 0;
for read_idx in 0..values.len() {
let eval_ctx =
EvalContext::from_exec_ctx(&ctx).with_value(&values[read_idx]);
if pred.evaluate(eval_ctx).await?.is_truthy() {
if write_idx != read_idx {
values.swap(write_idx, read_idx);
}
write_idx += 1;
}
}
values.truncate(write_idx);
}
if start_val > 0 {
if start_val >= values.len() {
return Ok(());
}
values.drain(..start_val);
}
if let Some(limit) = limit_val {
values.truncate(limit);
}
if !values.is_empty() {
yielder.emit(ValueBatch::new(values)).await;
}
return Ok(());
}
other => {
let limit_val: Option<usize> = match &limit_expr {
Some(expr) => Some(eval_limit_expr(&**expr, &ctx).await?),
None => None,
};
let start_val: usize = match &start_expr {
Some(expr) => eval_limit_expr(&**expr, &ctx).await?,
None => 0,
};
if limit_val == Some(0) || start_val > 0 {
return Ok(());
}
if let Some(ref pred) = predicate {
let eval_ctx = EvalContext::from_exec_ctx(&ctx).with_value(&other);
if !pred.evaluate(eval_ctx).await?.is_truthy() {
return Ok(());
}
}
yielder.emit(ValueBatch::new(vec![other])).await;
return Ok(());
}
};
let limit_val: Option<usize> = match &limit_expr {
Some(expr) => Some(eval_limit_expr(&**expr, &ctx).await?),
None => None,
};
let start_val: usize = match &start_expr {
Some(expr) => eval_limit_expr(&**expr, &ctx).await?,
None => 0,
};
let version_stamp: Option<u64> = match &version {
Some(expr) => {
let eval_ctx = EvalContext::from_exec_ctx(&ctx);
let v = expr.evaluate(eval_ctx).await?;
Some(
v.cast_to::<crate::val::Datetime>()
.map_err(|e| anyhow::anyhow!("{e}"))?
.to_version_stamp(ctx.txn().timestamp_impl().as_ref())?,
)
}
None => ctx.version_stamp(),
};
let table_def = db_ctx
.get_table_def(&table_name, version_stamp)
.await
.context("Failed to get table")?;
if table_def.is_none() {
Err(ControlFlow::Err(anyhow::Error::new(Error::TbNotFound {
name: table_name.clone(),
})))?;
}
let select_permission = if check_perms {
let catalog_perm = table_select_permission(table_def.as_deref());
convert_permission_to_physical_runtime(catalog_perm, &ctx)
.await
.context("Failed to convert permission")?
} else {
PhysicalPermission::Allow
};
if matches!(select_permission, PhysicalPermission::Deny) {
return Ok(());
}
if limit_val == Some(0) {
return Ok(());
}
let field_state =
build_field_state(&ctx, &table_name, check_perms, needed_fields.as_ref()).await?;
let order = if order_touches_restricted_select_field(
order.as_ref(),
&field_state.field_permissions,
) {
None
} else {
order
};
let needs_row_filtering =
ScanPipeline::compute_needs_row_filtering(&select_permission, predicate.as_ref());
let pre_skip = if !needs_row_filtering {
start_val
} else {
0
};
let effective_storage_limit = if !needs_row_filtering {
limit_val
} else {
None
};
let direction = determine_scan_direction(order.as_ref());
let pre_decode_filter = pre_decode_filter_for_execute(
&pre_decode_filter_status,
&field_state,
check_perms,
ctx.ctx().config.exec.idiom_recursion_limit,
);
let (mut source, applied_pre_skip) = {
resolve_table_scan_stream(
&ctx,
TableScanConfig {
ns_id: ns.namespace_id,
db_id: db.database_id,
table_name,
cond,
order: order.clone(),
with,
direction,
version,
storage_limit: effective_storage_limit,
pre_skip,
has_pushed_limit: effective_storage_limit.is_some(),
limit_hint: limit_val
.map(|l| (l + start_val).min(u32::MAX as usize) as u32),
knn_context: knn_context.clone(),
pre_decode_filter,
field_permissions: Arc::clone(&field_state.field_permissions),
},
)
.await?
};
let mut pipeline = ScanPipeline::new(
select_permission,
predicate,
field_state,
check_perms,
limit_val,
start_val.saturating_sub(applied_pre_skip),
);
while let Some(batch_result) = source.next().await {
if ctx.cancellation().is_cancelled() {
Err(ControlFlow::Err(anyhow::anyhow!(EngineError::QueryCancelled)))?;
}
let mut batch = batch_result?;
let cont = pipeline.process_batch(batch.values_mut(), &ctx).await?;
if !batch.is_empty() {
yielder.emit(batch).await;
}
if !cont {
break;
}
}
Ok(())
});
Ok(monitor_stream(Box::pin(stream), "Scan", &self.metrics))
}
}
fn order_touches_restricted_select_field(
order: Option<&Ordering>,
field_permissions: &[(crate::expr::Idiom, PhysicalPermission)],
) -> bool {
let Some(Ordering::Order(order_list)) = order else {
return false;
};
if field_permissions.is_empty() {
return false;
}
order_list
.iter()
.any(|o| field_permissions.iter().any(|(field, _)| o.value.starts_with(field.0.as_slice())))
}
fn index_touches_restricted_select_field(
cols: &[crate::expr::Idiom],
field_permissions: &[(crate::expr::Idiom, PhysicalPermission)],
) -> bool {
if field_permissions.is_empty() {
return false;
}
cols.iter().any(|col| {
field_permissions
.iter()
.any(|(field, _)| crate::exec::planner::util::paths_overlap(&col.0, &field.0))
})
}
struct TableScanConfig {
ns_id: NamespaceId,
db_id: DatabaseId,
table_name: TableName,
cond: Option<Cond>,
order: Option<Ordering>,
with: Option<With>,
direction: Direction,
version: Option<Arc<dyn PhysicalExpr>>,
storage_limit: Option<usize>,
pre_skip: usize,
has_pushed_limit: bool,
limit_hint: Option<u32>,
knn_context: Option<Arc<crate::exec::function::KnnContext>>,
pre_decode_filter: Option<Arc<crate::exec::pre_decode_filter::PreDecodeFilter>>,
field_permissions: Arc<Vec<(crate::expr::Idiom, PhysicalPermission)>>,
}
async fn resolve_table_scan_stream(
ctx: &ExecutionContext,
cfg: TableScanConfig,
) -> Result<(ValueBatchStream, usize), ControlFlow> {
let txn = ctx.txn();
let version_stamp: Option<u64> = match &cfg.version {
Some(expr) => {
let eval_ctx = EvalContext::from_exec_ctx(ctx);
let v = expr.evaluate(eval_ctx).await?;
Some(
v.cast_to::<crate::val::Datetime>()
.map_err(|e| anyhow::anyhow!("{e}"))?
.to_version_stamp(txn.timestamp_impl().as_ref())?,
)
}
None => ctx.version_stamp(),
};
let resolved_cond = match cfg.cond.as_ref() {
Some(c) => {
let ns_db = Some((cfg.ns_id, cfg.db_id));
let mut cond =
resolve_condition_params(c, &ctx.root().ctx, ns_db, SELECT_ITERATION_PARAMS).await;
let capabilities = ctx.root().ctx.get_capabilities();
let check_perms = match ctx.database() {
Ok(db_ctx) => should_check_perms(db_ctx, Action::View).unwrap_or(true),
Err(_) => true,
};
let restricted_fields = table_read_restricted_fields(
&txn,
cfg.ns_id,
cfg.db_id,
&cfg.table_name,
version_stamp,
check_perms,
)
.await;
fold_condition_expressions(
&mut cond,
ctx.function_registry(),
&capabilities,
restricted_fields.as_ref(),
);
resolve_projection_field_idioms(&mut cond, ctx.function_registry());
Some(cond)
}
None => None,
};
if let Some(c) = resolved_cond.as_ref()
&& matches!(&c.0, crate::expr::Expr::Literal(crate::expr::literal::Literal::Bool(false)))
{
let op = super::EmptyScan::new();
let stream = op.execute(ctx)?;
return Ok((stream, 0));
}
let access_path = if matches!(&cfg.with, Some(With::NoIndex)) {
None
} else {
let db_ctx =
ctx.database().context("DynamicScan index analysis requires database context")?;
let indexes = db_ctx
.get_table_indexes(&cfg.table_name, version_stamp)
.await
.context("Failed to fetch indexes")?;
let indexes: std::sync::Arc<[_]> = if indexes
.iter()
.any(|ix| index_touches_restricted_select_field(&ix.cols, &cfg.field_permissions))
{
indexes
.iter()
.filter(|ix| {
!index_touches_restricted_select_field(&ix.cols, &cfg.field_permissions)
})
.cloned()
.collect::<Vec<_>>()
.into()
} else {
indexes
};
let columns = if may_have_array_columns(&indexes) {
match txn.all_tb_fields(cfg.ns_id, cfg.db_id, &cfg.table_name, version_stamp).await {
Ok(fields) => with_array_columns_as_elements(indexes, &fields),
Err(e) => {
tracing::warn!(
table = %cfg.table_name,
error = %e,
"field list failed in DynamicScan; planning without b-tree indexes",
);
ElementColumns::unchanged(without_btree_indexes(indexes))
}
}
} else {
ElementColumns::unchanged(indexes)
};
let analyzer = columns.analyzer(cfg.with.as_ref());
let candidates = analyzer.analyze(resolved_cond.as_ref(), cfg.order.as_ref());
if candidates.is_empty() {
analyzer
.try_or_union(resolved_cond.as_ref(), cfg.direction)
.or_else(|| analyzer.try_in_expansion(resolved_cond.as_ref(), cfg.direction))
.or_else(|| {
analyzer.try_containment_expansion(resolved_cond.as_ref(), cfg.direction)
})
} else {
let path = select_access_path(candidates, cfg.with.as_ref(), cfg.direction);
if path.is_full_range_scan() {
analyzer.try_or_union(resolved_cond.as_ref(), cfg.direction).or(Some(path))
} else {
Some(path)
}
}
};
match access_path {
Some(AccessPath::BTreeScan {
index_ref,
access,
direction,
}) if !cfg.has_pushed_limit
|| (cfg
.order
.as_ref()
.is_none_or(|o| index_covers_ordering(&index_ref, &access, direction, o))
&& !crate::exec::index::access_path::access_fans_out(&index_ref.cols, &access)) =>
{
let fans_out =
crate::exec::index::access_path::access_fans_out(&index_ref.cols, &access);
let operator: Arc<dyn ExecOperator> = Arc::new(IndexScan::new(
index_ref,
access,
direction,
cfg.table_name,
None,
None,
cfg.version,
None,
None,
));
let stream = if fans_out {
crate::exec::operators::DistinctRecords::new(operator).execute(ctx)?
} else {
operator.execute(ctx)?
};
Ok((stream, 0))
}
Some(AccessPath::FullTextSearch {
index_ref,
query,
operator,
}) => {
let ft_op =
FullTextScan::new(index_ref, query, operator, cfg.table_name, cfg.version, None);
let stream = ft_op.execute(ctx)?;
Ok((stream, 0))
}
Some(AccessPath::KnnSearch {
index_ref,
vector,
k,
ef,
prefilter: _,
}) => {
let residual_cond =
resolved_cond.as_ref().and_then(strip_knn_and_matches_from_condition);
let matches_condition = match resolved_cond.as_ref().and_then(strip_knn_from_condition)
{
Some(c) => {
let mut allowlist = std::collections::HashSet::new();
if let Some(orig) = cfg.cond.as_ref() {
crate::exec::physical_expr::collect_cond_matches(&orig.0, &mut allowlist);
}
crate::exec::physical_expr::collect_cond_matches(&c.0, &mut allowlist);
let mut planner =
crate::exec::planner::Planner::new(ctx.ctx(), ctx.function_registry());
planner.set_matches_scope(Arc::new(crate::exec::physical_expr::MatchesScope {
allowlist,
executor_tables: Arc::from(vec![cfg.table_name.clone()]),
}));
Some(planner.physical_expr(c.0).await.map_err(|e| anyhow::anyhow!("{e}"))?)
}
None => None,
};
let knn_op = KnnScan::new(
index_ref,
vector,
k,
ef,
cfg.table_name,
cfg.version,
cfg.knn_context.clone(),
residual_cond,
None,
None,
)
.with_matches_condition(matches_condition);
let stream = knn_op.execute(ctx)?;
Ok((stream, 0))
}
Some(AccessPath::Union {
paths,
dedupe: _,
}) => {
let mut sub_operators: Vec<Arc<dyn ExecOperator>> = Vec::with_capacity(paths.len());
for path in paths {
sub_operators.push(create_index_operator(&path, &cfg, resolved_cond.as_ref()));
}
let union_op = super::UnionIndexScan::new(cfg.table_name.clone(), sub_operators, None);
let stream = union_op.execute(ctx)?;
Ok((stream, 0))
}
Some(AccessPath::EmptyScan) => {
let op = super::EmptyScan::new();
let stream = op.execute(ctx)?;
Ok((stream, 0))
}
_ => {
{
if txn.get_tb(cfg.ns_id, cfg.db_id, &cfg.table_name, None).await?.is_some_and(
|def| crate::kvs::lightweight::lightweight_relation(&def.table_type).is_some(),
) {
return Err(ControlFlow::Err(anyhow::anyhow!(
"a LIGHTWEIGHT relation cannot be scanned by this execution path; \
reference the table statically so the record-less scan can serve it"
)));
}
}
let range = RecordPrefix {
ns: cfg.ns_id,
db: cfg.db_id,
tb: Cow::Borrowed(&cfg.table_name),
}
.range()?;
let stream = kv_scan_stream(
txn,
range,
version_stamp,
cfg.storage_limit,
cfg.direction,
cfg.pre_skip,
cfg.limit_hint,
cfg.pre_decode_filter.clone(),
None,
);
Ok((stream, cfg.pre_skip))
}
}
}
fn create_index_operator(
path: &AccessPath,
cfg: &TableScanConfig,
resolved_cond: Option<&Cond>,
) -> Arc<dyn ExecOperator> {
match path {
AccessPath::BTreeScan {
index_ref,
access,
direction,
} => Arc::new(IndexScan::new(
index_ref.clone(),
access.clone(),
*direction,
cfg.table_name.clone(),
None,
None,
cfg.version.clone(),
None,
None,
)),
AccessPath::FullTextSearch {
index_ref,
query,
operator,
} => Arc::new(FullTextScan::new(
index_ref.clone(),
query.clone(),
operator.clone(),
cfg.table_name.clone(),
cfg.version.clone(),
None,
)),
AccessPath::KnnSearch {
index_ref,
vector,
k,
ef,
prefilter: _,
} => {
let residual_cond = resolved_cond.and_then(strip_knn_and_matches_from_condition);
Arc::new(KnnScan::new(
index_ref.clone(),
vector.clone(),
*k,
*ef,
cfg.table_name.clone(),
cfg.version.clone(),
cfg.knn_context.clone(),
residual_cond,
None,
None,
))
}
AccessPath::EmptyScan => Arc::new(super::EmptyScan::new()),
AccessPath::TableScan
| AccessPath::Union {
..
}
| AccessPath::BitmapFusion {
..
} => Arc::new(super::TableScan::new(
cfg.table_name.clone(),
cfg.direction,
None,
None,
None,
None,
None,
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ctx::Context;
use crate::exec::planner::expr_to_physical_expr;
async fn create_test_scan(table_name: &str, with_index_hints: bool) -> DynamicScan {
let ctx = std::sync::Arc::new(Context::new_test());
let registry = crate::exec::function::FunctionRegistry::with_builtins();
let source = expr_to_physical_expr(
crate::expr::Expr::Literal(crate::expr::literal::Literal::String(table_name.into())),
&ctx,
®istry,
)
.await
.expect("Failed to create physical expression");
DynamicScan::new(
source,
None,
None,
None,
if with_index_hints {
Some(With::NoIndex)
} else {
None
},
None,
None,
None,
None,
)
}
#[tokio::test]
async fn test_scan_struct_with_index_fields() {
let scan = create_test_scan("test_table", false).await;
assert!(scan.cond.is_none());
assert!(scan.order.is_none());
assert!(scan.with.is_none());
}
#[tokio::test]
async fn test_scan_struct_with_noindex_hint() {
let scan = create_test_scan("test_table", true).await;
assert!(scan.with.is_some());
assert!(matches!(scan.with, Some(With::NoIndex)));
}
#[tokio::test]
async fn test_scan_operator_name() {
let scan = create_test_scan("test_table", false).await;
assert_eq!(scan.name(), "DynamicScan");
}
#[tokio::test]
async fn test_scan_required_context() {
let scan = create_test_scan("test_table", false).await;
assert!(matches!(scan.required_context(), ContextLevel::Database));
}
#[test]
fn test_determine_scan_direction_no_order() {
let direction = determine_scan_direction(None);
assert!(matches!(direction, Direction::Forward));
}
#[test]
fn test_determine_scan_direction_random_order() {
use crate::expr::order::Ordering;
let order = Ordering::Random;
let direction = determine_scan_direction(Some(&order));
assert!(matches!(direction, Direction::Forward));
}
}