Skip to main content

akar_processor/processor/mapper/
map_scan.rs

1use 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    // Try Arrow fast path first (avoids Vec<Vec<Value>> intermediate)
57    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        // Fallback to legacy Vec<Vec<Value>> data path (MVCC-aware)
62        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            // Resolve the HNSW index by column (first index on the table whose
129            // indexed column matches), falling back to the first index on the
130            // table — same convention as the explicit CALL handler.
131            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            // Prefix the output field names with the node alias (like the
161            // ScanNode branch) so downstream Filter/OrderBy/Projection, which
162            // reference `alias.col`, can resolve their columns. The bare CALL
163            // path has `alias: None` and keeps unprefixed names (distance, _id).
164            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}