akar_processor/processor/mapper/
map_scan.rs1use super::ExecutionContext;
2use crate::expression_evaluator::ExpressionEvaluator;
3use crate::physical_operator::*;
4use akar_common::error::ProcessorError;
5use akar_common::types::Value;
6use akar_common::vector::DataChunk;
7use akar_parser::ast::Expression;
8use akar_planner::logical_operator::{LogicalOperator, LogicalScanNode};
9use std::sync::{Arc, Mutex};
10
11fn extract_zone_map_predicate(expr: &Expression, columns: &[String]) -> Option<(usize, String, Value)> {
12 if let Expression::BinaryOp(op, left, right) = expr {
13 let op_str = match op {
14 akar_parser::ast::BinaryOp::Equal => "=",
15 akar_parser::ast::BinaryOp::GreaterThan => ">",
16 akar_parser::ast::BinaryOp::LessThan => "<",
17 akar_parser::ast::BinaryOp::GreaterThanOrEqual => ">=",
18 akar_parser::ast::BinaryOp::LessThanOrEqual => "<=",
19 akar_parser::ast::BinaryOp::NotEqual => "!=",
20 _ => return None,
21 };
22 if let Expression::Variable(var_name) = &**left {
23 if let Expression::Constant(c) = &**right {
24 let col_name = var_name.split('.').next_back().unwrap_or(var_name);
25 if let Some(col_idx) = columns.iter().position(|c| c == col_name) {
26 let val = match c {
27 akar_parser::ast::Constant::Integer(i) => Value::Int64(*i),
28 akar_parser::ast::Constant::Float(f) => Value::Double(*f),
29 akar_parser::ast::Constant::String(s) => Value::String(s.clone()),
30 akar_parser::ast::Constant::Bool(b) => Value::Bool(*b),
31 akar_parser::ast::Constant::Null => Value::Null,
32 };
33 return Some((col_idx, op_str.to_string(), val));
34 }
35 }
36 }
37 }
38 None
39}
40
41pub fn map_and_execute_scan_node(
42 s: &LogicalScanNode,
43 next_op: Option<&LogicalOperator>,
44 current_input: Vec<DataChunk>,
45 ctx: &mut ExecutionContext,
46) -> Result<Vec<DataChunk>, ProcessorError> {
47 let mut pred_owned = None;
48 if let Some(LogicalOperator::Filter(f)) = next_op {
49 pred_owned = extract_zone_map_predicate(&f.expression, &s.columns);
50 }
51
52 let pred_ref = pred_owned
53 .as_ref()
54 .map(|(idx, op_str, val)| (*idx, op_str.as_str(), val));
55
56 let (arrow_data, columns, arrow_num_rows) = ctx.resolve_scan_arrow_data(&s.table_name);
58 let mut scan = if let Some(arrays) = arrow_data {
59 PhysicalScan::new(s.table_name.clone(), s.table_id, arrow_num_rows.max(1)).with_arrow_data(arrays, columns)
60 } else {
61 let (data, fallback_columns, fallback_num_rows) = ctx.resolve_scan_data(&s.table_name, pred_ref);
63 let mut s = PhysicalScan::new(s.table_name.clone(), s.table_id, fallback_num_rows.max(1));
64 if let Some(d) = data {
65 s = s.with_data(d, fallback_columns);
66 }
67 s
68 };
69 if let Some(ref fq) = s.fts_query {
70 scan = scan.with_fts_query(PhysicalFtsScan {
71 index_name: fq.index_name.clone(),
72 query_string: fq.query_string.clone(),
73 docs_table: fq.docs_table.clone(),
74 terms_table: fq.terms_table.clone(),
75 posting_table: fq.posting_table.clone(),
76 table_name: fq.table_name.clone(),
77 column_name: fq.column_name.clone(),
78 table_catalog: ctx
79 .table_catalog
80 .clone()
81 .ok_or_else(|| "Table catalog required for FTS scan".to_string())?,
82 });
83 }
84 if let Some(ref pred) = s.predicate {
85 scan = scan.with_predicate(pred.clone());
86 let registry = ctx
87 .function_registry
88 .clone()
89 .ok_or_else(|| "No function registry available for predicate evaluation".to_string())?;
90 scan = scan.with_evaluator(Arc::new(Mutex::new(ExpressionEvaluator::new(registry))));
91 }
92 let mut result = scan.execute(current_input)?;
93 let prefix = s.alias.as_ref().unwrap_or(&s.table_name);
94
95 for chunk in &mut result {
96 chunk.field_names = chunk.field_names.iter().map(|n| format!("{}.{}", prefix, n)).collect();
97 }
98 Ok(result)
99}
100
101pub fn map_and_execute_scan(
102 op: &LogicalOperator,
103 current_input: Vec<DataChunk>,
104 ctx: &mut ExecutionContext,
105) -> Result<Vec<DataChunk>, ProcessorError> {
106 match op {
107 LogicalOperator::ScanRel(s) => {
108 let (data, columns, _num_rows) = ctx.resolve_scan_data(&s.table_name, None);
109 let scan = PhysicalScanRel {
110 table_name: s.table_name.clone(),
111 table_id: s.table_id,
112 direction: s.direction.clone(),
113 table_data: data,
114 table_columns: columns,
115 };
116 let mut result = scan.execute(current_input)?;
117 let prefix = &s.table_name;
118 for chunk in &mut result {
119 chunk.field_names = chunk.field_names.iter().map(|n| format!("{}.{}", prefix, n)).collect();
120 }
121 Ok(result)
122 }
123 LogicalOperator::VectorSimilarityScan(vs) => {
124 let tc = ctx
125 .table_catalog
126 .clone()
127 .ok_or_else(|| "No table catalog available for VECTOR SIMILARITY SCAN".to_string())?;
128 let index_name = {
132 let mut by_column = None;
133 let mut first_on_table = None;
134 for entry in tc.all_vector_indexes() {
135 if entry.table_name == vs.table_name {
136 if by_column.is_none() && entry.column_name == vs.column_name {
137 by_column = Some(entry.name.clone());
138 }
139 if first_on_table.is_none() {
140 first_on_table = Some(entry.name.clone());
141 }
142 }
143 }
144 by_column.or(first_on_table).ok_or_else(|| {
145 format!(
146 "No vector index found on table '{}' for column '{}'",
147 vs.table_name, vs.column_name
148 )
149 })?
150 };
151 let scan = PhysicalVectorSimilarityScan {
152 index_name,
153 index_id: 0,
154 query_vector: vs.query_vector.clone(),
155 top_k: vs.top_k,
156 table_name: vs.table_name.clone(),
157 table_catalog: Some(tc),
158 };
159 let mut result = scan.execute(current_input)?;
160 if let Some(alias) = &vs.alias {
165 for chunk in &mut result {
166 chunk.field_names = chunk.field_names.iter().map(|n| format!("{alias}.{n}")).collect();
167 }
168 }
169 Ok(result)
170 }
171 LogicalOperator::ArtIndexRangeScan(ars) => {
172 let scan = PhysicalArtIndexRangeScan {
173 table_name: ars.table_name.clone(),
174 table_id: ars.table_id,
175 lower_bound: ars.lower_bound.clone(),
176 upper_bound: ars.upper_bound.clone(),
177 lower_inclusive: ars.lower_inclusive,
178 upper_inclusive: ars.upper_inclusive,
179 table_catalog: ctx.table_catalog.clone(),
180 };
181 let mut result = scan.execute(current_input)?;
182 let prefix = ars.alias.as_ref().unwrap_or(&ars.table_name);
183 for chunk in &mut result {
184 chunk.field_names = chunk.field_names.iter().map(|n| format!("{}.{}", prefix, n)).collect();
185 }
186 Ok(result)
187 }
188 LogicalOperator::IndexLookup(il) => {
189 let table_catalog = ctx
190 .table_catalog
191 .clone()
192 .ok_or_else(|| "No table catalog available for INDEX LOOKUP".to_string())?;
193 let lookup_op = PhysicalIndexLookup {
194 table_name: il.table_name.clone(),
195 table_id: il.table_id,
196 key_value: il.key_value.clone(),
197 table_catalog,
198 };
199 let result = lookup_op.execute(current_input)?;
200 Ok(result)
201 }
202 LogicalOperator::ExpressionsScan(_es) => Ok(vec![DataChunk::new(vec![], vec![])]),
203 LogicalOperator::PathPropertyProbe(p) => {
204 let properties = p
205 .properties
206 .iter()
207 .map(|(t, is_node, props)| crate::physical::scan_filter::PathPropertySpec {
208 table_name: t.clone(),
209 is_node: *is_node,
210 property_names: props.clone(),
211 })
212 .collect();
213
214 let probe = crate::physical::scan_filter::PhysicalPathPropertyProbe {
215 node_ids_col_idx: p.node_ids_col_idx,
216 edge_ids_col_idx: p.edge_ids_col_idx,
217 properties,
218 table_catalog: ctx
219 .table_catalog
220 .clone()
221 .ok_or_else(|| "table catalog required for PathPropertyProbe".to_string())?,
222 };
223
224 let input = if !p.children.is_empty() {
225 ctx.execute_children(&p.children)?
226 } else {
227 current_input
228 };
229
230 let result = probe.execute(input)?;
231 Ok(result)
232 }
233 _ => Err(format!("Not a scan operator: {:?}", op).into()),
234 }
235}