use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use common::future::stream::{self, Yielder};
use crate::catalog::Permission;
use crate::catalog::providers::TableProvider;
use crate::exec::permission::{
PhysicalPermission, check_permission_for_value, convert_permission_to_physical,
};
use crate::exec::pre_decode_filter::{PreDecodeFilter, PreDecodeFilterOutcome};
use crate::exec::topk_pushdown::TopKThresholdProbe;
use crate::exec::{EvalContext, ExecutionContext, PhysicalExpr, ValueBatch, ValueBatchStream};
use crate::expr::{ControlFlow, ControlFlowExt};
use crate::key::schema::RecordKey;
use crate::key::{KVKeyDecode, KVValue, RawRange};
use crate::kvs::{Direction, Transaction};
use crate::val::{TableName, Value};
type RawComputedField = (
String,
Arc<dyn PhysicalExpr>,
Option<crate::expr::Kind>,
Vec<String>,
Option<crate::iam::AuthLimit>,
);
pub(crate) struct ScanPipeline {
permission: PhysicalPermission,
predicate: Option<Arc<dyn PhysicalExpr>>,
field_state: FieldState,
check_perms: bool,
needs_processing: bool,
limit: Option<usize>,
start: usize,
skipped: usize,
emitted: usize,
materialised_bytes: usize,
}
impl ScanPipeline {
pub(crate) fn compute_needs_processing(
permission: &PhysicalPermission,
field_state: &FieldState,
check_perms: bool,
predicate: Option<&Arc<dyn PhysicalExpr>>,
) -> bool {
!matches!(permission, PhysicalPermission::Allow)
|| !field_state.computed_fields.is_empty()
|| (check_perms && !field_state.field_permissions.is_empty())
|| predicate.is_some()
}
pub(crate) fn compute_needs_row_filtering(
permission: &PhysicalPermission,
predicate: Option<&Arc<dyn PhysicalExpr>>,
) -> bool {
!matches!(permission, PhysicalPermission::Allow) || predicate.is_some()
}
pub(crate) fn new(
permission: PhysicalPermission,
predicate: Option<Arc<dyn PhysicalExpr>>,
field_state: FieldState,
check_perms: bool,
limit: Option<usize>,
start: usize,
) -> Self {
let needs_processing = Self::compute_needs_processing(
&permission,
&field_state,
check_perms,
predicate.as_ref(),
);
Self {
permission,
predicate,
field_state,
check_perms,
needs_processing,
limit,
start,
skipped: 0,
emitted: 0,
materialised_bytes: 0,
}
}
pub(crate) fn materialised_bytes(&self) -> usize {
self.materialised_bytes
}
fn has_limit(&self) -> bool {
self.limit.is_some() || self.start > 0
}
pub(crate) async fn process_batch(
&mut self,
batch: &mut Vec<Value>,
ctx: &ExecutionContext,
) -> Result<bool, ControlFlow> {
self.materialised_bytes = 0;
if self.needs_processing {
self.materialised_bytes = filter_and_process_batch(
batch,
&self.permission,
self.predicate.as_ref(),
ctx,
&self.field_state,
self.check_perms,
)
.await?;
}
if self.has_limit() && !batch.is_empty() {
if self.skipped < self.start {
let remaining_to_skip = self.start - self.skipped;
if batch.len() <= remaining_to_skip {
self.skipped += batch.len();
batch.clear();
return Ok(true);
}
self.skipped = self.start;
batch.drain(..remaining_to_skip);
}
if let Some(limit) = self.limit {
let remaining = limit.saturating_sub(self.emitted);
if batch.len() > remaining {
batch.truncate(remaining);
}
}
self.emitted += batch.len();
}
Ok(self.limit.is_none_or(|l| self.emitted < l))
}
}
pub(crate) fn determine_scan_direction(order: Option<&crate::expr::order::Ordering>) -> Direction {
use crate::expr::order::Ordering as OrderingType;
if let Some(OrderingType::Order(order_list)) = order
&& let Some(first) = order_list.0.first()
&& !first.direction
&& first.value.is_id()
{
Direction::Backward
} else {
Direction::Forward
}
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn kv_scan_stream(
txn: Arc<Transaction>,
range: RawRange,
version: Option<u64>,
storage_limit: Option<usize>,
direction: Direction,
pre_skip: usize,
limit_hint: Option<u32>,
pre_decode_filter: Option<Arc<PreDecodeFilter>>,
topk_probe: Option<Arc<TopKThresholdProbe>>,
) -> ValueBatchStream {
let skip = pre_skip.min(u32::MAX as usize) as u32;
let stream = stream::try_async_stream(async move |mut yielder: Yielder<_>| {
let mut cursor = txn
.open_vals_cursor_raw(range, direction, skip, version)
.await
.context("Failed to open scan cursor")?;
let mut first = true;
let mut yielded: usize = 0;
loop {
let mut batch_size = crate::kvs::NORMAL_BATCH_SIZE;
if first && let Some(h) = limit_hint {
batch_size = batch_size.min(h);
}
if let Some(cap) = storage_limit {
let remaining = cap.saturating_sub(yielded);
let remaining_u32 = remaining.min(u32::MAX as usize) as u32;
batch_size = batch_size.min(remaining_u32);
}
if batch_size == 0 {
break;
}
let mut decoded: Vec<Value> = Vec::with_capacity(batch_size as usize);
let mut decode_err: Option<ControlFlow> = None;
let pdf = pre_decode_filter.as_ref();
let topk_threshold = topk_probe.as_ref().and_then(|p| p.snapshot());
let mut topk_skipped: u64 = 0;
let stats = cursor
.for_each(batch_size, &mut |key, val| {
if let Some(pdf) = pdf
&& pdf.apply(key, val) == PreDecodeFilterOutcome::Reject
{
return Ok(std::ops::ControlFlow::Continue(()));
}
if let (Some(probe), Some(threshold)) =
(topk_probe.as_ref(), topk_threshold.as_deref())
&& probe.rejects(threshold, val)
{
topk_skipped += 1;
return Ok(std::ops::ControlFlow::Continue(()));
}
match decode_record(key, val) {
Ok(v) => {
decoded.push(v);
Ok(std::ops::ControlFlow::Continue(()))
}
Err(cf) => {
decode_err = Some(cf);
Ok(std::ops::ControlFlow::Break(()))
}
}
})
.await
.context("Failed to scan record")?;
if topk_skipped > 0
&& let Some(m) = topk_probe.as_ref().and_then(|p| p.metrics())
{
m.add_skipped_rows(topk_skipped);
}
if let Some(cf) = decode_err {
Err(cf)?;
}
first = false;
yielded += stats.rows as usize;
if stats.rows == 0 {
break;
}
super::common::ensure_below_memory_threshold()?;
if !decoded.is_empty() {
yielder.emit(ValueBatch::new(decoded)).await;
}
}
Ok(())
});
Box::pin(stream)
}
#[inline]
pub(crate) fn decode_record(key: &[u8], val: &[u8]) -> Result<Value, ControlFlow> {
let decoded_key = RecordKey::decode_key(key).context("Failed to decode record key")?;
let rid = crate::val::RecordId {
table: decoded_key.tb.into_owned(),
key: decoded_key.id.into_owned(),
};
let record = crate::catalog::Record::kv_decode_value(val, rid)
.context("Failed to deserialize record")?;
Ok(record.data)
}
macro_rules! check_perm {
($permission:expr, $value:expr, $ctx:expr) => {
match $permission {
PhysicalPermission::Allow => Ok::<bool, ControlFlow>(true),
PhysicalPermission::Deny => Ok(false),
PhysicalPermission::Conditional(expr) => {
if $ctx.root().skip_fetch_perms {
Ok(true)
} else {
let mut eval_ctx = EvalContext::from_exec_ctx($ctx).with_value_and_doc($value);
eval_ctx.skip_fetch_perms = true;
expr.evaluate(eval_ctx).await.map(|v| v.is_truthy()).map_err(|e| {
ControlFlow::Err(anyhow::anyhow!("Failed to check permission: {e}"))
})
}
}
}
};
}
pub(crate) async fn filter_and_process_batch(
batch: &mut Vec<Value>,
permission: &PhysicalPermission,
predicate: Option<&Arc<dyn PhysicalExpr>>,
ctx: &ExecutionContext,
state: &FieldState,
check_perms: bool,
) -> Result<usize, ControlFlow> {
let needs_perm_filter = !matches!(permission, PhysicalPermission::Allow);
if !needs_perm_filter
&& state.computed_fields.is_empty()
&& (!check_perms || state.field_permissions.is_empty())
&& let Some(pred) = predicate
{
let eval_ctx = EvalContext::from_exec_ctx(ctx);
let results = pred.evaluate_batch(eval_ctx, &batch[..]).await?;
let mut write_idx = 0;
for (read_idx, result) in results.into_iter().enumerate() {
if result.is_truthy() {
if write_idx != read_idx {
batch.swap(write_idx, read_idx);
}
write_idx += 1;
}
}
batch.truncate(write_idx);
return Ok(0);
}
let mut write_idx = 0;
for read_idx in 0..batch.len() {
if needs_perm_filter && !check_perm!(permission, &batch[read_idx], ctx)? {
continue;
}
if write_idx != read_idx {
batch.swap(write_idx, read_idx);
}
materialise_fields_with_permissions(ctx, state, &mut batch[write_idx], false, check_perms)
.await?;
if let Some(pred) = predicate {
let eval_ctx = EvalContext::from_exec_ctx(ctx).with_value_and_doc(&batch[write_idx]);
if !pred.evaluate(eval_ctx).await?.is_truthy() {
continue;
}
}
write_idx += 1;
}
let materialised = if state.computed_fields.is_empty() {
0
} else {
super::common::ensure_below_memory_threshold()?;
batch.iter().map(super::common::approx_value_size).sum()
};
batch.truncate(write_idx);
Ok(materialised)
}
pub(crate) async fn eval_limit_expr(
expr: &dyn PhysicalExpr,
ctx: &ExecutionContext,
) -> Result<usize, ControlFlow> {
let eval_ctx = EvalContext::from_exec_ctx(ctx);
let value = expr
.evaluate(eval_ctx)
.await
.map_err(|e| ControlFlow::Err(anyhow::anyhow!("Failed to evaluate LIMIT/START: {e}")))?;
match &value {
Value::Number(n) => {
let i = (*n).to_int();
if i >= 0 {
Ok(i as usize)
} else {
Err(ControlFlow::Err(anyhow::anyhow!(
"LIMIT/START must be a non-negative integer, got {i}"
)))
}
}
Value::None | Value::Null => Ok(0),
_ => Err(ControlFlow::Err(anyhow::anyhow!(
"LIMIT/START must be an integer, got {:?}",
value
))),
}
}
#[derive(Debug, Clone)]
pub(crate) struct FieldState {
pub(crate) computed_fields: Vec<ComputedFieldDef>,
pub(crate) field_permissions: Arc<Vec<(crate::expr::Idiom, PhysicalPermission)>>,
dep_map: Arc<HashMap<String, crate::expr::computed_deps::ComputedDeps>>,
permission_field_deps: Arc<HashSet<String>>,
permission_deps_complete: bool,
}
impl FieldState {
pub(crate) fn empty() -> Self {
Self {
computed_fields: Vec::new(),
field_permissions: Arc::new(Vec::new()),
dep_map: Arc::new(HashMap::new()),
permission_field_deps: Arc::new(HashSet::new()),
permission_deps_complete: true,
}
}
}
#[derive(Debug, Clone)]
pub(crate) struct ComputedFieldDef {
field_name: String,
expr: Arc<dyn PhysicalExpr>,
kind: Option<crate::expr::Kind>,
auth_limit: Option<crate::iam::AuthLimit>,
}
impl ComputedFieldDef {
pub(crate) fn field_name(&self) -> &str {
&self.field_name
}
#[cfg(test)]
pub(crate) fn for_test(field_name: impl Into<String>) -> Self {
Self {
field_name: field_name.into(),
expr: Arc::new(crate::exec::physical_expr::Literal(Value::None)),
kind: None,
auth_limit: None,
}
}
}
fn narrowing_auth_limit(
auth_limit: &crate::catalog::auth::AuthLimit,
) -> Result<Option<crate::iam::AuthLimit>, anyhow::Error> {
if auth_limit == &crate::catalog::auth::AuthLimit::new_no_limit() {
return Ok(None);
}
Ok(Some(crate::iam::AuthLimit::try_from(auth_limit)?))
}
pub(crate) async fn table_read_restricted_fields(
txn: &Transaction,
ns_id: crate::catalog::NamespaceId,
db_id: crate::catalog::DatabaseId,
table_name: &TableName,
version: Option<u64>,
check_perms: bool,
) -> Option<HashSet<String>> {
let Ok(table_def) = txn.get_tb(ns_id, db_id, table_name, version).await else {
return None;
};
if check_perms
&& matches!(
crate::catalog::table_select_permission(table_def.as_deref()),
Permission::Specific(_)
) {
return None;
}
let Ok(field_defs) = txn.all_tb_fields(ns_id, db_id, table_name, version).await else {
return None;
};
if check_perms
&& field_defs.iter().any(|fd| matches!(fd.select_permission, Permission::Specific(_)))
{
return None;
}
let mut restricted = HashSet::new();
for fd in field_defs.iter() {
if fd.computed.is_none() {
continue;
}
match fd.name.0.first() {
Some(crate::expr::part::Part::Field(name)) => {
restricted.insert(name.to_string());
}
_ => return None,
}
}
Some(restricted)
}
pub(crate) async fn build_field_state_raw(
planner: &crate::exec::planner::Planner<'_>,
ns_id: crate::catalog::NamespaceId,
db_id: crate::catalog::DatabaseId,
table_name: &TableName,
check_perms: bool,
version: Option<u64>,
) -> Result<FieldState, ControlFlow> {
let txn =
planner.txn().context("build_field_state_raw requires a planner with a transaction")?;
let field_defs = txn
.all_tb_fields(ns_id, db_id, table_name, version)
.await
.context("Failed to get field definitions")?;
let has_computed = field_defs.iter().any(|fd| fd.computed.is_some());
let has_field_perms = check_perms
&& field_defs.iter().any(|fd| !matches!(fd.select_permission, Permission::Full));
if !has_computed && !has_field_perms {
return Ok(FieldState::empty());
}
let mut raw_computed: Vec<RawComputedField> = Vec::new();
let mut dep_map: HashMap<String, crate::expr::computed_deps::ComputedDeps> = HashMap::new();
for fd in field_defs.iter() {
if let Some(ref expr) = fd.computed {
let field_name = fd.name.to_raw_string();
let deps = crate::expr::computed_deps::extract_computed_deps(expr);
dep_map.insert(field_name.clone(), deps.clone());
let physical_expr = planner.physical_expr(expr.clone()).await.with_context(|| {
format!("Computed field '{field_name}' has unsupported expression")
})?;
raw_computed.push((
field_name,
physical_expr,
fd.field_kind.clone(),
deps.fields,
narrowing_auth_limit(&fd.auth_limit)?,
));
}
}
let topo_input: Vec<(String, Vec<String>)> =
raw_computed.iter().map(|(name, _, _, deps, _)| (name.clone(), deps.clone())).collect();
let sorted_indices = crate::expr::computed_deps::topological_sort_computed_fields(&topo_input);
let mut computed_fields = Vec::with_capacity(sorted_indices.len());
for idx in sorted_indices {
let (field_name, expr, kind, _, auth_limit) = &raw_computed[idx];
computed_fields.push(ComputedFieldDef {
field_name: field_name.clone(),
expr: Arc::clone(expr),
kind: kind.clone(),
auth_limit: auth_limit.clone(),
});
}
let mut field_permissions: Vec<(crate::expr::Idiom, PhysicalPermission)> = Vec::new();
let mut permission_field_deps: HashSet<String> = HashSet::new();
let mut permission_deps_complete = true;
if check_perms {
for fd in field_defs.iter() {
if matches!(fd.select_permission, Permission::Full) {
continue;
}
if let Permission::Specific(ref expr) = fd.select_permission {
let deps = crate::expr::computed_deps::extract_computed_deps(expr);
if !deps.is_complete {
if permission_deps_complete {
crate::expr::computed_deps::warn_incomplete_perm_deps(
table_name.as_str(),
fd.name.to_raw_string().as_str(),
);
}
permission_deps_complete = false;
}
permission_field_deps.extend(deps.fields);
}
let physical_perm = convert_permission_to_physical(&fd.select_permission, planner)
.await
.context("Failed to convert field permission")?;
field_permissions.push((fd.name.clone(), physical_perm));
}
}
Ok(FieldState {
computed_fields,
field_permissions: Arc::new(field_permissions),
dep_map: Arc::new(dep_map),
permission_field_deps: Arc::new(permission_field_deps),
permission_deps_complete,
})
}
pub(crate) async fn build_field_state(
ctx: &ExecutionContext,
table_name: &TableName,
check_perms: bool,
needed_fields: Option<&std::collections::HashSet<String>>,
) -> Result<FieldState, ControlFlow> {
let db_ctx = ctx.database().context("build_field_state requires database context")?;
let version = ctx.version_stamp();
let cache_key = (table_name.clone(), check_perms);
if version.is_none() {
let cache = db_ctx.field_state_cache.read().await;
if let Some(cached) = cache.get(&cache_key) {
return Ok(filter_field_state_for_projection(cached, needed_fields));
}
}
let planner = crate::exec::planner::Planner::for_database(ctx.ctx(), ctx.txn(), db_ctx);
let full_state = build_field_state_raw(
&planner,
db_ctx.ns_ctx.ns.namespace_id,
db_ctx.db.database_id,
table_name,
check_perms,
version,
)
.await?;
let cached = Arc::new(full_state);
if version.is_none() {
db_ctx.field_state_cache.write().await.insert(cache_key, Arc::clone(&cached));
}
Ok(filter_field_state_for_projection(&cached, needed_fields))
}
pub(crate) fn filter_field_state_for_projection(
full_state: &FieldState,
needed_fields: Option<&std::collections::HashSet<String>>,
) -> FieldState {
let Some(needed) = needed_fields else {
return full_state.clone();
};
if !full_state.permission_deps_complete {
return full_state.clone();
}
let mut needed_with_perms: std::collections::HashSet<String> = needed.clone();
needed_with_perms.extend(full_state.permission_field_deps.iter().cloned());
let required = crate::expr::computed_deps::resolve_required_computed_fields(
&needed_with_perms,
&full_state.dep_map,
);
let computed_fields = if let Some(ref required_set) = required {
full_state
.computed_fields
.iter()
.filter(|cf| required_set.contains(&cf.field_name))
.cloned()
.collect()
} else {
full_state.computed_fields.clone()
};
FieldState {
computed_fields,
field_permissions: Arc::clone(&full_state.field_permissions),
dep_map: Arc::clone(&full_state.dep_map),
permission_field_deps: Arc::clone(&full_state.permission_field_deps),
permission_deps_complete: full_state.permission_deps_complete,
}
}
pub(crate) async fn compute_fields_for_value(
ctx: &ExecutionContext,
state: &FieldState,
value: &mut Value,
skip_fetch_perms: bool,
) -> Result<(), ControlFlow> {
let view = shadowed_stored_view(state, value);
compute_fields_for_value_against(ctx, state, value, skip_fetch_perms, view).await
}
fn shadowed_stored_view(state: &FieldState, value: &Value) -> Option<Value> {
let shadowed = match value {
Value::Object(obj) => state.dep_map.keys().any(|name| obj.contains_key(name.as_str())),
_ => false,
};
shadowed.then(|| stored_fields_view(state, value))
}
pub(crate) async fn compute_fields_for_value_against(
ctx: &ExecutionContext,
state: &FieldState,
value: &mut Value,
skip_fetch_perms: bool,
stored_view: Option<Value>,
) -> Result<(), ControlFlow> {
if state.computed_fields.is_empty() {
return Ok(());
}
let mut eval_ctx = EvalContext::from_exec_ctx(ctx);
eval_ctx.skip_fetch_perms = skip_fetch_perms;
eval_ctx.computing_record = match &*value {
Value::Object(obj) => match obj.get("id") {
Some(Value::RecordId(rid)) => {
Some(Arc::new(crate::exec::physical_expr::ComputingRecord {
rid: rid.clone(),
stored: stored_view,
}))
}
_ => None,
},
_ => None,
};
let field_ctx = ctx.clone().computing_field();
for cf in &state.computed_fields {
let limited_ctx = cf.auth_limit.as_ref().map(|limit| field_ctx.with_limited_auth(limit));
let narrowed_eval_ctx = limited_ctx.as_ref().unwrap_or(&field_ctx);
let narrowed_eval_ctx = {
let mut narrowed = EvalContext::from_exec_ctx(narrowed_eval_ctx);
narrowed.skip_fetch_perms = skip_fetch_perms;
narrowed.computing_record.clone_from(&eval_ctx.computing_record);
narrowed
};
let row_ctx = narrowed_eval_ctx.with_value_and_doc(value);
let computed_value = match cf.expr.evaluate(row_ctx).await {
Ok(v) => v,
Err(ControlFlow::Return(v)) => v,
Err(e) => return Err(e),
};
let final_value = if let Some(kind) = &cf.kind {
computed_value
.coerce_to_kind(kind)
.with_context(|| format!("Failed to coerce computed field '{}'", cf.field_name))?
} else {
computed_value
};
if let Value::Object(obj) = value {
obj.insert(cf.field_name.clone(), final_value);
} else {
return Err(ControlFlow::Err(anyhow::anyhow!("Value is not an object: {:?}", value)));
}
}
Ok(())
}
fn stored_fields_view(state: &FieldState, value: &Value) -> Value {
let mut view = value.clone();
if let Value::Object(obj) = &mut view {
for name in state.dep_map.keys() {
obj.remove(name.as_str());
}
}
view
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum PermissionPass {
All,
Stored,
Computed,
}
impl PermissionPass {
fn applies_to(self, state: &FieldState, idiom: &crate::expr::Idiom) -> bool {
let computed = match idiom.0.as_slice() {
[crate::expr::Part::Field(name)] => state.dep_map.contains_key(name.as_str()),
_ => state.dep_map.contains_key(&idiom.to_raw_string()),
};
match self {
PermissionPass::All => true,
PermissionPass::Stored => !computed,
PermissionPass::Computed => computed,
}
}
}
pub(crate) async fn materialise_fields_with_permissions(
ctx: &ExecutionContext,
state: &FieldState,
value: &mut Value,
skip_fetch_perms: bool,
check_perms: bool,
) -> Result<(), ControlFlow> {
compute_fields_for_value(ctx, state, value, skip_fetch_perms).await?;
if !check_perms || state.field_permissions.is_empty() {
return Ok(());
}
if state.computed_fields.is_empty() {
return filter_fields_by_permission(ctx, state, value).await;
}
let record = value.clone();
super::common::ensure_below_memory_threshold()?;
let cut = filter_fields_by_permission_inner(
ctx,
state,
value,
PermissionPass::Stored,
Some(&record),
Some(&record),
)
.await?;
if cut {
let stored_view = stored_fields_view(state, value);
compute_fields_for_value_against(ctx, state, value, skip_fetch_perms, Some(stored_view))
.await?;
super::common::ensure_below_memory_threshold()?;
}
let snapshot = (!cut).then_some(&record);
filter_fields_by_permission_inner(
ctx,
state,
value,
PermissionPass::Computed,
snapshot,
Some(&record),
)
.await?;
Ok(())
}
pub(crate) async fn filter_fields_by_permission(
ctx: &ExecutionContext,
state: &FieldState,
value: &mut Value,
) -> Result<(), ControlFlow> {
filter_fields_by_permission_inner(ctx, state, value, PermissionPass::All, None, None).await?;
Ok(())
}
async fn filter_fields_by_permission_inner(
ctx: &ExecutionContext,
state: &FieldState,
value: &mut Value,
pass: PermissionPass,
snapshot: Option<&Value>,
predicate_doc: Option<&Value>,
) -> Result<bool, ControlFlow> {
if state.field_permissions.is_empty() {
return Ok(false);
}
if !matches!(value, Value::Object(_)) {
return Ok(false);
}
let mut lazy_snapshot: Option<Value> = None;
let mut cut_any = false;
for (idiom, perm) in state.field_permissions.iter() {
if !pass.applies_to(state, idiom) {
continue;
}
match perm {
PhysicalPermission::Allow => continue,
PhysicalPermission::Deny => {
let original = match snapshot {
Some(snapshot) => snapshot,
None => &*lazy_snapshot.get_or_insert_with(|| value.clone()),
};
for path in original.each(&idiom.0).into_iter().rev() {
value.cut(&path.0);
cut_any = true;
}
}
PhysicalPermission::Conditional(_) => {
let original = match snapshot {
Some(snapshot) => snapshot,
None => &*lazy_snapshot.get_or_insert_with(|| value.clone()),
};
let doc = predicate_doc.unwrap_or(original);
for path in original.each(&idiom.0).into_iter().rev() {
let field_value = original.pick(&path.0);
let allowed = check_permission_for_value(perm, doc, Some(&field_value), ctx)
.await
.map_err(|e| {
ControlFlow::Err(anyhow::anyhow!(
"Failed to check field permission: {e}"
))
})?;
if !allowed {
value.cut(&path.0);
cut_any = true;
}
}
}
}
}
Ok(cut_any)
}
#[cfg(test)]
mod tests {
use std::borrow::Cow;
use super::*;
use crate::exec::operators::test_util::{TestDb, parse_idiom, physical_expr, root_ctx, val};
use crate::expr::computed_deps::ComputedDeps;
use crate::expr::order::{Order, OrderList, Ordering};
use crate::key::schema::RecordPrefix;
use crate::val::Number;
fn state(
computed_fields: Vec<ComputedFieldDef>,
field_permissions: Vec<(crate::expr::Idiom, PhysicalPermission)>,
) -> FieldState {
FieldState {
computed_fields,
field_permissions: Arc::new(field_permissions),
dep_map: Arc::new(HashMap::new()),
permission_field_deps: Arc::new(HashSet::new()),
permission_deps_complete: true,
}
}
fn computed(
field_name: &str,
expr: Arc<dyn PhysicalExpr>,
kind: Option<crate::expr::Kind>,
) -> ComputedFieldDef {
ComputedFieldDef {
field_name: field_name.to_owned(),
expr,
kind,
auth_limit: None,
}
}
async fn rows(srcs: &[&str]) -> Vec<Value> {
let mut out = Vec::with_capacity(srcs.len());
for src in srcs {
out.push(val(src).await);
}
out
}
fn pick(value: &Value, path: &str) -> Value {
value.pick(&parse_idiom(path).0)
}
fn ns(batch: &[Value]) -> Vec<Value> {
batch.iter().map(|v| pick(v, "n")).collect()
}
fn order_by(field: &str, ascending: bool) -> Ordering {
Ordering::Order(OrderList(vec![Order {
value: parse_idiom(field),
direction: ascending,
..Default::default()
}]))
}
#[tokio::test]
async fn an_unrestricted_scan_needs_neither_processing_nor_row_filtering() {
let empty = FieldState::empty();
assert!(!ScanPipeline::compute_needs_processing(
&PhysicalPermission::Allow,
&empty,
true,
None
));
assert!(!ScanPipeline::compute_needs_row_filtering(&PhysicalPermission::Allow, None));
}
#[tokio::test]
async fn a_table_permission_needs_processing_and_removes_rows() {
let ctx = root_ctx();
let empty = FieldState::empty();
for permission in [
PhysicalPermission::Deny,
PhysicalPermission::Conditional(physical_expr("public = true", &ctx).await),
] {
assert!(ScanPipeline::compute_needs_processing(&permission, &empty, true, None));
assert!(ScanPipeline::compute_needs_row_filtering(&permission, None));
}
}
#[tokio::test]
async fn computed_fields_need_processing_but_are_not_row_filtering() {
let with_computed = state(vec![ComputedFieldDef::for_test("total")], Vec::new());
assert!(ScanPipeline::compute_needs_processing(
&PhysicalPermission::Allow,
&with_computed,
false,
None
));
assert!(!ScanPipeline::compute_needs_row_filtering(&PhysicalPermission::Allow, None));
}
#[tokio::test]
async fn field_permissions_need_processing_only_when_perms_are_checked() {
let with_field_perms =
state(Vec::new(), vec![(parse_idiom("secret"), PhysicalPermission::Deny)]);
assert!(ScanPipeline::compute_needs_processing(
&PhysicalPermission::Allow,
&with_field_perms,
true,
None
));
assert!(!ScanPipeline::compute_needs_processing(
&PhysicalPermission::Allow,
&with_field_perms,
false,
None
));
assert!(!ScanPipeline::compute_needs_row_filtering(&PhysicalPermission::Allow, None));
}
#[tokio::test]
async fn a_predicate_needs_processing_and_removes_rows() {
let ctx = root_ctx();
let predicate = physical_expr("n > 2", &ctx).await;
let empty = FieldState::empty();
assert!(ScanPipeline::compute_needs_processing(
&PhysicalPermission::Allow,
&empty,
false,
Some(&predicate)
));
assert!(ScanPipeline::compute_needs_row_filtering(
&PhysicalPermission::Allow,
Some(&predicate)
));
}
fn limit_pipeline(limit: Option<usize>, start: usize) -> ScanPipeline {
ScanPipeline::new(PhysicalPermission::Allow, None, FieldState::empty(), false, limit, start)
}
async fn numbered(from: usize, count: usize) -> Vec<Value> {
let srcs: Vec<String> = (from..from + count).map(|n| format!("{{ n: {n} }}")).collect();
rows(&srcs.iter().map(String::as_str).collect::<Vec<_>>()).await
}
#[tokio::test]
async fn a_batch_wholly_inside_the_start_offset_is_discarded_and_iteration_continues() {
let ctx = root_ctx();
let mut pipeline = limit_pipeline(None, 5);
let mut batch = numbered(0, 4).await;
assert!(pipeline.process_batch(&mut batch, &ctx).await.unwrap());
assert!(batch.is_empty());
assert_eq!(pipeline.skipped, 4);
assert_eq!(pipeline.emitted, 0);
}
#[tokio::test]
async fn the_start_offset_is_consumed_across_batches_and_skips_only_its_prefix() {
let ctx = root_ctx();
let mut pipeline = limit_pipeline(None, 5);
let mut first = numbered(0, 4).await;
assert!(pipeline.process_batch(&mut first, &ctx).await.unwrap());
assert!(first.is_empty());
let mut second = numbered(4, 4).await;
assert!(pipeline.process_batch(&mut second, &ctx).await.unwrap());
assert_eq!(ns(&second), vec![Value::from(5), Value::from(6), Value::from(7)]);
assert_eq!(pipeline.skipped, 5);
assert_eq!(pipeline.emitted, 3);
}
#[tokio::test]
async fn the_batch_is_truncated_at_the_limit_and_iteration_stops_there() {
let ctx = root_ctx();
let mut pipeline = limit_pipeline(Some(2), 0);
let mut batch = numbered(0, 5).await;
assert!(!pipeline.process_batch(&mut batch, &ctx).await.unwrap());
assert_eq!(ns(&batch), vec![Value::from(0), Value::from(1)]);
assert_eq!(pipeline.emitted, 2);
}
#[tokio::test]
async fn the_limit_is_tracked_across_batches_and_only_goes_false_once_reached() {
let ctx = root_ctx();
let mut pipeline = limit_pipeline(Some(3), 0);
let mut first = numbered(0, 2).await;
assert!(pipeline.process_batch(&mut first, &ctx).await.unwrap());
assert_eq!(first.len(), 2);
let mut second = numbered(2, 2).await;
assert!(!pipeline.process_batch(&mut second, &ctx).await.unwrap());
assert_eq!(ns(&second), vec![Value::from(2)]);
assert_eq!(pipeline.emitted, 3);
}
#[tokio::test]
async fn a_zero_limit_emits_nothing_and_stops_immediately() {
let ctx = root_ctx();
let mut pipeline = limit_pipeline(Some(0), 0);
let mut batch = numbered(0, 3).await;
assert!(!pipeline.process_batch(&mut batch, &ctx).await.unwrap());
assert!(batch.is_empty());
}
#[tokio::test]
async fn without_limit_or_start_the_batch_passes_through_untouched() {
let ctx = root_ctx();
let mut pipeline = limit_pipeline(None, 0);
let mut batch = numbered(0, 3).await;
assert!(pipeline.process_batch(&mut batch, &ctx).await.unwrap());
assert_eq!(ns(&batch), vec![Value::from(0), Value::from(1), Value::from(2)]);
assert_eq!(pipeline.emitted, 0);
}
#[tokio::test]
async fn a_row_denied_by_the_table_permission_never_reaches_the_predicate() {
let ctx = root_ctx();
let empty = FieldState::empty();
let poison = physical_expr("THROW 'the predicate must not see this row'", &ctx).await;
let mut batch = rows(&["{ n: 1 }", "{ n: 2 }"]).await;
filter_and_process_batch(
&mut batch,
&PhysicalPermission::Deny,
Some(&poison),
&ctx,
&empty,
false,
)
.await
.expect("no row reaches the predicate, so it cannot fail");
assert!(batch.is_empty());
let permission = PhysicalPermission::Conditional(physical_expr("public", &ctx).await);
let guarded = physical_expr(
"IF public { true } ELSE { THROW 'the predicate saw a denied row' }",
&ctx,
)
.await;
let mut batch =
rows(&["{ n: 1, public: false }", "{ n: 2, public: true }", "{ n: 3, public: false }"])
.await;
filter_and_process_batch(&mut batch, &permission, Some(&guarded), &ctx, &empty, false)
.await
.expect("denied rows are cut before the predicate");
assert_eq!(ns(&batch), vec![Value::from(2)]);
}
#[tokio::test]
async fn a_computed_field_is_visible_to_the_where_predicate() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD score ON t TYPE int;
DEFINE FIELD doubled ON t TYPE int COMPUTED score * 2;",
)
.await;
let ctx = db.exec_ctx().await;
let field_state =
build_field_state(&ctx, &TableName::from("t"), false, None).await.unwrap();
assert_eq!(field_state.computed_fields.len(), 1);
let predicate = physical_expr("doubled > 4", &ctx).await;
let mut batch = rows(&["{ id: t:1, score: 3 }", "{ id: t:2, score: 1 }"]).await;
filter_and_process_batch(
&mut batch,
&PhysicalPermission::Allow,
Some(&predicate),
&ctx,
&field_state,
false,
)
.await
.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(pick(&batch[0], "id"), val("t:1").await);
assert_eq!(pick(&batch[0], "doubled"), Value::from(6));
}
#[tokio::test]
async fn the_materialised_size_counts_rows_the_predicate_rejects() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD wide ON t COMPUTED array::repeat(0.5f, 1024);",
)
.await;
let ctx = db.exec_ctx().await;
let field_state =
build_field_state(&ctx, &TableName::from("t"), false, None).await.unwrap();
let predicate = physical_expr("false", &ctx).await;
let mut pipeline = ScanPipeline::new(
PhysicalPermission::Allow,
Some(predicate),
field_state,
false,
None,
0,
);
let mut batch = rows(&["{ id: t:1 }", "{ id: t:2 }"]).await;
pipeline.process_batch(&mut batch, &ctx).await.unwrap();
assert!(batch.is_empty(), "the predicate rejects every row");
let wide = std::mem::size_of::<Value>() * 1024;
assert!(pipeline.materialised_bytes() >= 2 * wide, "{}", pipeline.materialised_bytes());
}
#[tokio::test]
async fn a_field_cut_by_a_field_permission_is_invisible_to_the_where_predicate() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD secret ON t TYPE string PERMISSIONS FOR select NONE;",
)
.await;
let ctx = db.exec_ctx().await;
let table = TableName::from("t");
let checked = build_field_state(&ctx, &table, true, None).await.unwrap();
assert_eq!(checked.field_permissions.len(), 1);
let predicate = physical_expr("secret = 'x'", &ctx).await;
let mut batch = rows(&["{ id: t:1, secret: 'x' }"]).await;
filter_and_process_batch(
&mut batch,
&PhysicalPermission::Allow,
Some(&predicate),
&ctx,
&checked,
true,
)
.await
.unwrap();
assert!(batch.is_empty());
let unchecked = build_field_state(&ctx, &table, false, None).await.unwrap();
let mut batch = rows(&["{ id: t:1, secret: 'x' }"]).await;
filter_and_process_batch(
&mut batch,
&PhysicalPermission::Allow,
Some(&predicate),
&ctx,
&unchecked,
false,
)
.await
.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(pick(&batch[0], "secret"), Value::from("x"));
let survives = physical_expr("id = t:1", &ctx).await;
let mut batch = rows(&["{ id: t:1, secret: 'x' }"]).await;
filter_and_process_batch(
&mut batch,
&PhysicalPermission::Allow,
Some(&survives),
&ctx,
&checked,
true,
)
.await
.unwrap();
assert_eq!(batch.len(), 1);
assert_eq!(pick(&batch[0], "secret"), Value::None);
}
#[tokio::test]
async fn the_predicate_only_fast_path_matches_the_slow_path_and_keeps_row_order() {
let ctx = root_ctx();
let predicate = physical_expr("n % 2 = 1", &ctx).await;
let empty = FieldState::empty();
let input =
["{ n: 0 }", "{ n: 1 }", "{ n: 2 }", "{ n: 3 }", "{ n: 4 }", "{ n: 5 }", "{ n: 6 }"];
let mut fast = rows(&input).await;
filter_and_process_batch(
&mut fast,
&PhysicalPermission::Allow,
Some(&predicate),
&ctx,
&empty,
true,
)
.await
.unwrap();
let allow_all = PhysicalPermission::Conditional(physical_expr("true", &ctx).await);
let mut slow = rows(&input).await;
filter_and_process_batch(&mut slow, &allow_all, Some(&predicate), &ctx, &empty, true)
.await
.unwrap();
let expected = vec![Value::from(1), Value::from(3), Value::from(5)];
assert_eq!(ns(&fast), expected);
assert_eq!(ns(&slow), expected);
}
#[tokio::test]
async fn computed_fields_are_evaluated_in_dependency_order() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD base ON t TYPE int;
DEFINE FIELD y ON t TYPE int COMPUTED z * 10;
DEFINE FIELD z ON t TYPE int COMPUTED base + 1;",
)
.await;
let ctx = db.exec_ctx().await;
let field_state =
build_field_state(&ctx, &TableName::from("t"), false, None).await.unwrap();
assert_eq!(field_state.computed_fields.len(), 2);
assert_eq!(field_state.computed_fields[0].field_name(), "z");
let mut row = val("{ id: t:1, base: 1 }").await;
compute_fields_for_value(&ctx, &field_state, &mut row, false).await.unwrap();
assert_eq!(pick(&row, "z"), Value::from(2));
assert_eq!(pick(&row, "y"), Value::from(20));
}
#[tokio::test]
async fn a_computed_field_is_coerced_to_its_declared_kind() {
let ctx = root_ctx();
let field_state = state(
vec![computed("ratio", physical_expr("1", &ctx).await, Some(crate::expr::Kind::Float))],
Vec::new(),
);
let mut row = val("{ id: t:1 }").await;
compute_fields_for_value(&ctx, &field_state, &mut row, false).await.unwrap();
assert!(
matches!(pick(&row, "ratio"), Value::Number(Number::Float(f)) if f == 1.0),
"expected a float, got {:?}",
pick(&row, "ratio")
);
}
#[tokio::test]
async fn a_failed_coercion_surfaces_as_an_error_naming_the_field() {
let ctx = root_ctx();
let field_state = state(
vec![computed(
"count",
physical_expr("'abc'", &ctx).await,
Some(crate::expr::Kind::Int),
)],
Vec::new(),
);
let mut row = val("{ id: t:1 }").await;
let err = compute_fields_for_value(&ctx, &field_state, &mut row, false).await.unwrap_err();
let message = err.to_string();
assert!(
message.contains("Failed to coerce computed field 'count'"),
"expected the coercion context, got {message}"
);
}
#[tokio::test]
async fn a_return_out_of_a_computed_body_becomes_the_field_value() {
let ctx = root_ctx();
let field_state = state(
vec![computed("answer", physical_expr("{ RETURN 5 }", &ctx).await, None)],
Vec::new(),
);
let mut row = val("{ id: t:1 }").await;
compute_fields_for_value(&ctx, &field_state, &mut row, false).await.unwrap();
assert_eq!(pick(&row, "answer"), Value::from(5));
}
#[tokio::test]
async fn computing_a_field_on_a_non_object_row_is_an_error() {
let ctx = root_ctx();
let field_state = state(vec![ComputedFieldDef::for_test("x")], Vec::new());
let mut row = Value::from(1);
let err = compute_fields_for_value(&ctx, &field_state, &mut row, false).await.unwrap_err();
assert!(
err.to_string().contains("Value is not an object"),
"expected the non-object message, got {err}"
);
let mut row = Value::from(1);
compute_fields_for_value(&ctx, &FieldState::empty(), &mut row, false).await.unwrap();
assert_eq!(row, Value::from(1));
}
#[tokio::test]
async fn computing_record_is_taken_from_the_rows_id_so_self_reads_stay_raw() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD c ON t TYPE int COMPUTED 7;
DEFINE FIELD self_c ON t TYPE any COMPUTED id.c;
DEFINE FIELD other_c ON t TYPE any COMPUTED (t:two).c;
CREATE t:one;
CREATE t:two;",
)
.await;
let ctx = db.exec_ctx().await;
let field_state =
build_field_state(&ctx, &TableName::from("t"), false, None).await.unwrap();
let mut row = val("{ id: t:one }").await;
compute_fields_for_value(&ctx, &field_state, &mut row, false).await.unwrap();
assert_eq!(pick(&row, "self_c"), Value::None);
assert_eq!(pick(&row, "other_c"), Value::from(7));
}
#[tokio::test]
async fn a_deny_field_permission_cuts_the_field() {
let ctx = root_ctx();
let field_state =
state(Vec::new(), vec![(parse_idiom("secret"), PhysicalPermission::Deny)]);
let mut row = val("{ id: t:1, secret: 'x', public: 'y' }").await;
filter_fields_by_permission(&ctx, &field_state, &mut row).await.unwrap();
assert_eq!(pick(&row, "secret"), Value::None);
assert_eq!(pick(&row, "public"), Value::from("y"));
let allow = state(Vec::new(), vec![(parse_idiom("secret"), PhysicalPermission::Allow)]);
let mut row = val("{ secret: 'x' }").await;
filter_fields_by_permission(&ctx, &allow, &mut row).await.unwrap();
assert_eq!(pick(&row, "secret"), Value::from("x"));
let mut scalar = Value::from(1);
filter_fields_by_permission(&ctx, &field_state, &mut scalar).await.unwrap();
assert_eq!(scalar, Value::from(1));
}
#[tokio::test]
async fn a_conditional_field_permission_decides_from_the_picked_field_value() {
let ctx = root_ctx();
let field_state = state(
Vec::new(),
vec![(
parse_idiom("score"),
PhysicalPermission::Conditional(physical_expr("$value > 10", &ctx).await),
)],
);
let mut kept = val("{ score: 42 }").await;
filter_fields_by_permission(&ctx, &field_state, &mut kept).await.unwrap();
assert_eq!(pick(&kept, "score"), Value::from(42));
let mut cut = val("{ score: 3 }").await;
filter_fields_by_permission(&ctx, &field_state, &mut cut).await.unwrap();
assert_eq!(pick(&cut, "score"), Value::None);
}
#[tokio::test]
async fn a_field_permission_predicate_reads_the_pre_mutation_snapshot() {
let ctx = root_ctx();
let field_state = state(
Vec::new(),
vec![
(parse_idiom("a"), PhysicalPermission::Deny),
(
parse_idiom("b"),
PhysicalPermission::Conditional(physical_expr("a = 1", &ctx).await),
),
],
);
let mut row = val("{ a: 1, b: 2 }").await;
filter_fields_by_permission(&ctx, &field_state, &mut row).await.unwrap();
assert_eq!(pick(&row, "a"), Value::None);
assert_eq!(pick(&row, "b"), Value::from(2));
}
#[tokio::test]
async fn every_element_a_wildcard_deny_targets_is_cut() {
let ctx = root_ctx();
let field_state =
state(Vec::new(), vec![(parse_idiom("items[*]"), PhysicalPermission::Deny)]);
let mut row = val("{ items: [1, 2, 3, 4, 5] }").await;
filter_fields_by_permission(&ctx, &field_state, &mut row).await.unwrap();
assert_eq!(pick(&row, "items"), val("[]").await);
}
#[tokio::test]
async fn a_wildcard_conditional_permission_cuts_every_rejected_element() {
let ctx = root_ctx();
let field_state = state(
Vec::new(),
vec![(
parse_idiom("items[*]"),
PhysicalPermission::Conditional(physical_expr("$value % 2 = 0", &ctx).await),
)],
);
let mut row = val("{ items: [1, 2, 3, 4, 5, 6] }").await;
filter_fields_by_permission(&ctx, &field_state, &mut row).await.unwrap();
assert_eq!(pick(&row, "items"), val("[2, 4, 6]").await);
}
async fn projection_state(permission_deps_complete: bool) -> FieldState {
let ctx = root_ctx();
let mut dep_map = HashMap::new();
dep_map.insert(
"flag".to_owned(),
ComputedDeps {
fields: vec!["score".to_owned()],
is_complete: true,
},
);
dep_map.insert(
"other".to_owned(),
ComputedDeps {
fields: Vec::new(),
is_complete: true,
},
);
FieldState {
computed_fields: vec![
computed("flag", physical_expr("score >= 10", &ctx).await, None),
computed("other", physical_expr("1", &ctx).await, None),
],
field_permissions: Arc::new(vec![
(
parse_idiom("secret"),
PhysicalPermission::Conditional(physical_expr("flag", &ctx).await),
),
(parse_idiom("hidden"), PhysicalPermission::Deny),
]),
dep_map: Arc::new(dep_map),
permission_field_deps: Arc::new(HashSet::from(["flag".to_owned()])),
permission_deps_complete,
}
}
fn computed_names(state: &FieldState) -> Vec<&str> {
state.computed_fields.iter().map(|cf| cf.field_name()).collect()
}
#[tokio::test]
async fn no_projection_keeps_every_computed_field() {
let full = projection_state(true).await;
let filtered = filter_field_state_for_projection(&full, None);
assert_eq!(computed_names(&filtered), vec!["flag", "other"]);
}
#[tokio::test]
async fn a_selective_projection_drops_the_computed_fields_it_does_not_need() {
let mut full = projection_state(true).await;
full.permission_field_deps = Arc::new(HashSet::new());
let needed = HashSet::from(["other".to_owned()]);
let filtered = filter_field_state_for_projection(&full, Some(&needed));
assert_eq!(computed_names(&filtered), vec!["other"]);
}
#[tokio::test]
async fn a_selective_projection_still_computes_fields_a_field_permission_reads() {
let full = projection_state(true).await;
let needed = HashSet::from(["secret".to_owned()]);
let filtered = filter_field_state_for_projection(&full, Some(&needed));
assert_eq!(computed_names(&filtered), vec!["flag"]);
}
#[tokio::test]
async fn incomplete_permission_dependencies_force_every_computed_field() {
let full = projection_state(false).await;
let needed = HashSet::from(["unrelated".to_owned()]);
let filtered = filter_field_state_for_projection(&full, Some(&needed));
assert_eq!(computed_names(&filtered), vec!["flag", "other"]);
}
#[tokio::test]
async fn field_permissions_are_never_filtered_by_the_projection() {
let full = projection_state(true).await;
let needed = HashSet::new();
let filtered = filter_field_state_for_projection(&full, Some(&needed));
assert_eq!(filtered.field_permissions.len(), 2);
assert_eq!(filter_field_state_for_projection(&full, None).field_permissions.len(), 2);
}
#[tokio::test]
async fn a_plain_table_resolves_to_the_empty_field_state() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD name ON t TYPE string;",
)
.await;
let ctx = db.exec_ctx().await;
let field_state = build_field_state(&ctx, &TableName::from("t"), true, None).await.unwrap();
assert!(field_state.computed_fields.is_empty());
assert!(field_state.field_permissions.is_empty());
}
#[tokio::test]
async fn computed_fields_and_field_permissions_come_from_the_catalog() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD score ON t TYPE int;
DEFINE FIELD doubled ON t TYPE int COMPUTED score * 2;
DEFINE FIELD secret ON t TYPE string PERMISSIONS FOR select NONE;
DEFINE FIELD gated ON t TYPE string PERMISSIONS FOR select WHERE score > 1;",
)
.await;
let ctx = db.exec_ctx().await;
let table = TableName::from("t");
let checked = build_field_state(&ctx, &table, true, None).await.unwrap();
assert_eq!(computed_names(&checked), vec!["doubled"]);
let listed: Vec<String> =
checked.field_permissions.iter().map(|(idiom, _)| idiom.to_raw_string()).collect();
assert_eq!(listed, vec!["gated".to_owned(), "secret".to_owned()]);
assert!(matches!(checked.field_permissions[0].1, PhysicalPermission::Conditional(_)));
assert!(matches!(checked.field_permissions[1].1, PhysicalPermission::Deny));
assert!(checked.permission_deps_complete);
assert!(checked.permission_field_deps.contains("score"));
let unchecked = build_field_state(&ctx, &table, false, None).await.unwrap();
assert_eq!(computed_names(&unchecked), vec!["doubled"]);
assert!(unchecked.field_permissions.is_empty());
}
#[tokio::test]
async fn field_state_is_cached_per_table_and_check_perms_flag() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD score ON t TYPE int;
DEFINE FIELD doubled ON t TYPE int COMPUTED score * 2;",
)
.await;
let ctx = db.exec_ctx().await;
let db_ctx = ctx.database().unwrap();
let table = TableName::from("t");
build_field_state(&ctx, &table, true, None).await.unwrap();
build_field_state(&ctx, &table, true, None).await.unwrap();
{
let cache = db_ctx.field_state_cache.read().await;
assert_eq!(cache.len(), 1);
assert!(cache.contains_key(&(table.clone(), true)));
}
build_field_state(&ctx, &table, false, None).await.unwrap();
let cache = db_ctx.field_state_cache.read().await;
assert_eq!(cache.len(), 2);
assert!(cache.contains_key(&(table, false)));
}
#[tokio::test]
async fn a_projection_filters_the_cached_state_without_narrowing_the_cache() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
DEFINE FIELD score ON t TYPE int;
DEFINE FIELD doubled ON t TYPE int COMPUTED score * 2;
DEFINE FIELD tripled ON t TYPE int COMPUTED score * 3;",
)
.await;
let ctx = db.exec_ctx().await;
let table = TableName::from("t");
let needed = HashSet::from(["doubled".to_owned()]);
let projected = build_field_state(&ctx, &table, false, Some(&needed)).await.unwrap();
assert_eq!(computed_names(&projected), vec!["doubled"]);
let full = build_field_state(&ctx, &table, false, None).await.unwrap();
assert_eq!(full.computed_fields.len(), 2);
}
#[tokio::test]
async fn a_non_negative_integer_limit_evaluates_to_itself() {
let ctx = root_ctx();
let expr = physical_expr("7", &ctx).await;
assert_eq!(eval_limit_expr(expr.as_ref(), &ctx).await.unwrap(), 7);
let zero = physical_expr("0", &ctx).await;
assert_eq!(eval_limit_expr(zero.as_ref(), &ctx).await.unwrap(), 0);
}
#[tokio::test]
async fn an_absent_limit_means_no_offset() {
let ctx = root_ctx();
for src in ["NONE", "NULL"] {
let expr = physical_expr(src, &ctx).await;
assert_eq!(
eval_limit_expr(expr.as_ref(), &ctx).await.unwrap(),
0,
"{src} should evaluate to 0"
);
}
}
#[tokio::test]
async fn a_negative_limit_is_rejected() {
let ctx = root_ctx();
let expr = physical_expr("-1", &ctx).await;
let err = eval_limit_expr(expr.as_ref(), &ctx).await.unwrap_err();
assert!(
err.to_string().contains("non-negative"),
"expected the non-negative message, got {err}"
);
}
#[tokio::test]
async fn a_non_numeric_limit_is_rejected() {
let ctx = root_ctx();
let expr = physical_expr("'ten'", &ctx).await;
let err = eval_limit_expr(expr.as_ref(), &ctx).await.unwrap_err();
assert!(
err.to_string().contains("must be an integer"),
"expected the integer message, got {err}"
);
}
#[tokio::test]
async fn ordering_by_id_descending_scans_backward() {
let ordering = order_by("id", false);
assert_eq!(determine_scan_direction(Some(&ordering)), Direction::Backward);
}
#[tokio::test]
async fn every_other_ordering_scans_forward() {
let id_asc = order_by("id", true);
assert_eq!(determine_scan_direction(Some(&id_asc)), Direction::Forward);
let name_desc = order_by("name", false);
assert_eq!(determine_scan_direction(Some(&name_desc)), Direction::Forward);
let name_then_id = Ordering::Order(OrderList(vec![
Order {
value: parse_idiom("name"),
direction: false,
..Default::default()
},
Order {
value: parse_idiom("id"),
direction: false,
..Default::default()
},
]));
assert_eq!(determine_scan_direction(Some(&name_then_id)), Direction::Forward);
assert_eq!(determine_scan_direction(None), Direction::Forward);
let empty = Ordering::Order(OrderList(Vec::new()));
assert_eq!(determine_scan_direction(Some(&empty)), Direction::Forward);
assert_eq!(determine_scan_direction(Some(&Ordering::Random)), Direction::Forward);
}
#[tokio::test]
async fn a_stored_record_decodes_with_its_id_taken_from_the_key() {
let db = TestDb::new(
"DEFINE TABLE t SCHEMALESS;
CREATE t:tobie SET name = 'Tobie', age = 30;",
)
.await;
let ctx = db.exec_ctx().await;
let db_ctx = ctx.database().unwrap();
let table = TableName::from("t");
let range = RecordPrefix {
ns: db_ctx.ns_ctx.ns.namespace_id,
db: db_ctx.db.database_id,
tb: Cow::Borrowed(&table),
}
.range()
.unwrap();
let raw = ctx.txn().scan_raw(range, 10, 0, None).await.unwrap();
assert_eq!(raw.len(), 1, "the table holds exactly one record");
let (key, value) = &raw[0];
let decoded = decode_record(key, value).unwrap();
assert_eq!(pick(&decoded, "name"), Value::from("Tobie"));
assert_eq!(pick(&decoded, "age"), Value::from(30));
assert_eq!(pick(&decoded, "id"), val("t:tobie").await);
assert!(decode_record(b"not-a-record-key", value).is_err());
}
}