use std::sync::Arc;
use futures::StreamExt;
use tracing::instrument;
use super::pipeline::{
build_field_state, determine_scan_direction, eval_limit_expr, kv_scan_stream,
};
use super::{FullTextScan, IndexScan, KnnScan};
use crate::catalog::{DatabaseId, NamespaceId, Permission};
use crate::err::Error;
use crate::exec::index::access_path::{AccessPath, select_access_path};
use crate::exec::index::analysis::IndexAnalyzer;
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_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::idx::planner::ScanDirection;
use crate::key::record;
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,
}
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,
}
}
pub(crate) fn with_pre_decode_filter(mut self, status: PreDecodeFilterStatus) -> Self {
self.pre_decode_filter_status = status;
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()));
}
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 = 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 ctx = ctx.clone();
let stream = async_stream::try_stream! {
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?;
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) {
return;
}
let results = super::record_id::execute_record_lookup(
&rid, version, check_perms, needed_fields.as_ref(), &ctx,
predicate.as_ref(), limit_val, start_val, None,
&pre_decode_filter_status,
).await?;
if !results.is_empty() {
yield ValueBatch { values: results };
}
return;
}
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;
}
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;
}
values.drain(..start_val);
}
if let Some(limit) = limit_val {
values.truncate(limit);
}
if !values.is_empty() {
yield ValueBatch { values };
}
return;
}
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;
}
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;
}
}
yield ValueBatch { values: vec![other] };
return;
}
};
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 = match &table_def {
Some(def) => def.permissions.select.clone(),
None => Permission::None,
};
convert_permission_to_physical_runtime(&catalog_perm, ctx.ctx())
.await
.context("Failed to convert permission")?
} else {
PhysicalPermission::Allow
};
if matches!(select_permission, PhysicalPermission::Deny) {
return;
}
if limit_val == Some(0) {
return;
}
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.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,
},
).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!(crate::err::Error::QueryCancelled),
))?;
}
let mut batch = batch_result?;
let cont = pipeline.process_batch(&mut batch.values, &ctx).await?;
if !batch.values.is_empty() {
yield ValueBatch { values: batch.values };
}
if !cont {
break;
}
}
};
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())))
}
struct TableScanConfig {
ns_id: NamespaceId,
db_id: DatabaseId,
table_name: TableName,
cond: Option<Cond>,
order: Option<Ordering>,
with: Option<With>,
direction: ScanDirection,
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>>,
}
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;
fold_condition_expressions(&mut cond, ctx.function_registry());
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 analyzer = IndexAnalyzer::new(indexes, 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)) =>
{
let operator = IndexScan::new(
index_ref,
access,
direction,
cfg.table_name,
None,
None,
cfg.version,
None,
None,
);
let stream = 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,
}) => {
let residual_cond = resolved_cond.as_ref().and_then(strip_knn_from_condition);
let knn_op = KnnScan::new(
index_ref,
vector,
k,
ef,
cfg.table_name,
cfg.version,
cfg.knn_context.clone(),
residual_cond,
None,
);
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))
}
_ => {
let beg = record::prefix(cfg.ns_id, cfg.db_id, &cfg.table_name)?;
let end = record::suffix(cfg.ns_id, cfg.db_id, &cfg.table_name)?;
let stream = kv_scan_stream(
txn,
beg,
end,
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,
} => {
let residual_cond = resolved_cond.and_then(strip_knn_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,
))
}
AccessPath::EmptyScan => Arc::new(super::EmptyScan::new()),
AccessPath::TableScan
| AccessPath::Union {
..
} => 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 source = expr_to_physical_expr(
crate::expr::Expr::Literal(crate::expr::literal::Literal::String(table_name.into())),
&ctx,
)
.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, ScanDirection::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, ScanDirection::Forward));
}
}