use std::sync::Arc;
use super::pipeline::{
FieldState, build_field_state, compute_fields_for_value, filter_fields_by_permission,
};
use crate::catalog::providers::TableProvider;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::exec::{ControlFlowExt, EvalContext, ExecutionContext, PhysicalExpr};
use crate::expr::ControlFlow;
use crate::kvs::{CachePolicy, Transaction};
use crate::val::{RecordId, RecordIdKey, TableName, Value};
pub(crate) const DEFAULT_SCAN_BATCH_SIZE: usize = 1000;
pub(crate) fn value_to_record_id_key(val: Value) -> RecordIdKey {
match val {
Value::Number(n) => RecordIdKey::Number(n.as_int()),
Value::String(s) => RecordIdKey::String(s),
Value::Uuid(u) => RecordIdKey::Uuid(u),
Value::Array(a) => RecordIdKey::Array(a),
Value::Object(o) => RecordIdKey::Object(o),
other => RecordIdKey::String(other.to_raw_string().into()),
}
}
pub(crate) fn extract_record_ids_into(val: Value, rids: &mut Vec<RecordId>) {
match val {
Value::RecordId(rid) => rids.push(rid),
Value::Object(mut obj) => {
if let Some(id_val) = obj.remove("id") {
extract_record_ids_into(id_val, rids);
}
}
Value::Array(arr) => {
for v in arr {
extract_record_ids_into(v, rids);
}
}
_ => {}
}
}
pub(crate) async fn evaluate_bound_key(
expr: &Arc<dyn PhysicalExpr>,
ctx: &ExecutionContext,
) -> Result<RecordIdKey, ControlFlow> {
let eval_ctx = EvalContext::from_exec_ctx(ctx);
let val = expr.evaluate(eval_ctx).await?;
Ok(value_to_record_id_key(val))
}
pub(crate) async fn resolve_version_stamp(
ctx: &ExecutionContext,
version_expr: Option<&Arc<dyn PhysicalExpr>>,
) -> Result<Option<u64>, ControlFlow> {
if let Some(stamp) = ctx.version_stamp() {
return Ok(Some(stamp));
}
let Some(expr) = version_expr else {
return Ok(None);
};
let eval_ctx = EvalContext::from_exec_ctx(ctx);
let v = expr.evaluate(eval_ctx).await?;
let stamp = v
.cast_to::<crate::val::Datetime>()
.map_err(|e| anyhow::anyhow!("{e}"))?
.to_version_stamp(ctx.txn().timestamp_impl().as_ref())?;
Ok(Some(stamp))
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn resolve_record_batch(
ctx: &ExecutionContext,
txn: &Transaction,
ns_id: NamespaceId,
db_id: DatabaseId,
rids: &[RecordId],
fetch_full: bool,
check_perms: bool,
version: Option<u64>,
cache_policy: CachePolicy,
perm_cache: &mut std::collections::HashMap<
crate::val::TableName,
crate::exec::permission::PhysicalPermission,
>,
) -> Result<Vec<Value>, ControlFlow> {
use crate::exec::permission::{PhysicalPermission, check_permission_for_value};
if !check_perms && !fetch_full {
return Ok(rids.iter().map(|rid| Value::RecordId(rid.clone())).collect());
}
if check_perms {
let db_ctx = ctx.database().context("permission resolution requires database context")?;
for rid in rids {
if perm_cache.contains_key(&rid.table) {
continue;
}
let table_def = db_ctx
.get_table_def(&rid.table, version)
.await
.context("Failed to get table definition")?;
let catalog_perm =
crate::exec::permission::resolve_select_permission(table_def.as_deref());
let perm = crate::exec::permission::convert_permission_to_physical_runtime(
catalog_perm,
ctx.ctx(),
)
.await
.context("Failed to convert permission")?;
perm_cache.insert(rid.table.clone(), perm);
}
}
let records = txn
.get_records(ns_id, db_id, rids, version, cache_policy)
.await
.context("Failed to fetch records")?;
let mut field_state_cache: std::collections::HashMap<TableName, FieldState> =
std::collections::HashMap::new();
let skip_fetch_perms = ctx.root().skip_fetch_perms;
let mut values = Vec::with_capacity(rids.len());
for (rid, record) in rids.iter().zip(records) {
if record.data.is_none() {
continue;
}
if check_perms {
let perm = perm_cache.get(&rid.table).map_or(&PhysicalPermission::Deny, |p| p);
let allowed = check_permission_for_value(perm, &record.data, None, ctx)
.await
.context("Failed to check permission")?;
if !allowed {
continue;
}
}
if fetch_full {
let mut value = match Arc::try_unwrap(record) {
Ok(rec) => rec.data,
Err(arc) => arc.data.clone(),
};
if !field_state_cache.contains_key(&rid.table) {
let fs = build_field_state(ctx, &rid.table, check_perms, None).await?;
field_state_cache.insert(rid.table.clone(), fs);
}
let field_state = &field_state_cache[&rid.table];
compute_fields_for_value(ctx, field_state, &mut value, skip_fetch_perms).await?;
if check_perms {
filter_fields_by_permission(ctx, field_state, &mut value).await?;
}
values.push(value);
} else {
values.push(Value::RecordId(rid.clone()));
}
}
Ok(values)
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn fetch_and_filter_records_batch(
ctx: &ExecutionContext,
txn: &Transaction,
ns_id: NamespaceId,
db_id: DatabaseId,
rids: &[RecordId],
select_permission: &crate::exec::permission::PhysicalPermission,
check_perms: bool,
version: Option<u64>,
cache_policy: CachePolicy,
) -> Result<Vec<Value>, ControlFlow> {
let records = txn
.get_records(ns_id, db_id, rids, version, cache_policy)
.await
.context("Failed to fetch records")?;
let mut values = Vec::with_capacity(rids.len());
for record in records {
if record.data.is_none() {
continue;
}
if check_perms {
let allowed = crate::exec::permission::check_permission_for_value(
select_permission,
&record.data,
None,
ctx,
)
.await
.context("Failed to check permission")?;
if !allowed {
continue;
}
}
let value = match Arc::try_unwrap(record) {
Ok(rec) => rec.data,
Err(arc) => arc.data.clone(),
};
values.push(value);
}
Ok(values)
}