use akar_common::arrow_vector::{ArrowVector, VectorAccess};
use akar_common::error::ProcessorError;
use akar_common::types::{PhysicalTypeID, Value};
use akar_common::vector::{DataChunk, ValueVector};
use akar_function::registry::{FunctionRegistry, ScalarFunction};
use akar_function::scalar::evaluate_scalar;
use akar_parser::ast::{BinaryOp, Constant, Expression, Query, UnaryOp};
use arrow::array::{Array, ArrayRef};
use std::sync::{Arc, Mutex};
pub type SubqueryFn = Arc<dyn Fn(&Query) -> Result<Vec<DataChunk>, ProcessorError> + Send + Sync>;
pub type SequenceFn = Arc<dyn Fn(&str, bool) -> Result<Value, ProcessorError> + Send + Sync>;
pub fn map_property_value(val: &Value, prop: &str) -> Value {
match val {
Value::Map(entries) => entries
.iter()
.find(|(k, _)| matches!(k, Value::String(s) if s == prop))
.map(|(_, v)| v.clone())
.unwrap_or(Value::Null),
Value::Struct(fields) => fields
.iter()
.find(|(name, _)| name == prop)
.map(|(_, v)| v.clone())
.unwrap_or(Value::Null),
_ => Value::Null,
}
}
impl std::fmt::Debug for ExpressionEvaluator {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ExpressionEvaluator").finish()
}
}
pub struct ExpressionEvaluator {
registry: Arc<Mutex<FunctionRegistry>>,
pub subquery_fn: Option<SubqueryFn>,
pub sequence_fn: Option<SequenceFn>,
}
impl ExpressionEvaluator {
pub fn new(registry: Arc<Mutex<FunctionRegistry>>) -> Self {
Self {
registry,
subquery_fn: None,
sequence_fn: None,
}
}
pub fn with_subquery_fn(mut self, f: SubqueryFn) -> Self {
self.subquery_fn = Some(f);
self
}
pub fn with_sequence_fn(mut self, f: SequenceFn) -> Self {
self.sequence_fn = Some(f);
self
}
fn evaluate_subquery(&self, query: &Query) -> Result<Vec<DataChunk>, ProcessorError> {
if let Some(ref f) = self.subquery_fn {
f(query)
} else {
Err("No subquery executor configured".into())
}
}
pub fn evaluate_to_arrow(&self, expr: &Expression, chunk: &DataChunk) -> Result<ArrowVector, ProcessorError> {
self.evaluate_arrow(expr, chunk)
}
pub fn evaluate(&self, expr: &Expression, chunk: &DataChunk) -> Result<ValueVector, ProcessorError> {
match expr {
Expression::Constant(c) => self.evaluate_constant(c, chunk.size),
Expression::Variable(name) => self.evaluate_variable(name, chunk),
Expression::PropertyAccess(obj, prop) => self.evaluate_property_access(obj, prop, chunk),
Expression::FunctionCall(name, args) => {
let refs: Vec<&Expression> = args.iter().collect();
self.evaluate_function_call(name, &refs, chunk)
}
Expression::BinaryOp(op, left, right) => self.evaluate_binary_op(op, left, right, chunk),
Expression::UnaryOp(op, inner) => self.evaluate_unary_op(op, inner, chunk),
Expression::List(items) => self.evaluate_list_literal(items, chunk),
Expression::Map(items) => self.evaluate_map_literal(items, chunk),
Expression::Parameter(_) => {
let mut v = ValueVector::new(akar_common::types::PhysicalTypeID::Int64, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
v.set_null(i, true);
}
Ok(v)
}
Expression::ExistsSubquery(query) => {
let result = self.evaluate_subquery(query)?;
let exists = !result.is_empty() && result.iter().any(|c| c.size > 0);
let mut v = ValueVector::new(akar_common::types::PhysicalTypeID::Bool, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
store_value_in_vector_simple(&mut v, i, &Value::Bool(exists))?;
}
Ok(v)
}
Expression::Case(case_expr) => self.evaluate_case(case_expr, chunk),
Expression::Star => {
Err("STAR expression should be expanded by the binder before reaching the evaluator".into())
}
Expression::ListPredicate {
quantifier,
list,
var_name,
predicate,
} => self.evaluate_list_predicate(quantifier, list, var_name, predicate, chunk),
Expression::Lambda { .. } => {
Err("Lambda expression should only appear as argument to list_transform/filter/reduce".into())
}
}
}
pub fn evaluate_arrow(&self, expr: &Expression, chunk: &DataChunk) -> Result<ArrowVector, ProcessorError> {
match expr {
Expression::Constant(c) => self.evaluate_arrow_constant(c, chunk.size),
Expression::Variable(name) => self.evaluate_arrow_variable(name, chunk),
Expression::PropertyAccess(obj, prop) => self.evaluate_arrow_property_access(obj, prop, chunk),
Expression::FunctionCall(name, args) => {
let refs: Vec<&Expression> = args.iter().collect();
self.evaluate_arrow_function_call(name, &refs, chunk)
}
Expression::BinaryOp(op, left, right) => self.evaluate_arrow_binary_op(op, left, right, chunk),
Expression::UnaryOp(op, inner) => self.evaluate_arrow_unary_op(op, inner, chunk),
Expression::List(items) => self.evaluate_arrow_list_literal(items, chunk),
Expression::Map(items) => self.evaluate_arrow_map_literal(items, chunk),
_ => {
let legacy = self.evaluate(expr, chunk)?;
Ok(ArrowVector::from_legacy(&legacy))
}
}
}
fn evaluate_arrow_list_literal(
&self,
items: &[Expression],
chunk: &DataChunk,
) -> Result<ArrowVector, ProcessorError> {
let num_rows = chunk.size;
let mut row_results: Vec<Value> = Vec::with_capacity(num_rows);
for row in 0..num_rows {
let mut list_values = Vec::with_capacity(items.len());
for item in items {
let item_vec = self.evaluate(item, chunk)?;
let val = if row < item_vec.size() {
item_vec.get_value(row).unwrap_or(Value::Null)
} else {
Value::Null
};
list_values.push(val);
}
row_results.push(Value::List(list_values));
}
build_arrow_from_values(&row_results, PhysicalTypeID::List, num_rows)
}
fn evaluate_arrow_map_literal(
&self,
items: &[(String, Expression)],
chunk: &DataChunk,
) -> Result<ArrowVector, ProcessorError> {
let num_rows = chunk.size;
let mut row_results: Vec<Value> = Vec::with_capacity(num_rows);
for row in 0..num_rows {
let mut entries = Vec::with_capacity(items.len());
for (key, item) in items {
let item_vec = self.evaluate(item, chunk)?;
let val = if row < item_vec.size() {
item_vec.get_value(row).unwrap_or(Value::Null)
} else {
Value::Null
};
entries.push((key.clone(), val));
}
row_results.push(Value::Struct(entries));
}
build_arrow_from_values(&row_results, PhysicalTypeID::Struct, num_rows)
}
fn evaluate_arrow_constant(&self, c: &Constant, size: usize) -> Result<ArrowVector, ProcessorError> {
match c {
Constant::Null => {
let mut builder = arrow::array::Int64Builder::with_capacity(size);
builder.append_nulls(size);
Ok(ArrowVector::new(Arc::new(builder.finish()), PhysicalTypeID::Int64))
}
Constant::Bool(b) => {
let mut builder = arrow::array::BooleanBuilder::with_capacity(size);
for _ in 0..size {
builder.append_value(*b);
}
Ok(ArrowVector::new(Arc::new(builder.finish()), PhysicalTypeID::Bool))
}
Constant::Integer(i) => {
let mut builder = arrow::array::Int64Builder::with_capacity(size);
let v = *i;
for _ in 0..size {
builder.append_value(v);
}
Ok(ArrowVector::new(Arc::new(builder.finish()), PhysicalTypeID::Int64))
}
Constant::Float(f) => {
let mut builder = arrow::array::Float64Builder::with_capacity(size);
let v = *f;
for _ in 0..size {
builder.append_value(v);
}
Ok(ArrowVector::new(Arc::new(builder.finish()), PhysicalTypeID::Double))
}
Constant::String(s) => {
let mut builder = arrow::array::StringBuilder::with_capacity(size, size * s.len().max(1));
for _ in 0..size {
builder.append_value(s);
}
Ok(ArrowVector::new(Arc::new(builder.finish()), PhysicalTypeID::String))
}
}
}
fn evaluate_arrow_variable(&self, name: &str, chunk: &DataChunk) -> Result<ArrowVector, ProcessorError> {
let idx = if let Ok(idx) = name.parse::<usize>() {
idx
} else if !chunk.field_names.is_empty() {
if let Some(idx) = chunk.field_names.iter().position(|n| n == name) {
idx
} else {
return Err(format!("Variable '{}' not found in field_names", name).into());
}
} else {
return Err(format!("Variable '{}' not found (chunk has no field_names)", name).into());
};
let field = chunk
.fields
.get(idx)
.ok_or_else(|| format!("Variable '{}' (index {}) not found in chunk fields", name, idx))?;
Ok(ArrowVector::new(field.clone(), chunk.field_types[idx]))
}
fn evaluate_arrow_property_access(
&self,
obj: &Expression,
prop: &str,
chunk: &DataChunk,
) -> Result<ArrowVector, ProcessorError> {
let qualified_prop = if let Expression::Variable(var_name) = obj {
format!("{}.{}", var_name, prop)
} else {
prop.to_string()
};
if !chunk.field_names.is_empty()
&& let Some(idx) = chunk.field_names.iter().position(|n| n == &qualified_prop)
{
let field = chunk
.fields
.get(idx)
.ok_or_else(|| format!("Property '{}' not found in chunk", prop))?;
return Ok(ArrowVector::new(field.clone(), chunk.field_types[idx]));
}
if let Expression::Variable(var_name) = obj
&& let Some(idx) = chunk.field_names.iter().position(|n| n == var_name)
{
let extracted: Vec<Value> = (0..chunk.size)
.map(|i| map_property_value(&chunk.get_value(idx, i).unwrap_or(Value::Null), prop))
.collect();
if extracted.iter().any(|v| !matches!(v, Value::Null)) {
let result_type = extracted
.iter()
.find(|v| !matches!(v, Value::Null))
.map(|v| v.physical_type())
.unwrap_or(PhysicalTypeID::Int64);
return build_arrow_from_values(&extracted, result_type, chunk.size);
}
}
if !chunk.field_names.is_empty()
&& let Some(idx) = chunk.field_names.iter().position(|n| n == prop)
{
let field = chunk
.fields
.get(idx)
.ok_or_else(|| format!("Property '{}' not found in chunk", prop))?;
return Ok(ArrowVector::new(field.clone(), chunk.field_types[idx]));
}
let legacy = self.evaluate(obj, chunk)?;
Ok(ArrowVector::from_legacy(&legacy))
}
fn evaluate_arrow_binary_op(
&self,
op: &BinaryOp,
left: &Expression,
right: &Expression,
chunk: &DataChunk,
) -> Result<ArrowVector, ProcessorError> {
match op {
BinaryOp::In | BinaryOp::NotIn => {
let legacy = self.evaluate_in_op(op, left, right, chunk)?;
return Ok(ArrowVector::from_legacy(&legacy));
}
BinaryOp::Concat | BinaryOp::StartsWith | BinaryOp::EndsWith | BinaryOp::Contains | BinaryOp::Like => {
let func_name = match op {
BinaryOp::Concat => "concat",
BinaryOp::StartsWith => "starts_with",
BinaryOp::EndsWith => "ends_with",
BinaryOp::Contains => "contains",
BinaryOp::Like => "like",
_ => unreachable!(),
};
return self.evaluate_arrow_function_call(func_name, &[left, right], chunk);
}
_ => {}
}
let kernel_name = match op {
BinaryOp::Add => "add",
BinaryOp::Subtract => "sub",
BinaryOp::Multiply => "mul",
BinaryOp::Divide => "div",
BinaryOp::Modulo => "mod",
BinaryOp::Equal => "eq",
BinaryOp::NotEqual => "neq",
BinaryOp::LessThan => "lt",
BinaryOp::LessThanOrEqual => "lt_eq",
BinaryOp::GreaterThan => "gt",
BinaryOp::GreaterThanOrEqual => "gt_eq",
BinaryOp::And => "and",
BinaryOp::Or => "or",
BinaryOp::Xor => "xor",
_ => return Err(format!("Unsupported binary op: {:?}", op).into()),
};
let left_arrow = self.evaluate_arrow(left, chunk)?;
let right_arrow = self.evaluate_arrow(right, chunk)?;
match self.apply_arrow_kernel(kernel_name, &left_arrow, &right_arrow) {
Ok(result) => Ok(result),
Err(_) => {
let legacy = self.evaluate_binary_op(op, left, right, chunk)?;
Ok(ArrowVector::from_legacy(&legacy))
}
}
}
fn evaluate_arrow_unary_op(
&self,
op: &UnaryOp,
inner: &Expression,
chunk: &DataChunk,
) -> Result<ArrowVector, ProcessorError> {
match op {
UnaryOp::Not => {
let inner_arrow = self.evaluate_arrow(inner, chunk)?;
self.apply_arrow_unary_kernel("not", &inner_arrow).or_else(|_| {
let legacy = self.evaluate_unary_op(&UnaryOp::Not, inner, chunk)?;
Ok(ArrowVector::from_legacy(&legacy))
})
}
UnaryOp::Negate => {
let inner_arrow = self.evaluate_arrow(inner, chunk)?;
self.apply_arrow_unary_kernel("negate", &inner_arrow).or_else(|_| {
let legacy = self.evaluate_unary_op(&UnaryOp::Negate, inner, chunk)?;
Ok(ArrowVector::from_legacy(&legacy))
})
}
UnaryOp::IsNull => {
if let Expression::Variable(name) = inner
&& !chunk.field_names.is_empty()
&& !chunk.field_names.iter().any(|n| n == name)
{
let prefix = format!("{}.", name);
let cols: Vec<usize> = chunk
.field_names
.iter()
.enumerate()
.filter(|(_, n)| n.starts_with(prefix.as_str()))
.map(|(i, _)| i)
.collect();
if !cols.is_empty() {
let size = chunk.size;
let mut acc: Option<arrow::array::BooleanArray> = None;
for &col in &cols {
let field = chunk
.fields
.get(col)
.ok_or_else(|| format!("Variable '{}' column {} not found", name, col))?;
let nulls: Vec<bool> = (0..size).map(|i| field.is_null(i)).collect();
let m = arrow::array::BooleanArray::from(nulls);
acc = Some(match acc {
Some(a) => arrow::compute::kernels::boolean::and(&a, &m)
.map_err(|e| format!("Arrow is_null over variable failed: {e}"))?,
None => m,
});
}
return Ok(ArrowVector::new(Arc::new(acc.unwrap()), PhysicalTypeID::Bool));
}
}
let inner_arrow = self.evaluate_arrow(inner, chunk)?;
self.apply_arrow_unary_kernel("is_null", &inner_arrow)
}
UnaryOp::IsNotNull => {
if let Expression::Variable(name) = inner
&& !chunk.field_names.is_empty()
&& !chunk.field_names.iter().any(|n| n == name)
{
let prefix = format!("{}.", name);
let cols: Vec<usize> = chunk
.field_names
.iter()
.enumerate()
.filter(|(_, n)| n.starts_with(prefix.as_str()))
.map(|(i, _)| i)
.collect();
if !cols.is_empty() {
let size = chunk.size;
let mut acc: Option<arrow::array::BooleanArray> = None;
for &col in &cols {
let field = chunk
.fields
.get(col)
.ok_or_else(|| format!("Variable '{}' column {} not found", name, col))?;
let nulls: Vec<bool> = (0..size).map(|i| field.is_null(i)).collect();
let m = arrow::array::BooleanArray::from(nulls);
acc = Some(match acc {
Some(a) => arrow::compute::kernels::boolean::and(&a, &m)
.map_err(|e| format!("Arrow is_not_null over variable failed: {e}"))?,
None => m,
});
}
let all_null = acc.unwrap();
let not_all = arrow::compute::kernels::boolean::not(&all_null)
.map_err(|e| format!("Arrow is_not_null over variable failed: {e}"))?;
return Ok(ArrowVector::new(Arc::new(not_all), PhysicalTypeID::Bool));
}
}
let inner_arrow = self.evaluate_arrow(inner, chunk)?;
self.apply_arrow_unary_kernel("is_not_null", &inner_arrow)
}
}
}
fn apply_arrow_kernel(
&self,
name: &str,
left: &ArrowVector,
right: &ArrowVector,
) -> Result<ArrowVector, ProcessorError> {
use arrow::compute::kernels::boolean::{and_kleene, or_kleene};
use arrow::compute::kernels::cmp::{eq, gt, gt_eq, lt, lt_eq, neq};
use arrow::compute::kernels::numeric::{add, div, mul, rem, sub};
let result: ArrayRef = match name {
"add" => Arc::new(add(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"sub" => Arc::new(sub(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"mul" => Arc::new(mul(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"div" => Arc::new(div(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"mod" => Arc::new(rem(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"eq" => Arc::new(eq(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"neq" => Arc::new(neq(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"lt" => Arc::new(lt(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"lt_eq" => Arc::new(lt_eq(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"gt" => Arc::new(gt(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"gt_eq" => Arc::new(gt_eq(&left.array, &right.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"and" => {
let l = left
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", left.array.data_type()))?;
let r = right
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", right.array.data_type()))?;
Arc::new(and_kleene(l, r).map_err(|e| format!("Arrow {name} failed: {e}"))?)
}
"or" => {
let l = left
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", left.array.data_type()))?;
let r = right
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", right.array.data_type()))?;
Arc::new(or_kleene(l, r).map_err(|e| format!("Arrow {name} failed: {e}"))?)
}
"xor" => {
let l = left
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", left.array.data_type()))?;
let r = right
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", right.array.data_type()))?;
let not_r =
arrow::compute::kernels::boolean::not(r).map_err(|e| format!("Arrow xor/not failed: {e}"))?;
let not_l =
arrow::compute::kernels::boolean::not(l).map_err(|e| format!("Arrow xor/not failed: {e}"))?;
let l_and_not_r = and_kleene(l, ¬_r).map_err(|e| format!("Arrow xor/and failed: {e}"))?;
let not_l_and_r = and_kleene(¬_l, r).map_err(|e| format!("Arrow xor/and failed: {e}"))?;
Arc::new(or_kleene(&l_and_not_r, ¬_l_and_r).map_err(|e| format!("Arrow xor/or failed: {e}"))?)
}
_ => return Err(format!("Unknown binary kernel: {name}").into()),
};
let phys_type = match name {
"eq" | "neq" | "lt" | "lt_eq" | "gt" | "gt_eq" | "and" | "or" | "xor" => PhysicalTypeID::Bool,
_ => left.physical_type,
};
Ok(ArrowVector::new(result, phys_type))
}
fn apply_arrow_unary_kernel(&self, name: &str, arr: &ArrowVector) -> Result<ArrowVector, ProcessorError> {
use arrow::compute::kernels::boolean::{is_not_null, is_null, not};
use arrow::compute::kernels::numeric::neg;
let result: ArrayRef = match name {
"not" => {
let arr_ref = arr
.array
.as_any()
.downcast_ref::<arrow::array::BooleanArray>()
.ok_or_else(|| format!("Arrow {name}: expected BooleanArray, got {:?}", arr.array.data_type()))?;
Arc::new(not(arr_ref).map_err(|e| format!("Arrow {name} failed: {e}"))?)
}
"negate" => Arc::new(neg(&*arr.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"is_null" => Arc::new(is_null(&*arr.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
"is_not_null" => Arc::new(is_not_null(&*arr.array).map_err(|e| format!("Arrow {name} failed: {e}"))?),
_ => return Err(format!("Unknown unary kernel: {name}").into()),
};
let phys_type = if matches!(name, "is_null" | "is_not_null" | "not") {
PhysicalTypeID::Bool
} else {
arr.physical_type
};
Ok(ArrowVector::new(result, phys_type))
}
fn evaluate_arrow_function_call(
&self,
name: &str,
args: &[&Expression],
chunk: &DataChunk,
) -> Result<ArrowVector, ProcessorError> {
if let Some(_lambda) = self.extract_lambda_arg(args) {
match name {
"list_transform" | "list_filter" | "list_reduce" => {
let legacy = self.evaluate_function_call(name, args, chunk)?;
return Ok(ArrowVector::from_legacy(&legacy));
}
_ => {}
}
}
let arg_arrows: Vec<ArrowVector> = args
.iter()
.map(|arg| self.evaluate_arrow(arg, chunk))
.collect::<Result<Vec<_>, _>>()?;
if arg_arrows.is_empty() {
return Err(format!("Function '{}' requires at least one argument", name).into());
}
let num_rows = arg_arrows[0].size();
let func = {
let reg = self.registry.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
reg.get_scalar(name).cloned()
};
let func = match func {
Some(f) => f,
None => return Err(format!("Unknown function: '{}'", name).into()),
};
if matches!(func, ScalarFunction::SequenceOp { .. }) {
let legacy = self.evaluate_function_call(name, args, chunk)?;
return Ok(ArrowVector::from_legacy(&legacy));
}
let mut first_error: Option<String> = None;
let mut row_results: Vec<Value> = Vec::with_capacity(num_rows);
for row in 0..num_rows {
let arg_values: Vec<Value> = arg_arrows
.iter()
.map(|arr| {
if row < arr.size() && !arr.is_null(row) {
arr.get_value(row).unwrap_or(Value::Null)
} else {
Value::Null
}
})
.collect();
if !name.eq_ignore_ascii_case("coalesce")
&& !name.eq_ignore_ascii_case("ifnull")
&& arg_values.iter().any(|v| matches!(v, Value::Null))
{
row_results.push(Value::Null);
continue;
}
match evaluate_scalar(&func, &arg_values) {
Ok(val) => row_results.push(val),
Err(e) => {
row_results.push(Value::Null);
if first_error.is_none() {
first_error = Some(e);
}
}
}
}
let result_type = row_results
.iter()
.find(|v| !matches!(v, Value::Null))
.map(|v| v.physical_type())
.unwrap_or(PhysicalTypeID::Int64);
let arrow_result = build_arrow_from_values(&row_results, result_type, num_rows)?;
if row_results.iter().all(|v| matches!(v, Value::Null))
&& let Some(e) = first_error
{
return Err(e.into());
}
Ok(arrow_result)
}
fn evaluate_constant(&self, c: &Constant, size: usize) -> Result<ValueVector, ProcessorError> {
let val: Value = match c {
Constant::Null => Value::Null,
Constant::Bool(b) => Value::Bool(*b),
Constant::Integer(i) => Value::Int64(*i),
Constant::Float(f) => Value::Double(*f),
Constant::String(s) => Value::String(s.clone()),
};
let physical_type = val.physical_type();
let mut v = ValueVector::new(physical_type, size);
v.resize(size);
for i in 0..size {
store_value_in_vector(&mut v, i, &val)?;
}
Ok(v)
}
fn evaluate_variable(&self, name: &str, chunk: &DataChunk) -> Result<ValueVector, ProcessorError> {
if let Ok(idx) = name.parse::<usize>() {
if idx >= chunk.fields.len() {
return Err(format!(
"Variable '{}' (index {}) not found in chunk with {} fields",
name,
idx,
chunk.fields.len()
)
.into());
}
let phys_type = chunk.field_types[idx];
let mut v = ValueVector::new(phys_type, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
if let Some(val) = chunk.get_value(idx, i) {
store_value_in_vector(&mut v, i, &val)?;
} else {
v.set_null(i, true);
}
}
return Ok(v);
}
if !chunk.field_names.is_empty() {
if let Some(idx) = chunk.field_names.iter().position(|n| n == name) {
let phys_type = chunk.field_types[idx];
let mut v = ValueVector::new(phys_type, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
if let Some(val) = chunk.get_value(idx, i) {
store_value_in_vector(&mut v, i, &val)?;
} else {
v.set_null(i, true);
}
}
return Ok(v);
}
return Err(format!(
"Variable '{}' not found in chunk field_names {:?}",
name, chunk.field_names
)
.into());
}
if !chunk.fields.is_empty() {
let phys_type = chunk.field_types[0];
let mut v = ValueVector::new(phys_type, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
if let Some(val) = chunk.get_value(0, i) {
store_value_in_vector(&mut v, i, &val)?;
} else {
v.set_null(i, true);
}
}
Ok(v)
} else {
Ok(ValueVector::new(akar_common::types::PhysicalTypeID::Int64, 0))
}
}
fn evaluate_property_access(
&self,
obj: &Expression,
prop: &str,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let qualified_prop = if let Expression::Variable(var_name) = obj {
format!("{}.{}", var_name, prop)
} else {
prop.to_string()
};
if !chunk.field_names.is_empty()
&& let Some(idx) = chunk.field_names.iter().position(|n| n == &qualified_prop)
{
if chunk.fields.get(idx).is_none() {
return Err(format!("Column '{}' (index {}) not found in chunk", prop, idx).into());
}
let phys_type = chunk.field_types[idx];
let mut v = ValueVector::new(phys_type, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
if let Some(val) = chunk.get_value(idx, i) {
store_value_in_vector(&mut v, i, &val)?;
} else {
v.set_null(i, true);
}
}
return Ok(v);
}
if let Expression::Variable(var_name) = obj
&& let Some(idx) = chunk.field_names.iter().position(|n| n == var_name)
{
let extracted: Vec<Value> = (0..chunk.size)
.map(|i| map_property_value(&chunk.get_value(idx, i).unwrap_or(Value::Null), prop))
.collect();
if extracted.iter().any(|v| !matches!(v, Value::Null)) {
let result_type = extracted
.iter()
.find(|v| !matches!(v, Value::Null))
.map(|v| v.physical_type())
.unwrap_or(PhysicalTypeID::Int64);
let mut v = ValueVector::new(result_type, chunk.size);
v.resize(chunk.size);
for (i, val) in extracted.iter().enumerate() {
if matches!(val, Value::Null) {
v.set_null(i, true);
} else {
store_value_in_vector(&mut v, i, val)?;
}
}
return Ok(v);
}
}
if !chunk.field_names.is_empty()
&& let Some(idx) = chunk.field_names.iter().position(|n| n == prop)
{
if chunk.fields.get(idx).is_none() {
return Err(format!("Column '{}' (index {}) not found in chunk", prop, idx).into());
}
let phys_type = chunk.field_types[idx];
let mut v = ValueVector::new(phys_type, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
if let Some(val) = chunk.get_value(idx, i) {
store_value_in_vector(&mut v, i, &val)?;
} else {
v.set_null(i, true);
}
}
return Ok(v);
}
self.evaluate(obj, chunk)
}
fn evaluate_function_call(
&self,
name: &str,
args: &[&Expression],
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
if let Some(lambda) = self.extract_lambda_arg(args) {
match name {
"list_transform" => return self.evaluate_list_transform(args, lambda, chunk),
"list_filter" => return self.evaluate_list_filter(args, lambda, chunk),
"list_reduce" => return self.evaluate_list_reduce(args, lambda, chunk),
_ => {}
}
}
let arg_vectors: Vec<ValueVector> = args
.iter()
.map(|arg| self.evaluate(arg, chunk))
.collect::<Result<Vec<_>, _>>()?;
if arg_vectors.is_empty() {
return Err(format!("Function '{}' requires at least one argument", name).into());
}
let num_rows = arg_vectors[0].size();
let func = {
let reg = self.registry.lock().map_err(|e| format!("Lock poisoned: {e}"))?;
reg.get_scalar(name).cloned()
};
let func = match func {
Some(f) => f,
None => return Err(format!("Unknown function: '{}'", name).into()),
};
if matches!(func, ScalarFunction::SequenceOp { .. }) {
return self.evaluate_sequence_op(name, &func, &arg_vectors, num_rows);
}
let mut row_results: Vec<Value> = Vec::with_capacity(num_rows);
let mut first_error: Option<String> = None;
for row in 0..num_rows {
let arg_values: Vec<Value> = arg_vectors
.iter()
.map(|vec| {
if row < vec.size() && !vec.is_null(row) {
vec.get_value(row).unwrap_or(Value::Null)
} else {
Value::Null
}
})
.collect();
if name == "AND" || name == "OR" {
let l = &arg_values[0];
let r = &arg_values[1];
if name == "AND" {
if matches!(l, Value::Bool(false)) || matches!(r, Value::Bool(false)) {
row_results.push(Value::Bool(false));
continue;
}
if matches!(l, Value::Null) || matches!(r, Value::Null) {
row_results.push(Value::Null);
continue;
}
} else {
if matches!(l, Value::Bool(true)) || matches!(r, Value::Bool(true)) {
row_results.push(Value::Bool(true));
continue;
}
if matches!(l, Value::Null) || matches!(r, Value::Null) {
row_results.push(Value::Null);
continue;
}
}
}
if !name.eq_ignore_ascii_case("coalesce")
&& !name.eq_ignore_ascii_case("ifnull")
&& arg_values.iter().any(|v| matches!(v, Value::Null))
{
row_results.push(Value::Null);
continue;
}
match evaluate_scalar(&func, &arg_values) {
Ok(val) => row_results.push(val),
Err(e) => {
row_results.push(Value::Null);
if first_error.is_none() {
first_error = Some(e);
}
}
}
}
let result_type = row_results
.iter()
.find(|v| !matches!(v, Value::Null))
.map(|v| v.physical_type())
.unwrap_or(akar_common::types::PhysicalTypeID::Int64);
let mut result_vec = ValueVector::new(result_type, num_rows);
result_vec.resize(num_rows);
for (row, val) in row_results.iter().enumerate() {
store_value_in_vector(&mut result_vec, row, val)?;
}
if row_results.iter().all(|v| matches!(v, Value::Null))
&& let Some(e) = first_error
{
return Err(e.into());
}
Ok(result_vec)
}
fn evaluate_binary_op(
&self,
op: &BinaryOp,
left: &Expression,
right: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let func_name = match op {
BinaryOp::Add => "+",
BinaryOp::Subtract => "-",
BinaryOp::Multiply => "*",
BinaryOp::Divide => "/",
BinaryOp::Modulo => "%",
BinaryOp::Equal => "=",
BinaryOp::NotEqual => "<>",
BinaryOp::LessThan => "<",
BinaryOp::LessThanOrEqual => "<=",
BinaryOp::GreaterThan => ">",
BinaryOp::GreaterThanOrEqual => ">=",
BinaryOp::And => "AND",
BinaryOp::Or => "OR",
BinaryOp::Xor => "XOR",
BinaryOp::Concat => "concat",
BinaryOp::In | BinaryOp::NotIn => {
return self.evaluate_in_op(op, left, right, chunk);
}
BinaryOp::StartsWith => "starts_with",
BinaryOp::EndsWith => "ends_with",
BinaryOp::Contains => "contains",
BinaryOp::Like => {
return self.evaluate_function_call("like", &[left, right], chunk);
}
};
self.evaluate_function_call(func_name, &[left, right], chunk)
}
fn evaluate_unary_op(
&self,
op: &UnaryOp,
inner: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
match op {
UnaryOp::Not => self.evaluate_function_call("NOT", std::slice::from_ref(&inner), chunk),
UnaryOp::Negate => self.evaluate_function_call("-", std::slice::from_ref(&inner), chunk),
UnaryOp::IsNull => {
let vec = self.evaluate(inner, chunk)?;
let num_rows = vec.size();
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::Bool, num_rows);
result.resize(num_rows);
for i in 0..num_rows {
let is_null = vec.is_null(i) || matches!(vec.get_value(i), Some(Value::Null) | None);
store_value_in_vector_simple(&mut result, i, &Value::Bool(is_null))?;
}
Ok(result)
}
UnaryOp::IsNotNull => {
let vec = self.evaluate(inner, chunk)?;
let num_rows = vec.size();
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::Bool, num_rows);
result.resize(num_rows);
for i in 0..num_rows {
let is_null = vec.is_null(i) || matches!(vec.get_value(i), Some(Value::Null) | None);
store_value_in_vector_simple(&mut result, i, &Value::Bool(!is_null))?;
}
Ok(result)
}
}
}
fn evaluate_in_op(
&self,
op: &BinaryOp,
left: &Expression,
right: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let left_arr = self.evaluate_arrow(left, chunk)?;
let num_rows = chunk.size;
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::Bool, num_rows);
result.resize(num_rows);
for row in 0..num_rows {
let lv = left_arr.get_value(row).unwrap_or(Value::Null);
if matches!(lv, Value::Null) {
result.set_null(row, true);
continue;
}
let (in_list, has_null) = match right {
Expression::List(items) => {
let mut matched = false;
let mut has_null_item = false;
for item in items {
let item_vec = self.evaluate(item, chunk)?;
let iv = item_vec.get_value(row).unwrap_or(Value::Null);
if matches!(iv, Value::Null) {
has_null_item = true;
} else if iv == lv {
matched = true;
break;
}
}
(matched, has_null_item)
}
_ => {
let right_arr = self.evaluate_arrow(right, chunk)?;
let rv = right_arr.get_value(row).unwrap_or(Value::Null);
match &rv {
Value::List(ritems) => {
let matched = ritems.contains(&lv);
let has_null_item = if matched {
false
} else {
ritems.iter().any(|item| matches!(item, Value::Null))
};
(matched, has_null_item)
}
_ => (lv == rv, false),
}
}
};
let result_val = if *op == BinaryOp::NotIn { !in_list } else { in_list };
if !in_list && has_null {
result.set_null(row, true);
} else {
store_value_in_vector_simple(&mut result, row, &Value::Bool(result_val))?;
}
}
Ok(result)
}
fn evaluate_case(
&self,
case_expr: &akar_parser::ast::CaseExpr,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let num_rows = chunk.size;
let subject_vec = if let Some(subj) = &case_expr.subject {
Some(self.evaluate(subj, chunk)?)
} else {
None
};
let result_type = {
let first_then = self.evaluate(&case_expr.alternatives[0].then, chunk)?;
first_then.physical_type()
};
let mut result = ValueVector::new(result_type, num_rows);
result.resize(num_rows);
for row in 0..num_rows {
let subject_val = subject_vec.as_ref().and_then(|sv| sv.get_value(row));
let mut matched = false;
for alt in &case_expr.alternatives {
let when_vec = self.evaluate(&alt.when, chunk)?;
let when_val = when_vec.get_value(row).unwrap_or(Value::Null);
let branch_taken = if let Some(ref sv) = subject_val {
when_val != Value::Null && when_val == *sv
} else {
matches!(when_val, Value::Bool(true))
};
if branch_taken {
let then_vec = self.evaluate(&alt.then, chunk)?;
let then_val = then_vec.get_value(row).unwrap_or(Value::Null);
store_value_in_vector(&mut result, row, &then_val)?;
matched = true;
break;
}
}
if !matched {
if let Some(else_e) = &case_expr.else_expr {
let else_vec = self.evaluate(else_e, chunk)?;
let else_val = else_vec.get_value(row).unwrap_or(Value::Null);
store_value_in_vector(&mut result, row, &else_val)?;
} else {
result.set_null(row, true);
}
}
}
Ok(result)
}
fn evaluate_list_literal(&self, items: &[Expression], chunk: &DataChunk) -> Result<ValueVector, ProcessorError> {
if items.is_empty() {
let mut v = ValueVector::new(akar_common::types::PhysicalTypeID::List, chunk.size);
v.resize(chunk.size);
for i in 0..chunk.size {
store_value_in_vector(&mut v, i, &Value::List(vec![]))?;
}
return Ok(v);
}
let num_rows = chunk.size;
let mut result_vec = ValueVector::new(akar_common::types::PhysicalTypeID::List, num_rows);
result_vec.resize(num_rows);
for row in 0..num_rows {
let mut list_values = Vec::with_capacity(items.len());
for item in items {
let item_vec = self.evaluate(item, chunk)?;
let val = if row < item_vec.size() {
item_vec.get_value(row).unwrap_or(Value::Null)
} else {
Value::Null
};
list_values.push(val);
}
store_value_in_vector(&mut result_vec, row, &Value::List(list_values))?;
}
Ok(result_vec)
}
fn evaluate_map_literal(
&self,
items: &[(String, Expression)],
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let num_rows = chunk.size;
let mut result_vec = ValueVector::new(akar_common::types::PhysicalTypeID::Struct, num_rows);
result_vec.resize(num_rows);
for row in 0..num_rows {
let mut map_values = Vec::with_capacity(items.len());
for (key, item) in items {
let item_vec = self.evaluate(item, chunk)?;
let val = if row < item_vec.size() {
item_vec.get_value(row).unwrap_or(Value::Null)
} else {
Value::Null
};
map_values.push((Value::String(key.clone()), val));
}
store_value_in_vector(&mut result_vec, row, &Value::Map(map_values))?;
}
Ok(result_vec)
}
fn evaluate_list_predicate(
&self,
quantifier: &akar_parser::ast::Quantifier,
list: &Expression,
_var_name: &str,
predicate: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let list_vec = self.evaluate(list, chunk)?;
let num_rows = chunk.size;
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::Bool, num_rows);
result.resize(num_rows);
for row in 0..num_rows {
let list_val = list_vec.get_value(row).unwrap_or(Value::Null);
let items = match &list_val {
Value::List(items) => items.as_slice(),
_ => {
store_value_in_vector(&mut result, row, &Value::Bool(false))?;
continue;
}
};
let mut true_count = 0u64;
for item in items {
let mut elem_vec = ValueVector::new(item.physical_type(), 1);
elem_vec.resize(1);
store_value_in_vector(&mut elem_vec, 0, item)?;
let mini_chunk = {
let arrow_fields = vec![&elem_vec]
.into_iter()
.map(|v| akar_common::arrow_vector::ArrowVector::from_legacy(v).array)
.collect::<Vec<_>>();
let arrow_field_types = vec![&elem_vec]
.into_iter()
.map(|v| v.physical_type())
.collect::<Vec<_>>();
DataChunk::new(arrow_fields, arrow_field_types)
};
let pred_vec = self.evaluate(predicate, &mini_chunk)?;
let pred_val = pred_vec.get_value(0).unwrap_or(Value::Null);
if matches!(pred_val, Value::Bool(true)) {
true_count += 1;
}
}
let elem_count = items.len() as u64;
let bool_result = match quantifier {
akar_parser::ast::Quantifier::Any => true_count > 0,
akar_parser::ast::Quantifier::All => !items.is_empty() && true_count == elem_count,
akar_parser::ast::Quantifier::None => true_count == 0,
akar_parser::ast::Quantifier::Single => true_count == 1,
};
store_value_in_vector(&mut result, row, &Value::Bool(bool_result))?;
}
Ok(result)
}
fn extract_lambda_arg<'a>(&self, args: &'a [&Expression]) -> Option<&'a Expression> {
args.iter().find(|a| matches!(a, Expression::Lambda { .. })).copied()
}
fn evaluate_list_transform(
&self,
args: &[&Expression],
lambda: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let list_expr = args
.iter()
.find(|a| !matches!(a, Expression::Lambda { .. }))
.ok_or("list_transform requires a list argument")?;
let (var_name, body) = match lambda {
Expression::Lambda { var_name, body } => (var_name, body),
_ => return Err("Expected lambda expression".into()),
};
let list_vec = self.evaluate(list_expr, chunk)?;
let num_rows = chunk.size;
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::List, num_rows);
result.resize(num_rows);
for row in 0..num_rows {
let list_val = list_vec.get_value(row).unwrap_or(Value::Null);
let items = match list_val {
Value::List(items) => items,
_ => {
store_value_in_vector(&mut result, row, &Value::List(vec![]))?;
continue;
}
};
let mut transformed: Vec<Value> = Vec::with_capacity(items.len());
for item in items {
let mut elem_vec = ValueVector::new(item.physical_type(), 1);
elem_vec.resize(1);
store_value_in_vector(&mut elem_vec, 0, &item)?;
let mut mini_chunk = {
let arrow_fields = vec![&elem_vec]
.into_iter()
.map(|v| akar_common::arrow_vector::ArrowVector::from_legacy(v).array)
.collect::<Vec<_>>();
let arrow_field_types = vec![&elem_vec]
.into_iter()
.map(|v| v.physical_type())
.collect::<Vec<_>>();
DataChunk::new(arrow_fields, arrow_field_types)
};
mini_chunk.field_names.push(var_name.clone());
let body_vec = self.evaluate(body, &mini_chunk)?;
let body_val = body_vec.get_value(0).unwrap_or(Value::Null);
transformed.push(body_val);
}
store_value_in_vector(&mut result, row, &Value::List(transformed))?;
}
Ok(result)
}
fn evaluate_list_filter(
&self,
args: &[&Expression],
lambda: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let list_expr = args
.iter()
.find(|a| !matches!(a, Expression::Lambda { .. }))
.ok_or("list_filter requires a list argument")?;
let (var_name, body) = match lambda {
Expression::Lambda { var_name, body } => (var_name, body),
_ => return Err("Expected lambda expression".into()),
};
let list_vec = self.evaluate(list_expr, chunk)?;
let num_rows = chunk.size;
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::List, num_rows);
result.resize(num_rows);
for row in 0..num_rows {
let list_val = list_vec.get_value(row).unwrap_or(Value::Null);
let items = match list_val {
Value::List(items) => items,
_ => {
store_value_in_vector(&mut result, row, &Value::List(vec![]))?;
continue;
}
};
let mut filtered: Vec<Value> = Vec::with_capacity(items.len());
for item in items {
let mut elem_vec = ValueVector::new(item.physical_type(), 1);
elem_vec.resize(1);
store_value_in_vector(&mut elem_vec, 0, &item)?;
let mut mini_chunk = {
let arrow_fields = vec![&elem_vec]
.into_iter()
.map(|v| akar_common::arrow_vector::ArrowVector::from_legacy(v).array)
.collect::<Vec<_>>();
let arrow_field_types = vec![&elem_vec]
.into_iter()
.map(|v| v.physical_type())
.collect::<Vec<_>>();
DataChunk::new(arrow_fields, arrow_field_types)
};
mini_chunk.field_names.push(var_name.clone());
let pred_vec = self.evaluate(body, &mini_chunk)?;
let pred_val = pred_vec.get_value(0).unwrap_or(Value::Null);
if matches!(pred_val, Value::Bool(true)) {
filtered.push(item);
} else if let Value::Int64(x) = pred_val {
if x != 0 {
filtered.push(item);
}
}
}
store_value_in_vector(&mut result, row, &Value::List(filtered))?;
}
Ok(result)
}
fn evaluate_list_reduce(
&self,
args: &[&Expression],
lambda: &Expression,
chunk: &DataChunk,
) -> Result<ValueVector, ProcessorError> {
let list_expr = args
.iter()
.find(|a| !matches!(a, Expression::Lambda { .. }))
.ok_or("list_reduce requires a list argument")?;
let initial_expr = args
.iter()
.filter(|a| !matches!(a, Expression::Lambda { .. }))
.nth(1) .ok_or("list_reduce requires an initial value argument")?;
let (var_name, body) = match lambda {
Expression::Lambda { var_name, body } => (var_name, body),
_ => return Err("Expected lambda expression".into()),
};
let acc_name = var_name.clone();
let elem_name = match body.as_ref() {
Expression::BinaryOp(_op, left, right) => {
let left_var = if let Expression::Variable(v) = left.as_ref() {
Some(v.clone())
} else {
None
};
let right_var = if let Expression::Variable(v) = right.as_ref() {
Some(v.clone())
} else {
None
};
if left_var.as_deref() == Some(&acc_name) {
right_var.unwrap_or_default()
} else if right_var.as_deref() == Some(&acc_name) {
left_var.unwrap_or_default()
} else {
String::new()
}
}
_ => String::new(),
};
let list_vec = self.evaluate(list_expr, chunk)?;
let initial_vec = self.evaluate(initial_expr, chunk)?;
let num_rows = chunk.size;
let mut result = ValueVector::new(akar_common::types::PhysicalTypeID::Int64, num_rows);
result.resize(num_rows);
let field_names_template = {
let mut names = vec![acc_name.clone()];
if !elem_name.is_empty() {
names.push(elem_name.clone());
}
names
};
for row in 0..num_rows {
let list_val = list_vec.get_value(row).unwrap_or(Value::Null);
let items = match list_val {
Value::List(items) => items,
_ => {
store_value_in_vector(&mut result, row, &Value::Null)?;
continue;
}
};
let mut acc = initial_vec.get_value(row).unwrap_or(Value::Null);
for item in items {
let mut acc_vec = ValueVector::new(acc.physical_type(), 1);
acc_vec.resize(1);
store_value_in_vector(&mut acc_vec, 0, &acc)?;
let mut elem_vec = ValueVector::new(item.physical_type(), 1);
elem_vec.resize(1);
store_value_in_vector(&mut elem_vec, 0, &item)?;
let mut mini_chunk = {
let arrow_fields = vec![&acc_vec, &elem_vec]
.into_iter()
.map(|v| akar_common::arrow_vector::ArrowVector::from_legacy(v).array)
.collect::<Vec<_>>();
let arrow_field_types = vec![&acc_vec, &elem_vec]
.into_iter()
.map(|v| v.physical_type())
.collect::<Vec<_>>();
DataChunk::new(arrow_fields, arrow_field_types)
};
mini_chunk.field_names = field_names_template.clone();
let body_vec = self.evaluate(body, &mini_chunk)?;
acc = body_vec.get_value(0).unwrap_or(Value::Null);
}
store_value_in_vector(&mut result, row, &acc)?;
}
Ok(result)
}
}
fn store_value_in_vector_simple(v: &mut ValueVector, row: usize, val: &Value) -> Result<(), String> {
match val {
Value::Null => {
v.set_null(row, true);
}
Value::Bool(x) => {
if v.physical_type() == akar_common::types::PhysicalTypeID::Bool {
v.data_mut()[row] = if *x { 1 } else { 0 };
v.set_null(row, false);
}
}
Value::Int64(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::UInt64(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::Double(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::String(s) => {
let bytes = s.as_bytes();
if bytes.len() > 255 {
return Err(format!(
"Cannot store string of {} bytes: inline string storage limit is 255 bytes",
bytes.len()
));
}
let offset = row * 256;
if offset < v.data().len() {
v.data_mut()[offset] = bytes.len() as u8;
if offset + 1 + bytes.len() <= v.data().len() {
v.data_mut()[offset + 1..offset + 1 + bytes.len()].copy_from_slice(bytes);
}
v.set_null(row, false);
}
}
_ => {
v.set_null(row, true);
}
}
Ok(())
}
fn store_value_in_vector(v: &mut ValueVector, row: usize, val: &Value) -> Result<(), String> {
match val {
Value::Null => {
v.set_null(row, true);
}
Value::Bool(x) => {
if v.physical_type() == akar_common::types::PhysicalTypeID::Bool {
v.data_mut()[row] = if *x { 1 } else { 0 };
v.set_null(row, false);
}
}
Value::Int64(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::UInt64(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::Date(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&(x.0 as i64).to_le_bytes());
v.set_null(row, false);
}
}
Value::Timestamp(x) | Value::TimestampNs(x) | Value::TimestampMs(x) | Value::TimestampSec(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.0.to_le_bytes());
v.set_null(row, false);
}
}
Value::TimestampTz(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.0.to_le_bytes());
v.set_null(row, false);
}
}
Value::DTime(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::Int32(x) => {
let offset = row * 4;
if offset + 4 <= v.data().len() {
v.data_mut()[offset..offset + 4].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::Double(x) => {
let offset = row * 8;
if offset + 8 <= v.data().len() {
v.data_mut()[offset..offset + 8].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::Float(x) => {
let offset = row * 4;
if offset + 4 <= v.data().len() {
v.data_mut()[offset..offset + 4].copy_from_slice(&x.to_le_bytes());
v.set_null(row, false);
}
}
Value::String(s) => {
let bytes = s.as_bytes();
if bytes.len() > 255 {
return Err(format!(
"Cannot store string of {} bytes: inline string storage limit is 255 bytes",
bytes.len()
));
}
let offset = row * 256;
if offset < v.data().len() {
v.data_mut()[offset] = bytes.len() as u8;
if offset + 1 + bytes.len() <= v.data().len() {
v.data_mut()[offset + 1..offset + 1 + bytes.len()].copy_from_slice(bytes);
}
v.set_null(row, false);
}
}
_ => {
v.set_null(row, true);
}
}
Ok(())
}
pub(crate) fn build_arrow_from_values(
values: &[Value],
phys_type: PhysicalTypeID,
num_rows: usize,
) -> Result<ArrowVector, ProcessorError> {
match phys_type {
PhysicalTypeID::Bool => {
let mut builder = arrow::array::BooleanBuilder::with_capacity(num_rows);
for v in values {
match v {
Value::Null => builder.append_null(),
Value::Bool(b) => builder.append_value(*b),
_ => builder.append_null(),
}
}
Ok(ArrowVector::new(Arc::new(builder.finish()), phys_type))
}
PhysicalTypeID::Int64 => {
let mut builder = arrow::array::Int64Builder::with_capacity(num_rows);
for v in values {
match v {
Value::Null => builder.append_null(),
Value::Int64(n) => builder.append_value(*n),
Value::Int32(n) => builder.append_value(*n as i64),
Value::Date(n) => builder.append_value(n.0 as i64),
Value::Timestamp(n) | Value::TimestampNs(n) | Value::TimestampMs(n) | Value::TimestampSec(n) => {
builder.append_value(n.0)
}
Value::TimestampTz(n) => builder.append_value(n.0),
Value::DTime(n) => builder.append_value(*n),
_ => builder.append_null(),
}
}
Ok(ArrowVector::new(Arc::new(builder.finish()), phys_type))
}
PhysicalTypeID::Int32 => {
let mut builder = arrow::array::Int32Builder::with_capacity(num_rows);
for v in values {
match v {
Value::Null => builder.append_null(),
Value::Int32(n) => builder.append_value(*n),
_ => builder.append_null(),
}
}
Ok(ArrowVector::new(Arc::new(builder.finish()), phys_type))
}
PhysicalTypeID::Double => {
let mut builder = arrow::array::Float64Builder::with_capacity(num_rows);
for v in values {
match v {
Value::Null => builder.append_null(),
Value::Double(n) => builder.append_value(*n),
_ => builder.append_null(),
}
}
Ok(ArrowVector::new(Arc::new(builder.finish()), phys_type))
}
PhysicalTypeID::Float => {
let mut builder = arrow::array::Float32Builder::with_capacity(num_rows);
for v in values {
match v {
Value::Null => builder.append_null(),
Value::Float(n) => builder.append_value(*n),
Value::Double(n) => builder.append_value(*n as f32),
_ => builder.append_null(),
}
}
Ok(ArrowVector::new(Arc::new(builder.finish()), phys_type))
}
PhysicalTypeID::String => {
let mut builder = arrow::array::StringBuilder::with_capacity(num_rows, num_rows * 16);
for v in values {
match v {
Value::Null => builder.append_null(),
Value::String(s) => builder.append_value(s),
_ => builder.append_null(),
}
}
Ok(ArrowVector::new(Arc::new(builder.finish()), phys_type))
}
PhysicalTypeID::List | PhysicalTypeID::Array | PhysicalTypeID::Struct => {
let array = akar_common::arrow_vector::arrow_array_from_values(values);
Ok(ArrowVector::new(array, phys_type))
}
_ => {
let mut builder = arrow::array::Int64Builder::with_capacity(num_rows);
builder.append_nulls(num_rows);
Ok(ArrowVector::new(Arc::new(builder.finish()), PhysicalTypeID::Int64))
}
}
}
impl ExpressionEvaluator {
fn evaluate_sequence_op(
&self,
name: &str,
func: &ScalarFunction,
arg_vectors: &[ValueVector],
num_rows: usize,
) -> Result<ValueVector, ProcessorError> {
let is_nextval = match func {
ScalarFunction::SequenceOp { is_nextval } => *is_nextval,
_ => return Err(format!("Internal error: expected SequenceOp for '{}'", name).into()),
};
let seq_fn = self
.sequence_fn
.as_ref()
.ok_or_else(|| format!("No sequence callback configured for '{}'", name))?;
let mut result_vec = ValueVector::new(akar_common::types::PhysicalTypeID::Int64, num_rows);
result_vec.resize(num_rows);
for row in 0..num_rows {
let seq_name = if row < arg_vectors[0].size() && !arg_vectors[0].is_null(row) {
match arg_vectors[0].get_value(row) {
Some(Value::String(s)) => s,
_ => return Err("nextval/currval requires a string argument (sequence name)".into()),
}
} else {
result_vec.set_null(row, true);
continue;
};
match seq_fn(&seq_name, is_nextval) {
Ok(val) => {
store_value_in_vector(&mut result_vec, row, &val)?;
}
Err(e) => {
result_vec.set_null(row, true);
if row == 0 {
return Err(e);
}
}
}
}
Ok(result_vec)
}
}
#[cfg(test)]
mod tests {
use super::*;
use akar_common::types::PhysicalTypeID;
use akar_common::vector::ValueVector;
use akar_function::registry::FunctionRegistry;
use hashbrown::HashMap;
fn make_registry() -> Arc<Mutex<FunctionRegistry>> {
Arc::new(Mutex::new(FunctionRegistry::new()))
}
fn make_chunk(values: &[i64]) -> DataChunk {
let mut v = ValueVector::new(PhysicalTypeID::Int64, values.len());
v.resize(values.len());
for (i, val) in values.iter().enumerate() {
v.set_i64(i, *val);
}
{
let arrow_fields = vec![akar_common::arrow_vector::ArrowVector::from_legacy(&v).array];
let arrow_field_types = vec![v.physical_type()];
DataChunk::new(arrow_fields, arrow_field_types)
}
}
#[test]
fn test_evaluate_constant_int() {
let eval = ExpressionEvaluator::new(make_registry());
let expr = Expression::Constant(Constant::Integer(42));
let chunk = make_chunk(&[]);
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 0);
}
#[test]
fn test_evaluate_constant_bool() {
let eval = ExpressionEvaluator::new(make_registry());
let expr = Expression::Constant(Constant::Bool(true));
let chunk = make_chunk(&[1, 2, 3]);
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 3);
for i in 0..3 {
assert!(!result.is_null(i));
}
}
#[test]
fn test_evaluate_variable() {
let eval = ExpressionEvaluator::new(make_registry());
let chunk = make_chunk(&[10, 20, 30]);
let expr = Expression::Variable("0".into());
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 3);
assert_eq!(result.get_i64(0), Some(10));
assert_eq!(result.get_i64(1), Some(20));
assert_eq!(result.get_i64(2), Some(30));
}
#[test]
fn test_evaluate_binary_equal() {
let eval = ExpressionEvaluator::new(make_registry());
let left = Box::new(Expression::Variable("0".into()));
let right = Box::new(Expression::Constant(Constant::Integer(0)));
let expr = Expression::BinaryOp(BinaryOp::Equal, left, right);
let chunk = make_chunk(&[0, 1, 2]);
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 3);
assert_eq!(result.get_value(0), Some(Value::Bool(true)));
assert_eq!(result.get_value(1), Some(Value::Bool(false)));
assert_eq!(result.get_value(2), Some(Value::Bool(false)));
}
#[test]
fn test_evaluate_binary_greater_than() {
let eval = ExpressionEvaluator::new(make_registry());
let left = Box::new(Expression::Variable("0".into()));
let right = Box::new(Expression::Constant(Constant::Integer(2)));
let expr = Expression::BinaryOp(BinaryOp::GreaterThan, left, right);
let chunk = make_chunk(&[3, 1, 5]);
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 3);
assert_eq!(result.get_value(0), Some(Value::Bool(true)));
assert_eq!(result.get_value(1), Some(Value::Bool(false)));
assert_eq!(result.get_value(2), Some(Value::Bool(true)));
}
#[test]
fn test_evaluate_binary_and() {
let eval = ExpressionEvaluator::new(make_registry());
let mut v0 = ValueVector::new(PhysicalTypeID::Bool, 2);
v0.resize(2);
store_value_in_vector(&mut v0, 0, &Value::Bool(true)).unwrap();
store_value_in_vector(&mut v0, 1, &Value::Bool(false)).unwrap();
let mut v1 = ValueVector::new(PhysicalTypeID::Bool, 2);
v1.resize(2);
store_value_in_vector(&mut v1, 0, &Value::Bool(true)).unwrap();
store_value_in_vector(&mut v1, 1, &Value::Bool(true)).unwrap();
let chunk = {
let arrow_fields = vec![
akar_common::arrow_vector::ArrowVector::from_legacy(&v0).array,
akar_common::arrow_vector::ArrowVector::from_legacy(&v1).array,
];
let arrow_field_types = vec![v0.physical_type(), v1.physical_type()];
DataChunk::new(arrow_fields, arrow_field_types)
};
let left = Box::new(Expression::Variable("0".into()));
let right = Box::new(Expression::Variable("1".into()));
let expr = Expression::BinaryOp(BinaryOp::And, left, right);
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 2);
assert_eq!(result.get_value(0), Some(Value::Bool(true)));
assert_eq!(result.get_value(1), Some(Value::Bool(false)));
}
#[test]
fn test_evaluate_function_call_string_length() {
let eval = ExpressionEvaluator::new(make_registry());
let chunk = akar_common::vector::DataChunk::new(vec![], vec![]);
let expr = Expression::FunctionCall(
"length".into(),
vec![Expression::Constant(Constant::String("hello".into()))],
);
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 0);
}
#[test]
fn test_evaluate_not() {
let eval = ExpressionEvaluator::new(make_registry());
let mut v = ValueVector::new(PhysicalTypeID::Bool, 3);
v.resize(3);
store_value_in_vector(&mut v, 0, &Value::Bool(true)).unwrap();
store_value_in_vector(&mut v, 1, &Value::Bool(false)).unwrap();
store_value_in_vector(&mut v, 2, &Value::Bool(true)).unwrap();
let chunk = {
let arrow_fields = vec![akar_common::arrow_vector::ArrowVector::from_legacy(&v).array];
let arrow_field_types = vec![v.physical_type()];
DataChunk::new(arrow_fields, arrow_field_types)
};
let expr = Expression::UnaryOp(UnaryOp::Not, Box::new(Expression::Variable("0".into())));
let result = eval.evaluate(&expr, &chunk).unwrap();
assert_eq!(result.size(), 3);
assert_eq!(result.get_value(0), Some(Value::Bool(false)));
assert_eq!(result.get_value(1), Some(Value::Bool(true)));
assert_eq!(result.get_value(2), Some(Value::Bool(false)));
}
#[test]
fn test_sequence_nextval_currval_with_callback() {
let state = Arc::new(Mutex::new(HashMap::new()));
state.lock().unwrap().insert("my_seq".to_string(), 10_i64);
let state_for_fn = state.clone();
let seq_fn: Arc<dyn Fn(&str, bool) -> Result<Value, ProcessorError> + Send + Sync> =
Arc::new(move |seq_name: &str, is_nextval: bool| {
let mut map = state_for_fn.lock().map_err(|e| format!("Lock error: {e}"))?;
let current = map
.get_mut(seq_name)
.ok_or_else(|| format!("Sequence '{}' not found", seq_name))?;
if is_nextval {
let out = *current;
*current += 2;
Ok(Value::Int64(out))
} else {
Ok(Value::Int64(*current))
}
});
let eval = ExpressionEvaluator::new(make_registry()).with_sequence_fn(seq_fn);
let chunk = make_chunk(&[1, 2, 3]);
let nextval_expr = Expression::FunctionCall(
"nextval".into(),
vec![Expression::Constant(Constant::String("my_seq".into()))],
);
let nextvals = eval.evaluate(&nextval_expr, &chunk).unwrap();
assert_eq!(nextvals.get_value(0), Some(Value::Int64(10)));
assert_eq!(nextvals.get_value(1), Some(Value::Int64(12)));
assert_eq!(nextvals.get_value(2), Some(Value::Int64(14)));
let currval_expr = Expression::FunctionCall(
"currval".into(),
vec![Expression::Constant(Constant::String("my_seq".into()))],
);
let curr = eval.evaluate(&currval_expr, &make_chunk(&[1])).unwrap();
assert_eq!(curr.get_value(0), Some(Value::Int64(16)));
}
#[test]
fn test_sequence_requires_callback() {
let eval = ExpressionEvaluator::new(make_registry());
let expr = Expression::FunctionCall(
"nextval".into(),
vec![Expression::Constant(Constant::String("my_seq".into()))],
);
let err = eval.evaluate(&expr, &make_chunk(&[1])).unwrap_err();
assert!(
err.to_string().contains("No sequence callback configured"),
"Unexpected error: {err}"
);
}
#[test]
fn test_sequence_requires_string_arg() {
let seq_fn: Arc<dyn Fn(&str, bool) -> Result<Value, ProcessorError> + Send + Sync> =
Arc::new(|_seq_name: &str, _is_nextval: bool| Ok(Value::Int64(1)));
let eval = ExpressionEvaluator::new(make_registry()).with_sequence_fn(seq_fn);
let expr = Expression::FunctionCall("nextval".into(), vec![Expression::Constant(Constant::Integer(42))]);
let err = eval.evaluate(&expr, &make_chunk(&[1])).unwrap_err();
assert!(
err.to_string().contains("requires a string argument"),
"Unexpected error: {err}"
);
}
}