1use std::sync::Arc;
16
17use rudb_catalog::{Catalog, Entry, FileStamp, QualifiedName, same_name};
18use rudb_common::bounds::Zones;
19use rudb_common::{
20 Error, Field, LogicalType, Result, Semantics, Session, ShowBehavior, Span, Stat, Value,
21};
22use rudb_functions::{
23 Columns, FILE_ROW_NUMBER, Footers, FunctionKind, Given, Resolved, TableFunction, csv_fields,
24 csv_given, files, is_file, is_pattern, kind_of, parquet_footers, parquet_outline, resolve,
25 resolve_pragma, resolve_table,
26};
27use rudb_kernels::{percentage, row_count};
28use rudb_parse::ast::{self, Ast, Distinct, LiteralKind, Nulls, Order, Quantifier, SetOp};
29use rudb_parse::{NONE, identifier_parts, parse_ast_with_case};
30use rudb_plan::{
31 Bound, BuildSide, ColumnBinding, ConjunctionOp, Expr, ExprRef, JoinKind, Node, NodeRef, Plan,
32 SetOpKind, Share, SortKey, WindowBound, WindowExclude, WindowFrame, WindowUnit,
33};
34
35use crate::expr::{describe, has_aggregate};
36use crate::fold;
37use crate::parameters::Parameters;
38use crate::scope::{Scope, Visible};
39
40pub fn bind(ast: &Ast, catalog: &Catalog) -> Result<Plan> {
47 bind_with(ast, catalog, &Parameters::new(), &Session::new())
48}
49
50pub fn bind_with(
59 ast: &Ast,
60 catalog: &Catalog,
61 parameters: &Parameters,
62 session: &Session,
63) -> Result<Plan> {
64 let query = match ast.statements.as_slice() {
65 [ast::Statement::Query(query)] => *query,
66 [] => return Err(Error::binder("no statement to bind")),
67 [_] => return Err(Error::not_implemented("a statement that is not a query")),
70 _ => return Err(Error::not_implemented("a script of more than one statement")),
71 };
72 let mut binder = Binder::with(catalog, parameters, session);
73 let (root, _) = binder.bind_query(ast, query)?;
74 let mut plan = binder.into_plan();
75 plan.set_root(root);
76 plan.validate()?;
77 Ok(plan)
78}
79
80pub fn bind_sql(query: &str, catalog: &Catalog) -> Result<Plan> {
86 bind_sql_with(query, catalog, &Session::new())
87}
88
89pub fn bind_sql_with(query: &str, catalog: &Catalog, session: &Session) -> Result<Plan> {
95 let ast = parse_ast_with_case(query, session.semantics().identifier_case())?;
96 bind_with(&ast, catalog, &Parameters::new(), session)
97}
98
99#[derive(Debug)]
101pub(crate) struct Aggregation {
102 pub(crate) index: u32,
104 pub(crate) groups: Vec<ExprRef>,
106 pub(crate) aggregates: Vec<ExprRef>,
108}
109
110#[derive(Debug)]
118pub(crate) struct WindowRun {
119 index: u32,
121 partition: Vec<ExprRef>,
123 order: Vec<SortKey>,
125 frame: WindowFrame,
127 calls: Vec<ExprRef>,
129}
130
131pub(crate) struct WindowCall<'a> {
136 pub(crate) name: &'a str,
138 pub(crate) args: &'a [ast::ExprRef],
140 pub(crate) distinct: bool,
142 pub(crate) filter: ast::ExprRef,
144 pub(crate) ignore_nulls: bool,
146 pub(crate) order: ast::Slice,
149 pub(crate) spec: ast::WindowRef,
151}
152
153struct WindowParts {
155 args: Vec<ExprRef>,
157 partition: Vec<ExprRef>,
159 order: Vec<SortKey>,
161 inner: Vec<SortKey>,
164 frame: WindowFrame,
166}
167
168#[derive(Debug)]
175struct Read {
176 fields: Vec<Field>,
178 rows: Stat<u64>,
180 distincts: Vec<(String, Stat<u64>)>,
182 zones: Option<Arc<dyn Zones>>,
184}
185
186impl Read {
187 fn uncounted(fields: Vec<Field>) -> Self {
189 Self { fields, rows: Stat::Unknown, distincts: Vec::new(), zones: None }
190 }
191}
192
193#[derive(Debug)]
195struct Materialized {
196 written: u32,
198 cte: u32,
200 name: String,
202 fields: Vec<Field>,
204}
205
206#[derive(Debug)]
207pub(crate) struct PendingSubquery {
208 pub(crate) node: NodeRef,
209 pub(crate) kind: JoinKind,
210 pub(crate) conditions: Vec<ExprRef>,
211 pub(crate) dependent: bool,
212 pub(crate) reads: Vec<ColumnBinding>,
218 pub(crate) index: u32,
224 pub(crate) inside_aggregate: bool,
231}
232
233#[derive(Debug, Clone, Copy, PartialEq, Eq)]
235enum Side {
236 Left,
237 Right,
238}
239
240#[derive(Debug)]
242pub(crate) struct Binder<'a> {
243 catalog: &'a Catalog,
244 pub(crate) parameters: &'a Parameters,
246 pub(crate) session: &'a Session,
248 pub(crate) semantics: Semantics,
250 plan: Plan,
251 next_index: u32,
252 pub(crate) current_span: Span,
254 pub(crate) aggregation: Option<Aggregation>,
256 want_ascending: bool,
259 pub(crate) upsert: bool,
262 pub(crate) insert_defaults: Option<Vec<(LogicalType, Option<String>)>>,
265 pub(crate) default_as_null: bool,
268 pub(crate) in_aggregate: bool,
270 pub(crate) in_filter: bool,
272 pub(crate) windows: Vec<WindowRun>,
274 pub(crate) sequences: Vec<QualifiedName>,
276 pub(crate) in_window: bool,
278 pub(crate) scalar_subqueries: Vec<PendingSubquery>,
280 pub(crate) joined_above: Vec<u32>,
286 pub(crate) outer_scopes: Vec<Scope>,
287 pub(crate) lateral_scopes: Vec<usize>,
294 pub(crate) correlations: Vec<Vec<ColumnBinding>>,
295 pub(crate) lambda_frames: Vec<crate::lambda::Frame>,
297 pub(crate) clause: &'static str,
299 pub(crate) outlined: bool,
308 expanding: Vec<String>,
310 materialized: Vec<Materialized>,
316 next_cte: u32,
318 started: Option<i64>,
320}
321
322impl<'a> Binder<'a> {
323 pub(crate) fn with(
324 catalog: &'a Catalog,
325 parameters: &'a Parameters,
326 session: &'a Session,
327 ) -> Self {
328 Self {
329 catalog,
330 parameters,
331 session,
332 semantics: session.semantics(),
333 plan: Plan::new(),
334 next_index: 0,
335 current_span: Span::new(0, 0),
336 aggregation: None,
337 want_ascending: false,
338 upsert: false,
339 insert_defaults: None,
340 default_as_null: false,
341 in_aggregate: false,
342 in_filter: false,
343 windows: Vec::new(),
344 sequences: Vec::new(),
345 in_window: false,
346 scalar_subqueries: Vec::new(),
347 joined_above: Vec::new(),
348 outer_scopes: Vec::new(),
349 lateral_scopes: Vec::new(),
350 correlations: Vec::new(),
351 lambda_frames: Vec::new(),
352 clause: "SELECT clause",
353 outlined: false,
354 expanding: Vec::new(),
355 materialized: Vec::new(),
356 next_cte: 0,
357 started: None,
358 }
359 }
360
361 pub(crate) fn catalog(&self) -> &Catalog {
362 self.catalog
363 }
364
365 pub(crate) fn instant(&mut self) -> i64 {
372 *self.started.get_or_insert_with(crate::context::micros_now)
373 }
374
375 pub(crate) fn plan(&self) -> &Plan {
376 &self.plan
377 }
378
379 pub(crate) fn plan_mut(&mut self) -> &mut Plan {
380 &mut self.plan
381 }
382
383 pub(crate) fn add_expr(&mut self, expr: Expr, ty: LogicalType) -> ExprRef {
384 self.plan.add_expr_at(expr, ty, self.current_span)
385 }
386
387 pub(crate) fn add_constant(&mut self, value: Value) -> ExprRef {
388 let ty = value.logical_type();
389 let reference = self.plan.add_value(value);
390 self.plan.add_expr_at(Expr::Constant(reference), ty, self.current_span)
391 }
392
393 pub(crate) fn add_node(&mut self, node: Node) -> NodeRef {
394 self.plan.add_node_at(node, self.current_span)
395 }
396
397 pub(crate) fn into_plan(self) -> Plan {
398 self.plan
399 }
400
401 pub(crate) fn fresh_index(&mut self) -> u32 {
403 let index = self.next_index;
404 self.next_index += 1;
405 index
406 }
407
408 fn column(&mut self, index: u32, position: usize, ty: LogicalType) -> ExprRef {
410 let binding = ColumnBinding::new(index, position as u32);
411 self.plan.add_expr(Expr::Column(binding), ty)
412 }
413
414 fn attach_scalar_subqueries(&mut self, mut input: NodeRef) -> NodeRef {
416 let subqueries = std::mem::take(&mut self.scalar_subqueries);
417 for pending in subqueries {
418 input = self.attach_subquery(input, pending);
419 }
420 input
421 }
422
423 fn attach_subquery(&mut self, input: NodeRef, pending: PendingSubquery) -> NodeRef {
430 let PendingSubquery {
431 node: mut right,
432 kind,
433 conditions,
434 dependent,
435 reads: _,
436 index: _,
437 inside_aggregate: _,
438 } = pending;
439 if kind == JoinKind::Single && !self.semantics.scalar_subquery_error_on_multiple_rows() {
440 right = self.add_node(Node::Limit {
441 input: right,
442 count: Bound::Rows(1),
443 offset: Bound::Rows(0),
444 });
445 }
446 let conditions = self.plan.add_expr_list(&conditions);
447 if dependent {
448 self.add_node(Node::DependentJoin { left: input, right, kind, conditions })
449 } else {
450 self.add_node(Node::Join {
451 left: input,
452 right,
453 kind,
454 conditions,
455 build: BuildSide::default(),
456 })
457 }
458 }
459
460 pub(crate) fn bind_query(
463 &mut self,
464 ast: &Ast,
465 query: ast::QueryRef,
466 ) -> Result<(NodeRef, Scope)> {
467 let span = ast.query_span(query);
468 let outer = std::mem::replace(&mut self.current_span, span);
469 let result =
470 self.bind_query_inner(ast, query).map_err(|error| error.with_fallback_span(span));
471 self.current_span = outer;
472 result
473 }
474
475 fn bind_query_inner(&mut self, ast: &Ast, query: ast::QueryRef) -> Result<(NodeRef, Scope)> {
476 let written = ast.query(query);
477 if written.ctes.is_empty() {
478 return self.bind_body(ast, &written);
479 }
480 let depth = self.materialized.len();
484 let result = self.bind_materialized(ast, &written);
485 self.materialized.truncate(depth);
486 result
487 }
488
489 fn bind_materialized(&mut self, ast: &Ast, written: &ast::Query) -> Result<(NodeRef, Scope)> {
495 let depth = self.materialized.len();
496 let held = ast.cte_list(written.ctes).to_vec();
497 let mut definitions = Vec::with_capacity(held.len());
498 for &index in &held {
499 definitions.push(self.bind_definition(ast, index)?);
500 }
501 let (mut node, scope) = self.bind_body(ast, written)?;
502 for (at, definition) in definitions.into_iter().enumerate().rev() {
503 let entry = &self.materialized[depth + at];
504 let cte = entry.cte;
505 let name = entry.name.clone();
506 let fields = entry.fields.clone();
507 let name = self.plan.intern(&name);
508 let columns = self.plan.add_fields(&fields);
509 node =
510 self.add_node(Node::MaterializedCte { definition, body: node, name, cte, columns });
511 }
512 Ok((node, scope))
513 }
514
515 fn bind_definition(&mut self, ast: &Ast, index: u32) -> Result<NodeRef> {
525 let held = ast.cte(index);
526 let name = ast.string(held.name).to_string();
527 let (node, mut scope) = self.bind_query(ast, held.query)?;
528 if !held.columns.is_empty() {
529 let names: Vec<&str> = ast.name(held.columns).collect();
530 scope.rename_prefix(&names);
531 }
532 let table = self.fresh_index();
533 let mut exprs = Vec::with_capacity(scope.len());
534 let mut names = Vec::with_capacity(scope.len());
535 for column in &scope.columns {
536 exprs.push(self.plan.add_expr(Expr::Column(column.binding), column.ty.clone()));
537 names.push(self.plan.intern(&column.name));
538 }
539 let exprs = self.plan.add_expr_list(&exprs);
540 let names = self.plan.add_name_list(&names);
541 let node = self.add_node(Node::Project { input: node, index: table, exprs, names });
542 let cte = self.next_cte;
543 self.next_cte += 1;
544 self.materialized.push(Materialized { written: index, cte, name, fields: scope.fields() });
545 Ok(node)
546 }
547
548 fn bind_body(&mut self, ast: &Ast, written: &ast::Query) -> Result<(NodeRef, Scope)> {
549 match written.body {
550 ast::QueryBody::Select(select) => self.bind_select(ast, select, written),
551 ast::QueryBody::SetOp { op, quantifier, by_name, left, right } => {
552 let operator = Operator { op, quantifier, by_name };
553 self.bind_set_op(ast, written, operator, left, right)
554 }
555 ast::QueryBody::Values(rows) => self.bind_values(ast, written, rows),
556 ast::QueryBody::Describe(inner) => self.bind_describe(ast, written, inner),
557 ast::QueryBody::Show { name, relation } => self.bind_show(ast, written, name, relation),
558 }
559 }
560
561 fn bind_show(
563 &mut self,
564 ast: &Ast,
565 query: &ast::Query,
566 name: ast::Slice,
567 relation: ast::QueryRef,
568 ) -> Result<(NodeRef, Scope)> {
569 let text = ast.name_text(name);
570 let parts: Vec<&str> = ast.name(name).collect();
571 let table_exists = self.catalog.resolve(&parts).is_ok();
572 let as_table = match self.semantics.show_behavior() {
573 ShowBehavior::Auto => table_exists,
574 ShowBehavior::Setting => false,
575 ShowBehavior::Table => true,
576 };
577 if as_table {
578 return self.bind_describe(ast, query, relation);
579 }
580 let shown = match self.session.iter().find(|(name, _)| name.eq_ignore_ascii_case(&text)) {
585 Some((_, value)) => value.to_string(),
586 None => match self.beyond(&text)? {
587 Some(Value::Varchar(declared)) => declared,
588 Some(other) => other.to_string(),
589 None => {
590 return Err(Error::catalog(format!(
591 "Setting with name \"{text}\" does not exist"
592 )));
593 }
594 },
595 };
596 let field = Field::new(text, LogicalType::Varchar);
597 let expr = self.plan.add_constant(Value::Varchar(shown));
598 let row = self.plan.add_expr_list(&[expr]);
599 let rows = self.plan.add_rows(&[row]);
600 let columns = self.plan.add_fields(std::slice::from_ref(&field));
601 let index = self.fresh_index();
602 let node = self.add_node(Node::Values { index, columns, rows });
603 let mut scope = Scope::empty();
604 scope.push(Visible {
605 table: String::new(),
606 name: field.name,
607 binding: ColumnBinding::new(index, 0),
608 ty: LogicalType::Varchar,
609 not_null: false,
610 key: None,
611 default: None,
612 qualified: false,
613 also: None,
614 });
615 Ok((node, scope))
616 }
617
618 fn bind_describe(
633 &mut self,
634 ast: &Ast,
635 query: &ast::Query,
636 inner: ast::QueryRef,
637 ) -> Result<(NodeRef, Scope)> {
638 let (_, described) = self.bind_query(ast, inner)?;
639 let fields: Vec<Field> = ["column_name", "column_type", "null", "key", "default", "extra"]
640 .iter()
641 .map(|name| Field::new(*name, LogicalType::Varchar))
642 .collect();
643 let mut slices = Vec::with_capacity(described.columns.len());
644 for column in described.columns.clone() {
645 let written = [
648 column.name.clone(),
649 column.ty.to_string(),
650 if column.not_null { "NO" } else { "YES" }.to_owned(),
651 ];
652 let mut items: Vec<ExprRef> = written
653 .into_iter()
654 .map(|text| self.plan.add_constant(Value::Varchar(text)))
655 .collect();
656 let mark = match column.key {
657 Some(mark) => self.plan.add_constant(Value::Varchar(mark.to_owned())),
658 None => {
659 let empty = self.plan.add_constant(Value::Null);
660 self.cast_to(empty, &LogicalType::Varchar)
661 }
662 };
663 items.push(mark);
664 if let Some(default) = &column.default {
665 items.push(self.plan.add_constant(Value::Varchar(default.clone())));
666 }
667 while items.len() < 6 {
668 let empty = self.plan.add_constant(Value::Null);
669 items.push(self.cast_to(empty, &LogicalType::Varchar));
670 }
671 slices.push(self.plan.add_expr_list(&items));
672 }
673 let rows = self.plan.add_rows(&slices);
674 let columns = self.plan.add_fields(&fields);
675 let index = self.fresh_index();
676 let mut node = self.add_node(Node::Values { index, columns, rows });
677 let mut scope = Scope::empty();
678 for (at, field) in fields.iter().enumerate() {
679 scope.push(Visible {
680 table: String::new(),
681 name: field.name.clone(),
682 binding: ColumnBinding::new(index, at as u32),
683 ty: field.ty.clone(),
684 not_null: false,
685 key: None,
686 default: None,
687 qualified: false,
688 also: None,
689 });
690 }
691 let keys = self.sort_keys(ast, query, &scope, &[])?;
692 if !keys.is_empty() {
693 let keys = self.plan.add_sort_keys(&keys);
694 node = self.add_node(Node::Sort { input: node, keys });
695 }
696 node = self.apply_limit(ast, query, node, &mut scope)?;
697 Ok((node, scope))
698 }
699
700 pub(crate) fn bind_default(&mut self, text: Option<&str>, ty: &LogicalType) -> Result<ExprRef> {
704 let Some(text) = text else {
705 let value = self.plan.add_value(Value::Null);
706 return Ok(self.plan.add_expr(Expr::Constant(value), ty.clone()));
707 };
708 let ast = rudb_parse::parse_ast(&format!("SELECT {text}"))?;
709 let found = match ast.statements.first() {
710 Some(&ast::Statement::Query(query)) => match ast.query(query).body {
711 ast::QueryBody::Select(select) => {
712 ast.target_list(ast.select(select).targets).first().map(|target| target.expr)
713 }
714 _ => None,
715 },
716 _ => None,
717 };
718 let Some(expr) = found else {
719 return Err(Error::internal(format!("a default that is not an expression: {text}")));
720 };
721 let expr = self.bind_expr(&ast, expr, &Scope::empty())?;
722 self.checked_cast_to(expr, ty, false)
723 }
724
725 fn passes_through(&self, expr: ExprRef, input: &Scope) -> bool {
731 let Expr::Column(binding) = *self.plan.expr(expr) else { return false };
732 input.columns.iter().any(|column| column.binding == binding && column.not_null)
733 }
734
735 fn key_through(&self, expr: ExprRef, input: &Scope) -> Option<&'static str> {
739 self.through(expr, input).and_then(|column| column.key)
740 }
741
742 fn through<'s>(&self, expr: ExprRef, input: &'s Scope) -> Option<&'s Visible> {
744 let Expr::Column(binding) = *self.plan.expr(expr) else { return None };
745 input.columns.iter().find(|column| column.binding == binding)
746 }
747
748 fn bind_values(
755 &mut self,
756 ast: &Ast,
757 query: &ast::Query,
758 rows: ast::Slice,
759 ) -> Result<(NodeRef, Scope)> {
760 let written = ast.rows(rows).to_vec();
761 let Some(first) = written.first() else {
762 return Err(Error::binder("VALUES needs at least one row"));
763 };
764 let width = first.len as usize;
765 for (at, row) in written.iter().enumerate() {
766 if row.len as usize != width {
767 return Err(Error::binder(format!(
768 "VALUES lists must all be the same length, expected {width} columns but row {} has {}",
769 at + 1,
770 row.len
771 )));
772 }
773 }
774 let empty = Scope::empty();
776 let defaults = self.insert_defaults.take();
777 let previous = std::mem::replace(&mut self.clause, "VALUES clause");
778 let mut bound: Vec<Vec<ExprRef>> = Vec::with_capacity(written.len());
779 for row in &written {
780 let mut items = Vec::with_capacity(width);
781 for (at, &expr) in ast.expr_list(*row).iter().enumerate() {
782 let column = defaults.as_ref().and_then(|defaults| defaults.get(at));
783 items.push(match (ast.expr(expr), column) {
784 (ast::Expr::Default, Some((ty, default))) => {
785 self.bind_default(default.as_deref(), ty)?
786 }
787 _ => self.bind_expr(ast, expr, &empty)?,
788 });
789 }
790 bound.push(items);
791 }
792 self.clause = previous;
793 let mut types = Vec::with_capacity(width);
794 for at in 0..width {
795 let mut ty = self.plan.expr_type(bound[0][at]).clone();
796 for row in &bound[1..] {
797 let other = self.plan.expr_type(row[at]).clone();
798 ty = ty.promote(&other).ok_or_else(|| {
799 Error::binder(format!(
800 "Cannot combine a value of type {ty} with a value of type {other} in column {} of a VALUES",
801 at + 1
802 ))
803 })?;
804 }
805 types.push(ty);
806 }
807 let mut slices = Vec::with_capacity(bound.len());
808 for row in &bound {
809 let items: Vec<ExprRef> = row
810 .iter()
811 .zip(&types)
812 .map(|(&expr, ty)| self.checked_cast_to(expr, ty, false))
813 .collect::<Result<_>>()?;
814 slices.push(self.plan.add_expr_list(&items));
815 }
816 let rows = self.plan.add_rows(&slices);
817 let fields: Vec<Field> = types
818 .iter()
819 .enumerate()
820 .map(|(at, ty)| Field::new(format!("col{at}"), ty.clone()))
821 .collect();
822 let columns = self.plan.add_fields(&fields);
823 let index = self.fresh_index();
824 let mut node = self.add_node(Node::Values { index, columns, rows });
825 let mut scope = Scope::empty();
826 for (at, field) in fields.iter().enumerate() {
827 scope.push(Visible {
828 table: String::new(),
829 name: field.name.clone(),
830 binding: ColumnBinding::new(index, at as u32),
831 ty: field.ty.clone(),
832 not_null: false,
833 key: None,
834 default: None,
835 qualified: false,
836 also: None,
837 });
838 }
839 let keys = self.sort_keys(ast, query, &scope, &[])?;
840 if !keys.is_empty() {
841 let keys = self.plan.add_sort_keys(&keys);
842 node = self.add_node(Node::Sort { input: node, keys });
843 }
844 node = self.apply_limit(ast, query, node, &mut scope)?;
845 Ok((node, scope))
846 }
847
848 fn bind_set_op(
849 &mut self,
850 ast: &Ast,
851 query: &ast::Query,
852 operator: Operator,
853 left: ast::QueryRef,
854 right: ast::QueryRef,
855 ) -> Result<(NodeRef, Scope)> {
856 let (left_node, left_scope) = self.bind_query(ast, left)?;
857 let (right_node, right_scope) = self.bind_query(ast, right)?;
858 let merged = if operator.by_name {
859 match_by_name(&left_scope, &right_scope)?
860 } else {
861 match_by_position(&left_scope, &right_scope)?
862 };
863 let left_node = self.conform(left_node, &left_scope, &merged, |column| column.left)?;
864 let right_node = self.conform(right_node, &right_scope, &merged, |column| column.right)?;
865 let index = self.fresh_index();
866 let kind = match operator.op {
867 SetOp::Union => SetOpKind::Union,
868 SetOp::Except => SetOpKind::Except,
869 SetOp::Intersect => SetOpKind::Intersect,
870 };
871 let all = operator.quantifier == Quantifier::All;
874 let mut node =
875 self.add_node(Node::SetOp { left: left_node, right: right_node, kind, all, index });
876 let mut scope = Scope::empty();
877 for (at, column) in merged.iter().enumerate() {
878 scope.push(Visible {
879 table: String::new(),
880 name: column.name.clone(),
881 binding: ColumnBinding::new(index, at as u32),
882 ty: column.ty.clone(),
883 not_null: false,
886 key: None,
887 default: None,
888 qualified: false,
889 also: None,
890 });
891 }
892 let keys = self.sort_keys(ast, query, &scope, &[])?;
896 if !keys.is_empty() {
897 let keys = self.plan.add_sort_keys(&keys);
898 node = self.add_node(Node::Sort { input: node, keys });
899 }
900 node = self.apply_limit(ast, query, node, &mut scope)?;
901 Ok((node, scope))
902 }
903
904 fn conform(
910 &mut self,
911 node: NodeRef,
912 scope: &Scope,
913 merged: &[Merged],
914 pick: impl Fn(&Merged) -> Option<usize>,
915 ) -> Result<NodeRef> {
916 let unchanged = merged.len() == scope.len()
917 && merged
918 .iter()
919 .enumerate()
920 .all(|(at, column)| pick(column) == Some(at) && column.ty == scope.columns[at].ty);
921 if unchanged {
922 return Ok(node);
923 }
924 let index = self.fresh_index();
925 let mut exprs = Vec::with_capacity(merged.len());
926 let mut names = Vec::with_capacity(merged.len());
927 for column in merged {
928 let expr = match pick(column) {
929 Some(at) => {
930 let held = &scope.columns[at];
931 self.plan.add_expr(Expr::Column(held.binding), held.ty.clone())
932 }
933 None => self.plan.add_constant(Value::Null),
934 };
935 exprs.push(self.checked_cast_to(expr, &column.ty, false)?);
936 names.push(self.plan.intern(&column.name));
937 }
938 let exprs = self.plan.add_expr_list(&exprs);
939 let names = self.plan.add_name_list(&names);
940 Ok(self.add_node(Node::Project { input: node, index, exprs, names }))
941 }
942
943 fn bind_select(
946 &mut self,
947 ast: &Ast,
948 select: ast::SelectRef,
949 query: &ast::Query,
950 ) -> Result<(NodeRef, Scope)> {
951 let written = ast.select(select);
952 self.want_ascending |= !written.group_by.is_empty() || written.group_by_all;
953 let outer_windows = std::mem::take(&mut self.windows);
957 let outer_joined_above = std::mem::take(&mut self.joined_above);
962 let (mut node, input) = self.bind_from(ast, written.from)?;
963 node = self.attach_scalar_subqueries(node);
964
965 if written.filter != NONE {
966 self.clause = "WHERE clause";
967 let predicate = self.bind_expr(ast, written.filter, &input)?;
968 let predicate = self.as_boolean(predicate, "WHERE")?;
969 node = self.attach_scalar_subqueries(node);
970 node = self.add_node(Node::Filter { input: node, predicate });
971 }
972
973 let targets = ast.target_list(written.targets).to_vec();
974 if targets.is_empty() {
975 return Err(Error::binder("a SELECT needs at least one expression to select"));
976 }
977
978 let group_items = self.group_items(ast, &written, &targets)?;
979 let aggregating = !group_items.is_empty()
980 || written.having != NONE
981 || targets.iter().any(|target| has_aggregate(ast, target.expr));
982 if aggregating {
983 self.clause = "GROUP BY clause";
984 let mut groups = Vec::with_capacity(group_items.len());
985 for item in &group_items {
986 groups.push(self.bind_expr(ast, *item, &input)?);
987 }
988 let index = self.fresh_index();
989 self.aggregation = Some(Aggregation { index, groups, aggregates: Vec::new() });
990 }
991
992 let mut above = Vec::new();
999
1000 self.clause = "SELECT clause";
1001 let (mut exprs, mut names) = self.bind_targets(ast, &targets, &input, &mut above)?;
1002 let visible = exprs.len();
1003
1004 let mut having = None;
1005 if written.having != NONE {
1006 self.clause = "HAVING clause";
1007 let before = self.scalar_subqueries.len();
1008 let predicate = self.bind_expr(ast, written.having, &input)?;
1009 self.lift_over_aggregate(before, &mut above, &input)?;
1010 let predicate = self.over_aggregate(predicate, &input)?;
1011 having = Some(self.as_boolean(predicate, "HAVING")?);
1012 }
1013
1014 let project = self.fresh_index();
1017 let mut output = Scope::empty();
1018 for (at, (expr, name)) in exprs.iter().zip(&names).enumerate() {
1019 output.push(Visible {
1020 table: String::new(),
1021 name: name.clone(),
1022 binding: ColumnBinding::new(project, at as u32),
1023 ty: self.plan.expr_type(*expr).clone(),
1024 not_null: self.passes_through(*expr, &input),
1025 key: self.key_through(*expr, &input),
1026 default: self.through(*expr, &input).and_then(|column| column.default.clone()),
1027 qualified: false,
1028 also: None,
1029 });
1030 }
1031
1032 self.clause = "ORDER BY clause";
1033 let mut extra = Vec::new();
1034 let keys = self.select_sort_keys(
1035 ast, query, &input, &output, project, &mut exprs, &mut names, &mut extra, &mut above,
1036 )?;
1037 self.joined_above = outer_joined_above;
1038 if !extra.is_empty() && written.distinct != Distinct::No {
1039 return Err(Error::binder(
1040 "For SELECT DISTINCT, ORDER BY expressions must appear in the select list",
1041 ));
1042 }
1043 let on = self.distinct_on(ast, written.distinct, &output)?;
1044
1045 node = self.attach_scalar_subqueries(node);
1046
1047 if let Some(aggregation) = self.aggregation.take() {
1048 let index = aggregation.index;
1049 let groups = self.plan.add_expr_list(&aggregation.groups);
1050 let aggregates = self.plan.add_expr_list(&aggregation.aggregates);
1051 node = self.add_node(Node::Aggregate { input: node, index, groups, aggregates });
1052 }
1053 if !above.is_empty() {
1054 debug_assert!(self.scalar_subqueries.is_empty(), "a query is waiting to be joined");
1055 self.scalar_subqueries = above;
1056 node = self.attach_scalar_subqueries(node);
1057 }
1058 if let Some(predicate) = having {
1059 node = self.add_node(Node::Filter { input: node, predicate });
1060 }
1061
1062 for run in std::mem::replace(&mut self.windows, outer_windows) {
1066 let partition = self.plan.add_expr_list(&run.partition);
1067 let order = self.plan.add_sort_keys(&run.order);
1068 let expressions = self.plan.add_expr_list(&run.calls);
1069 node = self.add_node(Node::Window {
1070 input: node,
1071 index: run.index,
1072 partition,
1073 order,
1074 frame: run.frame,
1075 expressions,
1076 });
1077 }
1078
1079 let interned: Vec<u32> = names.iter().map(|name| self.plan.intern(name)).collect();
1080 let exprs_slice = self.plan.add_expr_list(&exprs);
1081 let names_slice = self.plan.add_name_list(&interned);
1082 node = self.add_node(Node::Project {
1083 input: node,
1084 index: project,
1085 exprs: exprs_slice,
1086 names: names_slice,
1087 });
1088
1089 if written.distinct != Distinct::No {
1090 let on = self.plan.add_expr_list(&on);
1091 node = self.add_node(Node::Distinct { input: node, on });
1092 }
1093 if !keys.is_empty() {
1094 let keys = self.plan.add_sort_keys(&keys);
1095 node = self.add_node(Node::Sort { input: node, keys });
1096 }
1097 node = self.apply_limit(ast, query, node, &mut output)?;
1098
1099 if extra.is_empty() {
1100 output.columns.truncate(visible);
1101 return Ok((node, output));
1102 }
1103 let index = self.fresh_index();
1106 let mut kept = Vec::with_capacity(visible);
1107 let mut kept_names = Vec::with_capacity(visible);
1108 let mut scope = Scope::empty();
1109 for (at, name) in names.iter().enumerate().take(visible) {
1110 let ty = output.columns[at].ty.clone();
1111 let binding = output.columns[at].binding;
1115 kept.push(self.plan.add_expr(Expr::Column(binding), ty.clone()));
1116 kept_names.push(self.plan.intern(name));
1117 scope.push(Visible {
1118 table: String::new(),
1119 name: name.clone(),
1120 binding: ColumnBinding::new(index, at as u32),
1121 ty,
1122 not_null: output.columns[at].not_null,
1123 key: output.columns[at].key,
1124 default: output.columns[at].default.clone(),
1125 qualified: false,
1126 also: None,
1127 });
1128 }
1129 let exprs = self.plan.add_expr_list(&kept);
1130 let names = self.plan.add_name_list(&kept_names);
1131 node = self.add_node(Node::Project { input: node, index, exprs, names });
1132 Ok((node, scope))
1133 }
1134
1135 fn lift_over_aggregate(
1154 &mut self,
1155 before: usize,
1156 above: &mut Vec<PendingSubquery>,
1157 scope: &Scope,
1158 ) -> Result<()> {
1159 if self.aggregation.is_none() {
1160 return Ok(());
1161 }
1162 let mut lifted = Vec::new();
1163 for mut pending in self.scalar_subqueries.split_off(before) {
1164 let stays = pending.inside_aggregate
1165 || (pending.dependent && !self.lift_correlated(&mut pending));
1166 if stays {
1167 self.scalar_subqueries.push(pending);
1168 } else {
1169 self.joined_above.push(pending.index);
1170 lifted.push(pending);
1171 }
1172 }
1173 for pending in &mut lifted {
1178 let conditions = std::mem::take(&mut pending.conditions);
1179 let mut over = Vec::with_capacity(conditions.len());
1180 for condition in conditions {
1181 over.push(self.over_aggregate(condition, scope)?);
1182 }
1183 pending.conditions = over;
1184 }
1185 above.append(&mut lifted);
1186 Ok(())
1187 }
1188
1189 fn lift_correlated(&mut self, pending: &mut PendingSubquery) -> bool {
1206 let Some(index) = self.aggregation.as_ref().map(|aggregation| aggregation.index) else {
1207 return false;
1208 };
1209 let mut moved = Vec::with_capacity(pending.reads.len());
1210 for read in &pending.reads {
1211 let Some(at) = self.group_of(*read) else {
1212 return false;
1213 };
1214 moved.push((*read, ColumnBinding::new(index, at as u32)));
1215 }
1216 let mut rewrites = Vec::new();
1217 self.plan.subtree_columns(pending.node, &mut |reference, binding| {
1218 if let Some(&(_, to)) = moved.iter().find(|(from, _)| *from == binding) {
1219 rewrites.push((reference, to));
1220 }
1221 });
1222 for (reference, to) in rewrites {
1223 self.plan.rebind(reference, to);
1224 }
1225 pending.reads = moved.into_iter().map(|(_, to)| to).collect();
1226 true
1227 }
1228
1229 fn bind_targets(
1230 &mut self,
1231 ast: &Ast,
1232 targets: &[ast::Target],
1233 input: &Scope,
1234 above: &mut Vec<PendingSubquery>,
1235 ) -> Result<(Vec<ExprRef>, Vec<String>)> {
1236 let mut exprs = Vec::with_capacity(targets.len());
1237 let mut names = Vec::with_capacity(targets.len());
1238 for target in targets {
1239 if let ast::Expr::Star { qualifier, replacements } = ast.expr(target.expr) {
1240 let table = ast.name(qualifier).last().map(str::to_string);
1241 let expanded: Vec<Visible> =
1242 input.star(table.as_deref())?.into_iter().cloned().collect();
1243 let replacements = ast.target_list(replacements).to_vec();
1244 let mut used = vec![false; replacements.len()];
1245 for column in expanded {
1246 let found = replacements.iter().zip(&mut used).find(|(replacement, _)| {
1247 same_name(ast.string(replacement.alias), &column.name)
1248 });
1249 let before = self.scalar_subqueries.len();
1254 let (expr, name) = match found {
1255 Some((replacement, used)) => {
1256 *used = true;
1257 let expr = self.bind_expr(ast, replacement.expr, input)?;
1258 (expr, ast.string(replacement.alias).to_string())
1259 }
1260 None => (
1261 self.plan.add_expr(Expr::Column(column.binding), column.ty),
1262 column.name,
1263 ),
1264 };
1265 self.lift_over_aggregate(before, above, input)?;
1266 exprs.push(self.over_aggregate(expr, input)?);
1267 names.push(name);
1268 }
1269 if let Some((replacement, _)) =
1273 replacements.iter().zip(&used).find(|(_, used)| !**used)
1274 {
1275 return Err(missing_replacement(ast.string(replacement.alias), input));
1276 }
1277 continue;
1278 }
1279 let before = self.scalar_subqueries.len();
1280 let expr = self.bind_expr(ast, target.expr, input)?;
1281 self.lift_over_aggregate(before, above, input)?;
1282 exprs.push(self.over_aggregate(expr, input)?);
1283 names.push(if target.alias == NONE {
1284 self.output_name(ast, target.expr, input)
1285 } else {
1286 ast.string(target.alias).to_string()
1287 });
1288 }
1289 Ok((exprs, names))
1290 }
1291
1292 fn output_name(&self, ast: &Ast, target: ast::ExprRef, input: &Scope) -> String {
1298 if let ast::Expr::Column { name } = ast.expr(target) {
1299 let parts: Vec<&str> = ast.name(name).collect();
1300 if let Ok(found) = input.resolve(&parts) {
1301 let written = parts.last().copied().unwrap_or_default();
1304 if let Some(also) = &found.also {
1305 if !same_name(&found.name, written) && same_name(also, written) {
1306 return also.clone();
1307 }
1308 }
1309 return found.name.clone();
1310 }
1311 }
1312 describe(ast, target, self.semantics)
1313 }
1314
1315 fn group_items(
1317 &self,
1318 ast: &Ast,
1319 select: &ast::Select,
1320 targets: &[ast::Target],
1321 ) -> Result<Vec<ast::ExprRef>> {
1322 if select.group_by_all {
1323 return Ok(targets
1326 .iter()
1327 .filter(|target| !has_aggregate(ast, target.expr))
1328 .map(|target| target.expr)
1329 .collect());
1330 }
1331 let mut items = Vec::new();
1332 for &item in ast.expr_list(select.group_by) {
1333 items.push(self.output_reference(ast, item, targets, "GROUP BY")?.unwrap_or(item));
1334 }
1335 Ok(items)
1336 }
1337
1338 fn output_reference(
1340 &self,
1341 ast: &Ast,
1342 item: ast::ExprRef,
1343 targets: &[ast::Target],
1344 clause: &str,
1345 ) -> Result<Option<ast::ExprRef>> {
1346 match ast.expr(item) {
1347 ast::Expr::Literal { kind: LiteralKind::Number, text } => {
1348 let written = ast.string(text);
1349 let position: usize = written.parse().map_err(|_| {
1350 Error::binder(format!("{clause} term {written} is not a column"))
1351 })?;
1352 if position == 0 || position > targets.len() {
1353 return Err(Error::binder(format!(
1354 "{clause} term out of range - should be between 1 and {}",
1355 targets.len()
1356 )));
1357 }
1358 Ok(Some(targets[position - 1].expr))
1359 }
1360 ast::Expr::Column { name } => {
1361 let parts: Vec<&str> = ast.name(name).collect();
1362 let [written] = parts.as_slice() else { return Ok(None) };
1363 let mut found = None;
1364 for target in targets {
1365 if target.alias != NONE && same_name(ast.string(target.alias), written) {
1366 if found.is_some() {
1367 return Ok(None);
1368 }
1369 found = Some(target.expr);
1370 }
1371 }
1372 Ok(found)
1373 }
1374 _ => Ok(None),
1375 }
1376 }
1377
1378 #[allow(clippy::too_many_arguments)]
1382 fn select_sort_keys(
1383 &mut self,
1384 ast: &Ast,
1385 query: &ast::Query,
1386 input: &Scope,
1387 output: &Scope,
1388 project: u32,
1389 exprs: &mut Vec<ExprRef>,
1390 names: &mut Vec<String>,
1391 extra: &mut Vec<usize>,
1392 above: &mut Vec<PendingSubquery>,
1393 ) -> Result<Vec<SortKey>> {
1394 if query.order_by_all {
1395 return Ok(self.every_column(output));
1396 }
1397 let items = ast.order_list(query.order_by).to_vec();
1398 let mut keys = Vec::with_capacity(items.len());
1399 for item in items {
1400 self.check_order_literal(ast, item.expr)?;
1401 let position = match self.output_position(ast, item.expr, output)? {
1402 Some(position) => position,
1403 None => {
1404 let before = self.scalar_subqueries.len();
1405 let bound = self.bind_expr(ast, item.expr, input)?;
1406 self.lift_over_aggregate(before, above, input)?;
1407 let bound = self.over_aggregate(bound, input)?;
1408 match exprs.iter().position(|&held| self.same_expr(held, bound)) {
1409 Some(position) => position,
1410 None => {
1411 exprs.push(bound);
1412 names.push(describe(ast, item.expr, self.semantics));
1413 extra.push(exprs.len() - 1);
1414 exprs.len() - 1
1415 }
1416 }
1417 }
1418 };
1419 let ty = self.plan.expr_type(exprs[position]).clone();
1420 let expr = self.column(project, position, ty);
1421 keys.push(self.sort_key(expr, item));
1422 }
1423 Ok(keys)
1424 }
1425
1426 fn sort_keys(
1428 &mut self,
1429 ast: &Ast,
1430 query: &ast::Query,
1431 output: &Scope,
1432 targets: &[ast::Target],
1433 ) -> Result<Vec<SortKey>> {
1434 if query.order_by_all {
1435 return Ok(self.every_column(output));
1436 }
1437 let items = ast.order_list(query.order_by).to_vec();
1438 let mut keys = Vec::with_capacity(items.len());
1439 for item in items {
1440 self.check_order_literal(ast, item.expr)?;
1441 let expr = match self.output_position(ast, item.expr, output)? {
1442 Some(position) => {
1443 let column = &output.columns[position];
1444 let (binding, ty) = (column.binding, column.ty.clone());
1445 self.plan.add_expr(Expr::Column(binding), ty)
1446 }
1447 None => {
1448 let _ = targets;
1449 self.bind_expr(ast, item.expr, output)?
1450 }
1451 };
1452 keys.push(self.sort_key(expr, item));
1453 }
1454 Ok(keys)
1455 }
1456
1457 fn every_column(&mut self, output: &Scope) -> Vec<SortKey> {
1458 let columns: Vec<(ColumnBinding, LogicalType)> =
1459 output.columns.iter().map(|column| (column.binding, column.ty.clone())).collect();
1460 columns
1461 .into_iter()
1462 .map(|(binding, ty)| {
1463 let expr = self.plan.add_expr(Expr::Column(binding), ty);
1464 let expr = self.by_position(expr);
1465 let descending = self.semantics.default_descending();
1466 SortKey { expr, descending, nulls_first: self.semantics.nulls_first(descending) }
1467 })
1468 .collect()
1469 }
1470
1471 fn sort_key(&mut self, expr: ExprRef, item: ast::OrderItem) -> SortKey {
1473 let expr = self.by_position(expr);
1474 let descending = match item.order {
1475 Order::Unstated => self.semantics.default_descending(),
1476 Order::Ascending => false,
1477 Order::Descending => true,
1478 };
1479 let nulls_first = match item.nulls {
1480 Nulls::First => true,
1481 Nulls::Last => false,
1482 Nulls::Unstated => self.semantics.nulls_first(descending),
1483 };
1484 SortKey { expr, descending, nulls_first }
1485 }
1486
1487 fn output_position(
1489 &self,
1490 ast: &Ast,
1491 item: ast::ExprRef,
1492 output: &Scope,
1493 ) -> Result<Option<usize>> {
1494 match ast.expr(item) {
1495 ast::Expr::Literal { kind: LiteralKind::Number, text } => {
1496 let written = ast.string(text);
1497 if written.contains(['.', 'e', 'E']) {
1498 return Ok(None);
1499 }
1500 let position: usize = written.parse().map_err(|_| {
1501 Error::binder(format!("ORDER BY term {written} is not a column"))
1502 })?;
1503 if position == 0 || position > output.len() {
1504 return Err(Error::binder(format!(
1505 "ORDER BY term out of range - should be between 1 and {}",
1506 output.len()
1507 )));
1508 }
1509 Ok(Some(position - 1))
1510 }
1511 ast::Expr::Column { name } => {
1512 let parts: Vec<&str> = ast.name(name).collect();
1513 let [written] = parts.as_slice() else { return Ok(None) };
1514 Ok(output.position_of(None, written))
1515 }
1516 _ => Ok(None),
1517 }
1518 }
1519
1520 fn check_order_literal(&self, ast: &Ast, item: ast::ExprRef) -> Result<()> {
1522 if !self.semantics.order_by_non_integer_literal()
1523 && matches!(
1524 ast.expr(item),
1525 ast::Expr::Literal { kind, text }
1526 if kind != LiteralKind::Number
1527 || ast.string(text).contains(['.', 'e', 'E'])
1528 )
1529 {
1530 return Err(Error::binder(
1531 "ORDER BY non-integer literal has no effect.\n* SET order_by_non_integer_literal=true to allow this behavior.",
1532 ));
1533 }
1534 Ok(())
1535 }
1536
1537 fn distinct_on(
1539 &mut self,
1540 ast: &Ast,
1541 distinct: Distinct,
1542 output: &Scope,
1543 ) -> Result<Vec<ExprRef>> {
1544 let Distinct::On(items) = distinct else {
1545 return Ok(Vec::new());
1546 };
1547 let items = ast.expr_list(items).to_vec();
1548 let mut on = Vec::with_capacity(items.len());
1549 for item in items {
1550 let Some(position) = self.output_position(ast, item, output)? else {
1551 return Err(Error::not_implemented(
1552 "DISTINCT ON an expression that is not in the select list",
1553 ));
1554 };
1555 let column = &output.columns[position];
1556 let (binding, ty) = (column.binding, column.ty.clone());
1557 on.push(self.plan.add_expr(Expr::Column(binding), ty));
1558 }
1559 Ok(on)
1560 }
1561
1562 fn apply_limit(
1569 &mut self,
1570 ast: &Ast,
1571 query: &ast::Query,
1572 input: NodeRef,
1573 scope: &mut Scope,
1574 ) -> Result<NodeRef> {
1575 let waiting = self.scalar_subqueries.len();
1576 if query.limit_percent {
1577 let percent = self.share(ast, query.limit)?;
1578 let offset = self.skipped(ast, query.offset)?;
1579 let node = |binder: &mut Self, input| match percent {
1580 Some(percent) => binder.add_node(Node::LimitPercent { input, percent, offset }),
1581 None => binder.limited(input, Bound::All, offset),
1584 };
1585 return self.over_subqueries(waiting, input, scope, node);
1586 }
1587 let count = self.count_bound(ast, query.limit, "LIMIT")?;
1588 let offset = self.skipped(ast, query.offset)?;
1589 let node = |binder: &mut Self, input| binder.limited(input, count, offset);
1590 self.over_subqueries(waiting, input, scope, node)
1591 }
1592
1593 fn skipped(&mut self, ast: &Ast, written: ast::ExprRef) -> Result<Bound> {
1598 Ok(match self.count_bound(ast, written, "OFFSET")? {
1599 Bound::All => Bound::Rows(0),
1600 named => named,
1601 })
1602 }
1603
1604 fn over_subqueries(
1611 &mut self,
1612 waiting: usize,
1613 input: NodeRef,
1614 scope: &mut Scope,
1615 node: impl FnOnce(&mut Self, NodeRef) -> NodeRef,
1616 ) -> Result<NodeRef> {
1617 let joined = self.scalar_subqueries.split_off(waiting);
1618 if joined.is_empty() {
1619 return Ok(node(self, input));
1620 }
1621 let mut input = input;
1622 for pending in joined {
1623 input = self.attach_subquery(input, pending);
1624 }
1625 let limit = node(self, input);
1626 Ok(self.reproject(limit, scope))
1627 }
1628
1629 fn limited(&mut self, input: NodeRef, count: Bound, offset: Bound) -> NodeRef {
1632 if count == Bound::All && offset == Bound::Rows(0) {
1633 return input;
1634 }
1635 self.add_node(Node::Limit { input, count, offset })
1636 }
1637
1638 fn reproject(&mut self, node: NodeRef, scope: &mut Scope) -> NodeRef {
1644 let index = self.fresh_index();
1645 let mut exprs = Vec::with_capacity(scope.columns.len());
1646 let mut names = Vec::with_capacity(scope.columns.len());
1647 for column in &scope.columns {
1648 exprs.push(self.plan.add_expr(Expr::Column(column.binding), column.ty.clone()));
1649 names.push(self.plan.intern(&column.name));
1650 }
1651 for (at, column) in scope.columns.iter_mut().enumerate() {
1652 column.binding = ColumnBinding::new(index, at as u32);
1653 }
1654 let exprs = self.plan.add_expr_list(&exprs);
1655 let names = self.plan.add_name_list(&names);
1656 self.add_node(Node::Project { input: node, index, exprs, names })
1657 }
1658
1659 fn share(&mut self, ast: &Ast, written: ast::ExprRef) -> Result<Option<Share>> {
1676 if written == NONE {
1677 return Ok(None);
1678 }
1679 self.clause = "LIMIT clause";
1680 let scope = Scope::empty();
1681 let bound = self.bind_expr(ast, written, &scope)?;
1682 let Some(value) = fold::value_of(&self.plan, bound)? else {
1683 return Ok(Some(Share::Read(bound)));
1684 };
1685 if value.is_null() {
1686 return Ok(None);
1687 }
1688 let percent = percentage(&value)?;
1689 if !(0.0..=100.0).contains(&percent) {
1690 return Err(Error::out_of_range(
1691 "Limit percent out of range, should be between 0% and 100%",
1692 ));
1693 }
1694 Ok(Some(Share::Percent(percent)))
1695 }
1696
1697 fn count_bound(&mut self, ast: &Ast, written: ast::ExprRef, clause: &str) -> Result<Bound> {
1716 if written == NONE {
1717 return Ok(Bound::All);
1718 }
1719 self.clause = "LIMIT clause";
1720 let scope = Scope::empty();
1721 let bound = self.bind_expr(ast, written, &scope)?;
1722 let Some(value) = fold::value_of(&self.plan, bound)? else {
1723 return Ok(Bound::Read(bound));
1724 };
1725 if value.is_null() {
1728 return Ok(Bound::All);
1729 }
1730 row_count(&value, clause).map(Bound::Rows)
1731 }
1732
1733 fn bind_from(&mut self, ast: &Ast, from: ast::Slice) -> Result<(NodeRef, Scope)> {
1736 let sources = ast.source_list(from).to_vec();
1737 let Some((first, rest)) = sources.split_first() else {
1738 return Ok((self.add_node(Node::Dummy), Scope::empty()));
1741 };
1742 let (mut node, mut scope) = self.bind_source(ast, *first)?;
1743 for source in rest {
1744 let (right, right_scope, correlations) = self.bind_lateral(ast, *source, &scope)?;
1745 node = if correlations.is_empty() {
1746 self.add_node(Node::CrossProduct { left: node, right })
1747 } else {
1748 let conditions = self.plan.add_expr_list(&[]);
1749 self.add_node(Node::DependentJoin {
1750 left: node,
1751 right,
1752 kind: JoinKind::Inner,
1753 conditions,
1754 })
1755 };
1756 scope = scope.concat(right_scope);
1757 }
1758 Ok((node, scope))
1759 }
1760
1761 fn bind_lateral(
1773 &mut self,
1774 ast: &Ast,
1775 source: ast::SourceRef,
1776 left: &Scope,
1777 ) -> Result<(NodeRef, Scope, Vec<ColumnBinding>)> {
1778 self.lateral_scopes.push(self.outer_scopes.len());
1779 self.outer_scopes.push(left.clone());
1780 self.correlations.push(Vec::new());
1781 let bound = self.bind_source(ast, source);
1782 let read = self.correlations.pop().expect("correlation frame");
1783 self.outer_scopes.pop();
1784 self.lateral_scopes.pop();
1785 let (node, scope) = bound?;
1786
1787 let mut here = Vec::new();
1788 for binding in read {
1789 if left.columns.iter().any(|column| column.binding == binding) {
1790 here.push(binding);
1791 } else if let Some(enclosing) = self.correlations.last_mut() {
1792 if !enclosing.contains(&binding) {
1793 enclosing.push(binding);
1794 }
1795 }
1796 }
1797 Ok((node, scope, here))
1807 }
1808
1809 fn bind_source(&mut self, ast: &Ast, source: ast::SourceRef) -> Result<(NodeRef, Scope)> {
1810 match ast.source(source) {
1811 ast::Source::Table { name, alias, columns } => {
1812 self.bind_table(ast, name, alias, columns)
1813 }
1814 ast::Source::Function { name, args, alias, columns, pragma } => {
1815 self.bind_table_function(ast, name, args, alias, columns, pragma)
1816 }
1817 ast::Source::Subquery { query, alias, columns } => {
1818 let (node, mut scope) = self.bind_query(ast, query)?;
1819 let label = if alias == NONE {
1820 "unnamed_subquery".to_string()
1821 } else {
1822 ast.string(alias).to_string()
1823 };
1824 scope.relabel(&label);
1825 if !columns.is_empty() {
1826 let names: Vec<&str> = ast.name(columns).collect();
1827 scope.rename(&names, &label)?;
1828 }
1829 Ok((node, scope))
1830 }
1831 ast::Source::Values { rows, alias, columns } => {
1832 let bare = ast::Query::bare(ast::QueryBody::Values(rows));
1833 let (node, mut scope) = self.bind_values(ast, &bare, rows)?;
1834 let label =
1835 if alias == NONE { String::new() } else { ast.string(alias).to_string() };
1836 scope.relabel(&label);
1837 if !columns.is_empty() {
1838 let names: Vec<&str> = ast.name(columns).collect();
1839 scope.rename(&names, &label)?;
1840 }
1841 Ok((node, scope))
1842 }
1843 ast::Source::Cte { cte, alias, columns } => {
1844 self.bind_cte_scan(ast, cte, alias, columns)
1845 }
1846 ast::Source::Join { left, right, kind, natural, on, using } => {
1847 self.bind_join(ast, left, right, kind, natural, on, using)
1848 }
1849 }
1850 }
1851
1852 fn bind_cte_scan(
1859 &mut self,
1860 ast: &Ast,
1861 written: u32,
1862 alias: ast::StrRef,
1863 columns: ast::Slice,
1864 ) -> Result<(NodeRef, Scope)> {
1865 let Some(held) = self.materialized.iter().rev().find(|held| held.written == written) else {
1866 let name = ast.string(ast.cte(written).name);
1867 return Err(Error::binder(format!("Table with name {name} does not exist!")));
1868 };
1869 let cte = held.cte;
1870 let fields = held.fields.clone();
1871 let text = held.name.clone();
1872 let label = if alias == NONE { text.clone() } else { ast.string(alias).to_string() };
1873 let name = self.plan.intern(&text);
1874 let index = self.fresh_index();
1875 let mut scope = Scope::empty();
1876 for (at, field) in fields.iter().enumerate() {
1877 scope.push(Visible {
1878 table: label.clone(),
1879 name: field.name.clone(),
1880 binding: ColumnBinding::new(index, at as u32),
1881 ty: field.ty.clone(),
1882 not_null: field.not_null,
1883 key: None,
1884 default: None,
1885 qualified: false,
1886 also: None,
1887 });
1888 }
1889 if !columns.is_empty() {
1890 let names: Vec<&str> = ast.name(columns).collect();
1891 scope.rename(&names, &label)?;
1892 }
1893 let columns = self.plan.add_fields(&fields);
1894 let node = self.add_node(Node::CteScan { index, cte, name, columns });
1895 Ok((node, scope))
1896 }
1897
1898 fn bind_table(
1899 &mut self,
1900 ast: &Ast,
1901 name: ast::Slice,
1902 alias: ast::StrRef,
1903 columns: ast::Slice,
1904 ) -> Result<(NodeRef, Scope)> {
1905 let parts: Vec<&str> = ast.name(name).collect();
1906 let catalog = self.catalog;
1907 let resolved = match catalog.resolve(&parts) {
1910 Ok(resolved) => resolved,
1911 Err(missing) => {
1912 return self.bind_replacement_scan(ast, &parts, alias, columns, missing);
1913 }
1914 };
1915 if catalog.entry(&resolved)? == Entry::View {
1916 return self.bind_view(ast, &resolved, alias, columns);
1917 }
1918 let label =
1919 if alias == NONE { resolved.table.clone() } else { ast.string(alias).to_string() };
1920 self.bind_catalog_table(ast, &resolved, label, columns)
1921 }
1922
1923 pub(crate) fn bind_catalog_table(
1926 &mut self,
1927 ast: &Ast,
1928 resolved: &QualifiedName,
1929 label: String,
1930 columns: ast::Slice,
1931 ) -> Result<(NodeRef, Scope)> {
1932 let table = self.catalog.table(resolved)?;
1933 let fields: Vec<Field> = table.columns().to_vec();
1934 let excluded = self.upsert && same_name(&label, "excluded");
1937 let mut marks = vec![None; fields.len()];
1940 for key in table.keys() {
1941 for &column in &key.columns {
1942 if key.primary || marks[column].is_none() {
1943 marks[column] = Some(if key.primary { "PRI" } else { "UNI" });
1944 }
1945 }
1946 }
1947 let index = self.fresh_index();
1948 let mut scope = Scope::empty();
1949 for (at, field) in fields.iter().enumerate() {
1950 scope.push(Visible {
1951 table: label.clone(),
1952 name: field.name.clone(),
1953 binding: ColumnBinding::new(index, at as u32),
1954 ty: field.ty.clone(),
1955 not_null: field.not_null,
1956 key: marks[at],
1957 default: table.default(at).map(str::to_owned),
1958 qualified: excluded,
1959 also: None,
1960 });
1961 }
1962 if !columns.is_empty() {
1963 let names: Vec<&str> = ast.name(columns).collect();
1964 scope.rename(&names, &label)?;
1965 }
1966 let resolved = if excluded { &QualifiedName::excluded() } else { resolved };
1967 let catalog_name = self.plan.intern(&resolved.catalog);
1968 let schema = self.plan.intern(&resolved.schema);
1969 let table_name = self.plan.intern(&resolved.table);
1970 let alias = self.plan.intern(&label);
1971 let columns = self.plan.add_fields(&fields);
1972 if let Some(zones) = table.rows().zones().filter(|_| !excluded) {
1977 self.plan.set_zones(index, zones);
1978 }
1979 if let Some(frequencies) = table.frequencies().filter(|_| !excluded) {
1980 self.plan.set_frequencies(index, frequencies);
1981 }
1982 for (column, distinct) in table.distincts() {
1983 if !excluded {
1984 self.plan.measure_distinct(index, &column, distinct);
1985 }
1986 }
1987 if self.want_ascending && !excluded {
1988 for column in table.ascending() {
1989 self.plan.mark_ascending(index, &column);
1990 }
1991 }
1992 if !excluded {
1993 for (column, bytes) in table.widths() {
1994 self.plan.measure_width(index, &column, bytes);
1995 }
1996 }
1997 let node = self.add_node(Node::Get {
1998 catalog: catalog_name,
1999 schema,
2000 table: table_name,
2001 alias,
2002 index,
2003 columns,
2004 });
2005 Ok((node, scope))
2006 }
2007
2008 fn bind_view(
2020 &mut self,
2021 ast: &Ast,
2022 name: &QualifiedName,
2023 alias: ast::StrRef,
2024 columns: ast::Slice,
2025 ) -> Result<(NodeRef, Scope)> {
2026 let view = self.catalog.view(name)?;
2027 let full = name.to_string();
2028 if self.expanding.contains(&full) {
2029 return Err(Error::binder(format!(
2033 "infinite recursion detected: attempting to recursively bind view \"\"{}\"\"",
2034 name.table
2035 )));
2036 }
2037 let body = parse_ast_with_case(view.sql(), self.semantics.identifier_case())?;
2038 let query = match body.statements.as_slice() {
2039 [ast::Statement::Query(query)] => *query,
2040 _ => return Err(Error::binder(format!("view \"{}\" is not a query", name.table))),
2043 };
2044 self.expanding.push(full);
2045 let bound = self.bind_query(&body, query);
2046 self.expanding.pop();
2047 let (node, mut scope) = bound?;
2048
2049 let aliases: Vec<&str> = view.aliases().iter().map(String::as_str).collect();
2050 if !aliases.is_empty() {
2051 scope.rename(&aliases, "unnamed_subquery")?;
2052 }
2053 view.remember(scope.fields());
2060 let label = if alias == NONE { name.table.clone() } else { ast.string(alias).to_string() };
2061 scope.relabel(&label);
2062 if !columns.is_empty() {
2063 let names: Vec<&str> = ast.name(columns).collect();
2064 scope.rename(&names, &label)?;
2065 }
2066 Ok((node, scope))
2067 }
2068
2069 fn bind_table_function(
2077 &mut self,
2078 ast: &Ast,
2079 name: ast::Slice,
2080 args: ast::Slice,
2081 alias: ast::StrRef,
2082 columns: ast::Slice,
2083 pragma: bool,
2084 ) -> Result<(NodeRef, Scope)> {
2085 let renamed = columns;
2088 let parts: Vec<&str> = ast.name(name).collect();
2089 let function_name = *parts.last().unwrap_or(&"");
2093 if let Some(schema) = parts.iter().rev().nth(1) {
2094 if !schema.eq_ignore_ascii_case("main") && !schema.eq_ignore_ascii_case("system") {
2095 return Err(Error::catalog(format!(
2096 "Table Function with name {} does not exist!",
2097 parts.join(".")
2098 )));
2099 }
2100 }
2101 let Some(called) = TableFunction::lookup(function_name) else {
2105 if pragma {
2106 if args.is_empty() && self.catalog.resolve(&parts).is_ok() {
2112 return self.bind_table(ast, name, alias, columns);
2113 }
2114 let spelled = function_name.strip_prefix("pragma_").unwrap_or(function_name);
2115 return Err(Error::catalog(format!(
2116 "Pragma Function with name {spelled} does not exist!"
2117 )));
2118 }
2119 return Err(Error::catalog(format!(
2120 "Table Function with name {function_name} does not exist!"
2121 )));
2122 };
2123 let written = ast.target_list(args).to_vec();
2124 let empty = Scope::empty();
2125 let previous = std::mem::replace(&mut self.clause, "table function arguments");
2126 let mut bound = Vec::new();
2127 let mut written_options = Vec::new();
2128 for argument in written {
2129 let expr = self.bind_expr(ast, argument.expr, &empty)?;
2130 if argument.alias == NONE {
2131 bound.push(expr);
2132 } else {
2133 let name = ast.string(argument.alias).to_string();
2134 let (parameter, value) = self.named_argument(called, &name, expr)?;
2135 written_options.push((parameter, value, expr));
2136 }
2137 }
2138 self.clause = previous;
2139 let options = Options::of(&written_options)?;
2140
2141 let given: Vec<LogicalType> =
2144 bound.iter().map(|&expr| self.plan.expr_type(expr).clone()).collect();
2145 let resolved = if pragma {
2146 resolve_pragma(function_name, &given)?
2147 } else {
2148 resolve_table(function_name, &given)?
2149 };
2150 let mut cast: Vec<ExprRef> = bound
2151 .iter()
2152 .zip(&resolved.arguments)
2153 .map(|(&expr, ty)| self.checked_cast_to(expr, ty, false))
2154 .collect::<Result<_>>()?;
2155
2156 if resolved.function.answered_when_bound() {
2157 let Columns::Fixed(fields) = resolved.columns else {
2158 return Err(Error::internal("a pragma that resolved to a file"));
2159 };
2160 let [argument] = cast[..] else {
2161 return Err(Error::internal("a pragma that resolved to more than one name"));
2162 };
2163 return self.bind_pragma(ast, resolved.function, &fields, argument, alias, columns);
2164 }
2165 let mut measured = Stat::Unknown;
2168 let mut counted: Vec<(String, Stat<u64>)> = Vec::new();
2169 let mut bounded: Option<Arc<dyn Zones>> = None;
2170 let fields = match resolved.columns {
2171 Columns::Fixed(fields) => fields,
2172 columns => {
2173 let paths = self.file_paths(cast[0], resolved.function.name())?;
2178 let mut mirrorable = None;
2179 if resolved.function == TableFunction::ReadParquet && !options.file_row_number {
2180 if let Some((path, stamp)) = mirror_target(&paths) {
2181 if let Some(name) =
2182 self.catalog.mirror(&path, options.binary_as_string, stamp)
2183 {
2184 let name = name.clone();
2185 let label = if alias == NONE {
2186 resolved.function.name().to_string()
2187 } else {
2188 ast.string(alias).to_string()
2189 };
2190 return self.bind_catalog_table(ast, &name, label, renamed);
2191 }
2192 mirrorable = Some(path);
2193 }
2194 }
2195 let mut fields = match columns {
2196 Columns::Csv => csv_fields(&paths, options.given)?,
2199 _ => {
2200 let footers = self.footers(&paths, mirrorable.as_deref())?;
2201 if let Some(path) = mirrorable.as_deref() {
2202 self.want_mirror(path, options.binary_as_string, &footers.rows);
2203 }
2204 measured = footers.rows;
2205 counted = footers.distincts;
2206 bounded = footers.zones;
2207 footers.fields
2208 }
2209 };
2210 if options.all_varchar {
2211 for field in &mut fields {
2216 field.ty = LogicalType::Varchar;
2217 }
2218 }
2219 if options.binary_as_string {
2220 for field in &mut fields {
2225 if field.ty == LogicalType::Blob {
2226 field.ty = LogicalType::Varchar;
2227 }
2228 }
2229 }
2230 if options.file_row_number {
2231 if fields.iter().any(|field| field.name == FILE_ROW_NUMBER) {
2237 return Err(Error::binder(format!(
2238 "Duplicate column name \"{FILE_ROW_NUMBER}\": the file already has a \
2239 column of that name, so file_row_number cannot add one"
2240 )));
2241 }
2242 fields.push(Field::required(FILE_ROW_NUMBER.to_string(), LogicalType::BigInt));
2243 }
2244 cast = paths.iter().map(|path| self.path_constant(path)).collect();
2245 fields
2246 }
2247 };
2248 let label = if alias == NONE {
2249 resolved.function.name().to_string()
2250 } else {
2251 ast.string(alias).to_string()
2252 };
2253 let names: Vec<&str> = ast.name(columns).collect();
2254 self.table_function_source(
2255 resolved.function,
2256 &cast,
2257 &written_options,
2258 Read { fields, rows: measured, distincts: counted, zones: bounded },
2259 &label,
2260 &names,
2261 )
2262 }
2263
2264 fn bind_pragma(
2277 &mut self,
2278 ast: &Ast,
2279 function: TableFunction,
2280 fields: &[Field],
2281 argument: ExprRef,
2282 alias: ast::StrRef,
2283 columns: ast::Slice,
2284 ) -> Result<(NodeRef, Scope)> {
2285 let written = self.pragma_name(argument, function)?;
2286 let parts = identifier_parts(&written);
2287 let spelled: Vec<&str> = parts.iter().map(String::as_str).collect();
2288 let name = self.catalog.resolve(&spelled)?;
2289 let described = self.described(ast, &name)?;
2290 let mut rows = Vec::with_capacity(described.len());
2291 for (at, field) in described.iter().enumerate() {
2292 let items = if matches!(function, TableFunction::PragmaShow) {
2293 self.describing(field)
2294 } else {
2295 self.table_info(at, field)
2296 };
2297 rows.push(self.plan.add_expr_list(&items));
2298 }
2299 let rows = self.plan.add_rows(&rows);
2300 let held = self.plan.add_fields(fields);
2301 let index = self.fresh_index();
2302 let node = self.add_node(Node::Values { index, columns: held, rows });
2303 let label =
2304 if alias == NONE { function.name().to_string() } else { ast.string(alias).to_string() };
2305 let mut scope = Scope::empty();
2306 for (at, field) in fields.iter().enumerate() {
2307 scope.push(Visible {
2308 table: label.clone(),
2309 name: field.name.clone(),
2310 binding: ColumnBinding::new(index, at as u32),
2311 ty: field.ty.clone(),
2312 not_null: false,
2313 key: None,
2314 default: None,
2315 qualified: false,
2316 also: None,
2317 });
2318 }
2319 if !columns.is_empty() {
2320 let names: Vec<&str> = ast.name(columns).collect();
2321 scope.rename(&names, &label)?;
2322 }
2323 Ok((node, scope))
2324 }
2325
2326 fn pragma_name(&self, argument: ExprRef, function: TableFunction) -> Result<String> {
2336 let Expr::Constant(reference) = *self.plan.expr(argument) else {
2337 return Err(Error::not_implemented(format!(
2338 "{}() given a name that is not a constant",
2339 function.name()
2340 )));
2341 };
2342 match self.plan.value(reference) {
2343 Value::Varchar(name) => Ok(name.clone()),
2344 Value::Null => Ok("NULL".to_string()),
2345 other => {
2346 Err(Error::internal(format!("a pragma name bound as VARCHAR arrived as {other}")))
2347 }
2348 }
2349 }
2350
2351 fn described(&mut self, ast: &Ast, name: &QualifiedName) -> Result<Vec<Field>> {
2362 if self.catalog.entry(name)? == Entry::Table {
2363 return Ok(self.catalog.table(name)?.columns().to_vec());
2364 }
2365 let (_, scope) = self.bind_view(ast, name, NONE, ast::Slice::default())?;
2366 Ok(scope.fields())
2367 }
2368
2369 fn describing(&mut self, field: &Field) -> Vec<ExprRef> {
2371 let written = [
2372 field.name.clone(),
2373 field.ty.to_string(),
2374 if field.not_null { "NO" } else { "YES" }.to_owned(),
2375 ];
2376 let mut items: Vec<ExprRef> =
2377 written.into_iter().map(|text| self.plan.add_constant(Value::Varchar(text))).collect();
2378 for _ in 0..3 {
2379 let empty = self.plan.add_constant(Value::Null);
2380 items.push(self.cast_to(empty, &LogicalType::Varchar));
2381 }
2382 items
2383 }
2384
2385 fn table_info(&mut self, at: usize, field: &Field) -> Vec<ExprRef> {
2391 let cid = self.plan.add_constant(Value::Integer(i32::try_from(at).unwrap_or(i32::MAX)));
2392 let name = self.plan.add_constant(Value::Varchar(field.name.clone()));
2393 let ty = self.plan.add_constant(Value::Varchar(field.ty.to_string()));
2394 let not_null = self.plan.add_constant(Value::Boolean(field.not_null));
2395 let default = self.plan.add_constant(Value::Null);
2396 let default = self.cast_to(default, &LogicalType::Varchar);
2397 let key = self.plan.add_constant(Value::Boolean(false));
2398 vec![cid, name, ty, not_null, default, key]
2399 }
2400
2401 fn named_argument(
2415 &mut self,
2416 function: TableFunction,
2417 name: &str,
2418 expr: ExprRef,
2419 ) -> Result<(&'static str, Value)> {
2420 let known = function
2421 .parameters()
2422 .iter()
2423 .find(|(parameter, _)| parameter.eq_ignore_ascii_case(name));
2424 let Some((parameter, wanted)) = known else {
2425 let candidates: Vec<String> = function
2426 .parameters()
2427 .iter()
2428 .map(|(parameter, ty)| format!(" {parameter} {ty}"))
2429 .collect();
2430 return Err(Error::binder(format!(
2431 "Invalid named parameter \"{name}\" for function {}\nCandidates:\n{}\n",
2432 function.name(),
2433 candidates.join("\n")
2434 )));
2435 };
2436 let Expr::Constant(reference) = *self.plan.expr(expr) else {
2437 return Err(Error::not_implemented(format!(
2438 "the named parameter {parameter} with a value that is not a constant"
2439 )));
2440 };
2441 let value = self.plan.value(reference).clone();
2442 if value == Value::Null {
2443 return Err(Error::binder(null_parameter(function, parameter)));
2444 }
2445 let given = self.plan.expr_type(expr).clone();
2446 if given != *wanted {
2447 return Err(Error::not_implemented(format!(
2448 "the named parameter {parameter} given a {given} where a {wanted} was wanted"
2449 )));
2450 }
2451 Ok((parameter, value))
2452 }
2453
2454 fn bind_replacement_scan(
2465 &mut self,
2466 ast: &Ast,
2467 parts: &[&str],
2468 alias: ast::StrRef,
2469 columns: ast::Slice,
2470 missing: Error,
2471 ) -> Result<(NodeRef, Scope)> {
2472 let [path] = parts else { return Err(missing) };
2473 let path = *path;
2474 let extension = path.rsplit_once('.').map(|(_, after)| after).unwrap_or_default();
2475 let Some(function) = Self::reader_for(extension) else {
2476 if is_file(path) {
2477 return Err(Error::binder(format!(
2482 "No extension found that is capable of reading the file \"{path}\"\n* If this \
2483 file is a supported file format you can explicitly use the reader functions, \
2484 such as read_csv, read_json or read_parquet"
2485 )));
2486 }
2487 return Err(missing);
2488 };
2489 let paths = files(path)?;
2494 let label = if alias == NONE {
2500 if is_pattern(path) {
2501 path.to_string()
2502 } else {
2503 let file = path.rsplit_once('/').map_or(path, |(_, file)| file);
2504 file.rsplit_once('.').map_or(file, |(stem, _)| stem).to_string()
2505 }
2506 } else {
2507 ast.string(alias).to_string()
2508 };
2509 let mut mirrorable = None;
2510 if function == TableFunction::ReadParquet {
2511 if let Some((canonical, stamp)) = mirror_target(&paths) {
2512 if let Some(name) = self.catalog.mirror(&canonical, false, stamp) {
2513 let name = name.clone();
2514 return self.bind_catalog_table(ast, &name, label, columns);
2515 }
2516 mirrorable = Some(canonical);
2517 }
2518 }
2519 let read = match function {
2520 TableFunction::ReadParquet => {
2521 let footers = self.footers(&paths, mirrorable.as_deref())?;
2522 if let Some(canonical) = mirrorable.as_deref() {
2523 self.want_mirror(canonical, false, &footers.rows);
2524 }
2525 Read {
2526 fields: footers.fields,
2527 rows: footers.rows,
2528 distincts: footers.distincts,
2529 zones: footers.zones,
2530 }
2531 }
2532 _ => Read::uncounted(csv_fields(&paths, Given::default())?),
2533 };
2534 let arguments: Vec<ExprRef> = paths.iter().map(|path| self.path_constant(path)).collect();
2535 let names: Vec<&str> = ast.name(columns).collect();
2536 self.table_function_source(function, &arguments, &[], read, &label, &names)
2537 }
2538
2539 fn footers(&self, paths: &[String], mirrorable: Option<&str>) -> Result<Footers> {
2545 if let Some(path) = mirrorable.filter(|_| self.outlined) {
2546 let outline = parquet_outline(path)?;
2547 if outline.rows.value().is_some() {
2548 return Ok(outline);
2549 }
2550 }
2551 parquet_footers(paths)
2552 }
2553
2554 fn want_mirror(&mut self, path: &str, binary_as_string: bool, rows: &Stat<u64>) {
2558 if let Some(&rows) = rows.value() {
2559 self.plan.want_mirror(path, binary_as_string, rows);
2560 }
2561 }
2562
2563 fn path_constant(&mut self, path: &str) -> ExprRef {
2565 let value = self.plan.add_value(Value::Varchar(path.to_string()));
2566 self.plan.add_expr(Expr::Constant(value), LogicalType::Varchar)
2567 }
2568
2569 fn reader_for(extension: &str) -> Option<TableFunction> {
2576 if extension.eq_ignore_ascii_case("parquet") {
2577 return Some(TableFunction::ReadParquet);
2578 }
2579 if extension.eq_ignore_ascii_case("csv") || extension.eq_ignore_ascii_case("tsv") {
2580 return Some(TableFunction::ReadCsv);
2581 }
2582 None
2583 }
2584
2585 fn table_function_source(
2595 &mut self,
2596 function: TableFunction,
2597 args: &[ExprRef],
2598 written: &[(&'static str, Value, ExprRef)],
2599 read: Read,
2600 label: &str,
2601 names: &[&str],
2602 ) -> Result<(NodeRef, Scope)> {
2603 let Read { fields, rows, distincts, zones } = read;
2604 let index = self.fresh_index();
2605 if rows.is_known() {
2610 self.plan.measure(index, rows);
2611 }
2612 for (column, distinct) in distincts {
2613 self.plan.measure_distinct(index, &column, distinct);
2614 }
2615 if let Some(zones) = zones {
2616 self.plan.set_zones(index, zones);
2617 }
2618 let mut scope = Scope::empty();
2619 for (at, field) in fields.iter().enumerate() {
2620 scope.push(Visible {
2621 table: label.to_string(),
2622 name: field.name.clone(),
2623 binding: ColumnBinding::new(index, at as u32),
2624 ty: field.ty.clone(),
2625 not_null: false,
2628 key: None,
2629 default: None,
2630 qualified: false,
2631 also: None,
2632 });
2633 }
2634 if !names.is_empty() {
2635 scope.rename(names, label)?;
2636 } else if matches!(function, TableFunction::Range | TableFunction::GenerateSeries) {
2637 for column in &mut scope.columns {
2640 column.also = Some(std::mem::replace(&mut column.name, label.to_string()));
2641 }
2642 }
2643 let function = self.plan.intern(function.name());
2644 let args = self.plan.add_expr_list(args);
2645 let named: Vec<u32> =
2646 written.iter().map(|(parameter, _, _)| self.plan.intern(parameter)).collect();
2647 let settings: Vec<ExprRef> = written.iter().map(|(_, _, expr)| *expr).collect();
2648 let options = self.plan.add_name_list(&named);
2649 let settings = self.plan.add_expr_list(&settings);
2650 let columns = self.plan.add_fields(&fields);
2651 let node = self.add_node(Node::TableFunction {
2652 index,
2653 function,
2654 args,
2655 options,
2656 settings,
2657 columns,
2658 });
2659 Ok((node, scope))
2660 }
2661
2662 fn file_paths(&self, expr: ExprRef, name: &str) -> Result<Vec<String>> {
2669 let mut paths = Vec::new();
2670 for pattern in self.file_patterns(expr, name)? {
2671 paths.extend(files(&pattern)?);
2672 }
2673 Ok(paths)
2674 }
2675
2676 fn file_patterns(&self, expr: ExprRef, name: &str) -> Result<Vec<String>> {
2695 let Some(value) = fold::value_of(&self.plan, expr)? else {
2696 return Err(Error::not_implemented(
2697 "a table function file name that is not a constant",
2698 ));
2699 };
2700 match value {
2701 Value::Varchar(path) => Ok(vec![path]),
2702 Value::Null => Err(Error::parser(format!("{name} cannot take NULL list as parameter"))),
2704 Value::List { values, .. } if values.is_empty() => {
2709 Err(Error::io(format!("\"{name}\" needs at least one file to read")))
2710 }
2711 Value::List { values, .. } => values
2712 .iter()
2713 .map(|value| match value {
2714 Value::Varchar(path) => Ok(path.clone()),
2715 _ => Err(Error::parser(format!(
2716 "{name} reader cannot take NULL input as parameter"
2717 ))),
2718 })
2719 .collect(),
2720 other => {
2721 Err(Error::internal(format!("a file name bound as VARCHAR arrived as {other}")))
2722 }
2723 }
2724 }
2725
2726 fn side_of(
2744 &self,
2745 pending: &PendingSubquery,
2746 left_tables: &[u32],
2747 right_tables: &[u32],
2748 ) -> Option<Side> {
2749 let mut needs_left = false;
2750 let mut needs_right = false;
2751 let mut note = |binding: ColumnBinding| {
2752 needs_left |= left_tables.contains(&binding.table);
2753 needs_right |= right_tables.contains(&binding.table);
2754 };
2755 for &binding in &pending.reads {
2756 note(binding);
2757 }
2758 for &condition in &pending.conditions {
2763 self.plan.read_columns(condition, &mut |_, binding| note(binding));
2764 }
2765 match (needs_left, needs_right) {
2766 (true, true) => None,
2767 (_, true) => Some(Side::Right),
2768 _ => Some(Side::Left),
2769 }
2770 }
2771
2772 #[allow(clippy::too_many_arguments)]
2793 fn bind_pair_dependent_join(
2794 &mut self,
2795 kind: ast::JoinKind,
2796 independent: bool,
2797 left: NodeRef,
2798 right: NodeRef,
2799 pair: Vec<PendingSubquery>,
2800 conditions: Vec<ExprRef>,
2801 scope: Scope,
2802 ) -> Result<(NodeRef, Scope)> {
2803 if kind != ast::JoinKind::Inner {
2804 return Err(Error::not_implemented(
2805 "a subquery that reads both sides of that join, written in the condition of a join \
2806 that is not an inner join"
2807 .to_string(),
2808 ));
2809 }
2810 if !independent {
2813 return Err(Error::not_implemented(
2814 "a subquery that reads both sides of that join, written in the condition of a join \
2815 whose right side is lateral"
2816 .to_string(),
2817 ));
2818 }
2819 let mut node = self.add_node(Node::CrossProduct { left, right });
2820 for pending in pair {
2821 node = self.attach_subquery(node, pending);
2822 }
2823 let mut conditions = conditions.into_iter();
2827 let mut predicate = conditions.next().expect("a join condition was bound");
2828 for next in conditions {
2829 let children = self.plan.add_expr_list(&[predicate, next]);
2830 let conjunction = Expr::Conjunction { op: ConjunctionOp::And, children };
2831 predicate = self.plan.add_expr(conjunction, LogicalType::Boolean);
2832 }
2833 let node = self.add_node(Node::Filter { input: node, predicate });
2834 Ok((node, scope))
2835 }
2836
2837 #[allow(clippy::too_many_arguments)]
2838 fn bind_join(
2839 &mut self,
2840 ast: &Ast,
2841 left: ast::SourceRef,
2842 right: ast::SourceRef,
2843 kind: ast::JoinKind,
2844 natural: bool,
2845 on: ast::ExprRef,
2846 using: ast::Slice,
2847 ) -> Result<(NodeRef, Scope)> {
2848 let (left_node, left_scope) = self.bind_source(ast, left)?;
2849 let (right_node, right_scope, correlated) = self.bind_lateral(ast, right, &left_scope)?;
2850 if !correlated.is_empty()
2854 && !matches!(kind, ast::JoinKind::Inner | ast::JoinKind::Cross | ast::JoinKind::Left)
2855 {
2856 return Err(Error::binder(
2857 "The combining JOIN type must be INNER or LEFT for a LATERAL reference",
2858 ));
2859 }
2860 let split = left_scope.len();
2861 let left_tables: Vec<u32> =
2867 left_scope.columns.iter().map(|column| column.binding.table).collect();
2868 let right_tables: Vec<u32> =
2869 right_scope.columns.iter().map(|column| column.binding.table).collect();
2870 let mut scope = left_scope.concat(right_scope);
2871
2872 let merged: Vec<String> = if natural {
2875 let mut names = Vec::new();
2876 for (at, column) in scope.columns.iter().enumerate().take(split) {
2877 if scope.columns[split..].iter().any(|right| same_name(&right.name, &column.name))
2878 && !names.iter().any(|held: &String| same_name(held, &column.name))
2879 {
2880 let _ = at;
2881 names.push(column.name.clone());
2882 }
2883 }
2884 names
2885 } else {
2886 let mut names: Vec<String> = Vec::new();
2892 for name in ast.name(using) {
2893 if !names.iter().any(|held| same_name(held, name)) {
2894 names.push(name.to_string());
2895 }
2896 }
2897 names
2898 };
2899
2900 let mut conditions = Vec::new();
2901 let mut dropped = Vec::new();
2902 for name in &merged {
2903 let left_at = scope.columns[..split]
2904 .iter()
2905 .position(|column| same_name(&column.name, name))
2906 .ok_or_else(|| {
2907 Error::binder(format!(
2908 "column \"{name}\" specified in USING clause does not exist in left table"
2909 ))
2910 })?;
2911 let right_at = scope.columns[split..]
2912 .iter()
2913 .position(|column| same_name(&column.name, name))
2914 .map(|at| at + split)
2915 .ok_or_else(|| {
2916 Error::binder(format!(
2917 "column \"{name}\" specified in USING clause does not exist in right table"
2918 ))
2919 })?;
2920 let left_column = &scope.columns[left_at];
2921 let (left_binding, left_type) = (left_column.binding, left_column.ty.clone());
2922 let right_column = &scope.columns[right_at];
2923 let (right_binding, right_type) = (right_column.binding, right_column.ty.clone());
2924 let left_expr = self.plan.add_expr(Expr::Column(left_binding), left_type);
2925 let right_expr = self.plan.add_expr(Expr::Column(right_binding), right_type);
2926 conditions.push(self.compare(rudb_plan::CompareOp::Equal, left_expr, right_expr)?);
2927 dropped.push(right_at);
2928 }
2929 dropped.sort_unstable();
2932 for at in dropped.into_iter().rev() {
2933 scope.remove(at);
2934 }
2935
2936 let mut left_node = left_node;
2937 let mut right_node = right_node;
2938 let mut pair = Vec::new();
2939 if on != NONE {
2940 if !merged.is_empty() {
2941 return Err(Error::binder("a join cannot have both ON and USING"));
2942 }
2943 self.clause = "JOIN condition";
2944 let waiting = self.scalar_subqueries.len();
2945 let predicate = self.bind_expr(ast, on, &scope)?;
2946 conditions.push(self.as_boolean(predicate, "JOIN")?);
2947 for pending in self.scalar_subqueries.split_off(waiting) {
2948 match self.side_of(&pending, &left_tables, &right_tables) {
2949 Some(Side::Right) => right_node = self.attach_subquery(right_node, pending),
2950 Some(Side::Left) => left_node = self.attach_subquery(left_node, pending),
2951 None => pair.push(pending),
2952 }
2953 }
2954 }
2955
2956 if kind == ast::JoinKind::Cross && !conditions.is_empty() {
2957 return Err(Error::binder("a CROSS JOIN cannot have a condition"));
2958 }
2959 if !pair.is_empty() {
2960 return self.bind_pair_dependent_join(
2961 kind,
2962 correlated.is_empty(),
2963 left_node,
2964 right_node,
2965 pair,
2966 conditions,
2967 scope,
2968 );
2969 }
2970 if correlated.is_empty()
2974 && conditions.is_empty()
2975 && matches!(kind, ast::JoinKind::Cross | ast::JoinKind::Inner)
2976 {
2977 let node = self.add_node(Node::CrossProduct { left: left_node, right: right_node });
2978 return Ok((node, scope));
2979 }
2980 if matches!(kind, ast::JoinKind::Semi | ast::JoinKind::Anti) {
2989 scope.truncate(split);
2990 }
2991 let kind = match kind {
2992 ast::JoinKind::Inner | ast::JoinKind::Cross => JoinKind::Inner,
2993 ast::JoinKind::Left => JoinKind::Left,
2994 ast::JoinKind::Right => JoinKind::Right,
2995 ast::JoinKind::Full => JoinKind::Full,
2996 ast::JoinKind::Semi => JoinKind::Semi,
2997 ast::JoinKind::Anti => JoinKind::Anti,
2998 ast::JoinKind::Positional => JoinKind::Positional,
2999 };
3000 let conditions = self.plan.add_expr_list(&conditions);
3001 let node = if correlated.is_empty() {
3002 self.add_node(Node::Join {
3003 left: left_node,
3004 right: right_node,
3005 kind,
3006 conditions,
3007 build: BuildSide::default(),
3008 })
3009 } else {
3010 self.add_node(Node::DependentJoin {
3011 left: left_node,
3012 right: right_node,
3013 kind,
3014 conditions,
3015 })
3016 };
3017 Ok((node, scope))
3018 }
3019
3020 fn bind_filter(
3028 &mut self,
3029 ast: &Ast,
3030 filter: ast::ExprRef,
3031 scope: &Scope,
3032 ) -> Result<Option<ExprRef>> {
3033 if filter == NONE {
3034 return Ok(None);
3035 }
3036 let bound = self.bind_expr(ast, filter, scope)?;
3037 Ok(Some(self.checked_cast_to(bound, &LogicalType::Boolean, false)?))
3038 }
3039
3040 pub(crate) fn bind_aggregate(
3045 &mut self,
3046 ast: &Ast,
3047 name: &str,
3048 args: &[ast::ExprRef],
3049 distinct: bool,
3050 filter: ast::ExprRef,
3051 scope: &Scope,
3052 ) -> Result<ExprRef> {
3053 let frames = std::mem::take(&mut self.lambda_frames);
3054 let bound = self.bind_aggregate_over_rows(ast, name, args, distinct, filter, scope);
3055 self.lambda_frames = frames;
3056 bound
3057 }
3058
3059 fn bind_aggregate_over_rows(
3060 &mut self,
3061 ast: &Ast,
3062 name: &str,
3063 args: &[ast::ExprRef],
3064 distinct: bool,
3065 filter: ast::ExprRef,
3066 scope: &Scope,
3067 ) -> Result<ExprRef> {
3068 if self.in_filter {
3069 return Err(Error::binder("aggregate functions are not allowed in FILTER"));
3070 }
3071 if self.in_aggregate {
3072 return Err(Error::binder(format!(
3073 "aggregate function calls cannot be nested, and {name}() is inside one"
3074 )));
3075 }
3076 if self.aggregation.is_none() {
3077 return Err(Error::binder(format!(
3078 "aggregate function calls cannot be used in the {}",
3079 self.clause
3080 )));
3081 }
3082 self.in_aggregate = true;
3087 self.in_filter = true;
3088 let filter = self.bind_filter(ast, filter, scope);
3089 self.in_filter = false;
3090 self.in_aggregate = false;
3091 let filter = filter?;
3092
3093 self.in_aggregate = true;
3094 let mut bound = Vec::with_capacity(args.len());
3095 let mut failure = None;
3096 for &arg in args {
3097 match self.bind_expr(ast, arg, scope) {
3098 Ok(expr) => bound.push(expr),
3099 Err(error) => {
3100 failure = Some(error);
3101 break;
3102 }
3103 }
3104 }
3105 self.in_aggregate = false;
3106 if let Some(error) = failure {
3107 return Err(error);
3108 }
3109
3110 let types: Vec<LogicalType> =
3111 bound.iter().map(|&arg| self.plan.expr_type(arg).clone()).collect();
3112 let resolved = resolve(name, &types)?;
3113 if resolved.name == "string_agg"
3116 && bound.len() == 2
3117 && !matches!(fold::value_of(&self.plan, bound[1]), Ok(Some(_)))
3118 {
3119 return Err(Error::binder(
3120 "The \"separator\" argument in function \"string_agg\" must be a constant expression",
3121 ));
3122 }
3123 let mut cast = Vec::with_capacity(bound.len());
3124 for (arg, wanted) in bound.iter().zip(&resolved.arguments) {
3125 cast.push(self.checked_cast_to(*arg, wanted, false)?);
3126 }
3127 let args = self.plan.add_expr_list(&cast);
3128 let name = self.plan.intern(resolved.name);
3129 let ty = resolved.returns;
3130 let call = self.plan.add_expr(Expr::Aggregate { name, args, distinct, filter }, ty.clone());
3131
3132 let existing = self.aggregation.as_ref().map(|held| held.aggregates.clone());
3135 let existing = existing.unwrap_or_default();
3136 let at = match existing.iter().position(|&held| self.same_expr(held, call)) {
3137 Some(at) => at,
3138 None => {
3139 let aggregation = self.aggregation.as_mut().expect("checked above");
3140 aggregation.aggregates.push(call);
3141 aggregation.aggregates.len() - 1
3142 }
3143 };
3144 let aggregation = self.aggregation.as_ref().expect("checked above");
3145 let (index, groups) = (aggregation.index, aggregation.groups.len());
3146 Ok(self.column(index, groups + at, ty))
3147 }
3148
3149 pub(crate) fn bind_window(
3160 &mut self,
3161 ast: &Ast,
3162 written: &WindowCall<'_>,
3163 scope: &Scope,
3164 ) -> Result<ExprRef> {
3165 let frames = std::mem::take(&mut self.lambda_frames);
3166 let bound = self.bind_window_over_rows(ast, written, scope);
3167 self.lambda_frames = frames;
3168 bound
3169 }
3170
3171 fn bind_window_over_rows(
3172 &mut self,
3173 ast: &Ast,
3174 written: &WindowCall<'_>,
3175 scope: &Scope,
3176 ) -> Result<ExprRef> {
3177 let WindowCall { name, args, distinct, filter, ignore_nulls, spec, .. } = *written;
3178 if self.in_aggregate {
3179 return Err(Error::binder(
3180 "aggregate function calls cannot contain window function calls",
3181 ));
3182 }
3183 if self.in_window {
3184 return Err(Error::binder("window function calls cannot be nested"));
3185 }
3186 let clause = if self.clause == "JOIN condition" { "WHERE clause" } else { self.clause };
3190 if clause != "SELECT clause" && clause != "ORDER BY clause" {
3191 return Err(Error::binder(format!("{clause} cannot contain window functions!")));
3192 }
3193
3194 let starred = args.iter().any(|&arg| {
3198 matches!(ast.expr(arg), ast::Expr::Star { qualifier, replacements }
3199 if qualifier.is_empty() && replacements.is_empty())
3200 });
3201 let (name, args): (&str, &[ast::ExprRef]) = if starred {
3202 if !same_name(name, "count") || args.len() != 1 {
3203 return Err(Error::binder(format!("* is not allowed in {name}()")));
3204 }
3205 ("count_star", &[])
3206 } else if same_name(name, "count") && args.is_empty() {
3207 ("count_star", &[])
3210 } else {
3211 (name, args)
3212 };
3213
3214 let held = ast.window(spec);
3215 self.in_window = true;
3216 let parts = self.window_parts(ast, written, args, held, scope);
3217 let filter = if parts.is_ok() { self.bind_filter(ast, filter, scope) } else { Ok(None) };
3222 self.in_window = false;
3223 let parts = parts?;
3224 let filter = filter?;
3225 let offsets = [parts.frame.start, parts.frame.end]
3228 .iter()
3229 .any(|end| matches!(end, WindowBound::Preceding(_) | WindowBound::Following(_)));
3230 if parts.frame.unit == WindowUnit::Range && offsets && parts.order.len() != 1 {
3231 return Err(Error::binder("RANGE frames must have only one ORDER BY expression"));
3232 }
3233
3234 let types: Vec<LogicalType> =
3235 parts.args.iter().map(|&arg| self.plan.expr_type(arg).clone()).collect();
3236 let resolved = window_signature(name, &types)?;
3237 if resolved.name == "fill" {
3240 let keys: Vec<LogicalType> =
3241 parts.order.iter().map(|key| self.plan.expr_type(key.expr).clone()).collect();
3242 refuse_fill(&types[0], &keys, distinct, ignore_nulls)?;
3243 }
3244 if distinct && kind_of(resolved.name) == Some(FunctionKind::Window) {
3248 return Err(Error::binder(format!(
3249 "DISTINCT is not implemented for the window function \"\"{name}\"\""
3250 )));
3251 }
3252 if filter.is_some() && kind_of(resolved.name) == Some(FunctionKind::Window) {
3255 return Err(Error::binder(format!(
3256 "FILTER is not implemented for the window function \"\"{name}\"\""
3257 )));
3258 }
3259 if !parts.inner.is_empty() && kind_of(resolved.name) == Some(FunctionKind::Window) {
3267 let counts = matches!(resolved.name, "first_value" | "last_value" | "nth_value");
3268 if !counts {
3269 if parts.frame.exclude != WindowExclude::NoOthers {
3270 return Err(Error::binder(format!(
3271 "EXCLUDE is not supported for the window function \"\"{}\"\"",
3272 resolved.name
3273 )));
3274 }
3275 return Err(Error::not_implemented(format!(
3276 "ORDER BY inside the arguments of the window function \"{}\"",
3277 resolved.name
3278 )));
3279 }
3280 }
3281 let mut cast = Vec::with_capacity(parts.args.len());
3282 for (arg, wanted) in parts.args.iter().zip(&resolved.arguments) {
3283 cast.push(self.checked_cast_to(*arg, wanted, false)?);
3284 }
3285 let args = self.plan.add_expr_list(&cast);
3286 let order = self.plan.add_sort_keys(&parts.inner);
3287 let name = self.plan.intern(resolved.name);
3288 let ty = resolved.returns;
3289 let call = self.plan.add_expr(
3290 Expr::Window { name, args, distinct, filter, ignore_nulls, order },
3291 ty.clone(),
3292 );
3293
3294 let at = self.window_run(parts.partition, parts.order, parts.frame, call);
3295 let index = self.windows.last().expect("the run was just filed").index;
3296 Ok(self.column(index, at, ty))
3297 }
3298
3299 fn window_run(
3306 &mut self,
3307 partition: Vec<ExprRef>,
3308 order: Vec<SortKey>,
3309 frame: WindowFrame,
3310 call: ExprRef,
3311 ) -> usize {
3312 let matches = self.windows.last().is_some_and(|run| {
3313 run.frame == frame
3314 && run.partition.len() == partition.len()
3315 && run.order.len() == order.len()
3316 && run.partition.iter().zip(&partition).all(|(&l, &r)| self.same_expr(l, r))
3317 && run.order.iter().zip(&order).all(|(l, r)| {
3318 l.descending == r.descending
3319 && l.nulls_first == r.nulls_first
3320 && self.same_expr(l.expr, r.expr)
3321 })
3322 });
3323 if !matches {
3324 let index = self.fresh_index();
3325 self.windows.push(WindowRun { index, partition, order, frame, calls: Vec::new() });
3326 }
3327 let calls = self.windows.last().expect("a run is open").calls.clone();
3330 if let Some(at) = calls.iter().position(|&held| self.same_expr(held, call)) {
3331 return at;
3332 }
3333 let run = self.windows.last_mut().expect("a run is open");
3334 run.calls.push(call);
3335 run.calls.len() - 1
3336 }
3337
3338 fn window_parts(
3344 &mut self,
3345 ast: &Ast,
3346 written: &WindowCall<'_>,
3347 args: &[ast::ExprRef],
3348 held: ast::WindowSpec,
3349 scope: &Scope,
3350 ) -> Result<WindowParts> {
3351 let mut bound = Vec::with_capacity(args.len());
3352 for &arg in args {
3353 let expr = self.bind_expr(ast, arg, scope)?;
3354 bound.push(self.over_aggregate(expr, scope)?);
3355 }
3356 let mut inner = Vec::new();
3360 for item in ast.order_list(written.order).to_vec() {
3361 let expr = self.bind_expr(ast, item.expr, scope)?;
3362 let expr = self.over_aggregate(expr, scope)?;
3363 inner.push(self.sort_key(expr, item));
3364 }
3365 let mut partition = Vec::new();
3366 for &key in ast.expr_list(held.partition) {
3367 let expr = self.bind_expr(ast, key, scope)?;
3368 partition.push(self.over_aggregate(expr, scope)?);
3369 }
3370 let mut order = Vec::new();
3371 for item in ast.order_list(held.order).to_vec() {
3372 let expr = self.bind_expr(ast, item.expr, scope)?;
3373 let expr = self.over_aggregate(expr, scope)?;
3374 order.push(self.sort_key(expr, item));
3375 }
3376 let frame = WindowFrame {
3377 unit: match held.unit {
3378 ast::WindowUnit::Rows => WindowUnit::Rows,
3379 ast::WindowUnit::Range => WindowUnit::Range,
3380 ast::WindowUnit::Groups => WindowUnit::Groups,
3381 },
3382 start: self.window_bound(ast, held.start, scope)?,
3383 end: self.window_bound(ast, held.end, scope)?,
3384 exclude: match held.exclude {
3385 ast::WindowExclude::NoOthers => WindowExclude::NoOthers,
3386 ast::WindowExclude::CurrentRow => WindowExclude::CurrentRow,
3387 ast::WindowExclude::Group => WindowExclude::Group,
3388 ast::WindowExclude::Ties => WindowExclude::Ties,
3389 },
3390 };
3391 Ok(WindowParts { args: bound, partition, order, inner, frame })
3392 }
3393
3394 fn window_bound(
3396 &mut self,
3397 ast: &Ast,
3398 bound: ast::WindowBound,
3399 scope: &Scope,
3400 ) -> Result<WindowBound> {
3401 let offset = |binder: &mut Self, written| {
3402 let expr = binder.bind_expr(ast, written, scope)?;
3403 binder.over_aggregate(expr, scope)
3404 };
3405 Ok(match bound {
3406 ast::WindowBound::UnboundedPreceding => WindowBound::UnboundedPreceding,
3407 ast::WindowBound::CurrentRow => WindowBound::CurrentRow,
3408 ast::WindowBound::UnboundedFollowing => WindowBound::UnboundedFollowing,
3409 ast::WindowBound::Preceding(written) => WindowBound::Preceding(offset(self, written)?),
3410 ast::WindowBound::Following(written) => WindowBound::Following(offset(self, written)?),
3411 })
3412 }
3413
3414 fn group_of(&self, read: ColumnBinding) -> Option<usize> {
3420 self.aggregation.as_ref()?.groups.iter().position(
3421 |group| matches!(*self.plan.expr(*group), Expr::Column(binding) if binding == read),
3422 )
3423 }
3424
3425 fn ungrouped_correlation(&self, binding: ColumnBinding) -> Option<ColumnBinding> {
3431 let pending =
3432 self.scalar_subqueries.iter().find(|pending| pending.index == binding.table)?;
3433 pending.reads.iter().copied().find(|read| self.group_of(*read).is_none())
3434 }
3435
3436 fn is_window_output(&self, binding: ColumnBinding) -> bool {
3438 self.windows.iter().any(|run| run.index == binding.table)
3439 }
3440
3441 fn is_correlation(&self, binding: ColumnBinding) -> bool {
3447 self.correlations.last().is_some_and(|frame| frame.contains(&binding))
3448 }
3449
3450 fn name_of(&self, binding: ColumnBinding, scope: &Scope) -> String {
3456 std::iter::once(scope)
3457 .chain(self.outer_scopes.iter().rev())
3458 .flat_map(|visible| visible.columns.iter())
3459 .find(|column| column.binding == binding)
3460 .map_or_else(|| "a column".to_string(), |column| format!("\"{}\"", column.name))
3461 }
3462
3463 pub(crate) fn over_aggregate(&mut self, expr: ExprRef, scope: &Scope) -> Result<ExprRef> {
3469 let Some(aggregation) = self.aggregation.as_ref() else {
3470 return Ok(expr);
3471 };
3472 let index = aggregation.index;
3473 let groups = aggregation.groups.clone();
3474 for (at, group) in groups.iter().enumerate() {
3475 if self.same_expr(expr, *group) {
3476 let ty = self.plan.expr_type(*group).clone();
3477 return Ok(self.column(index, at, ty));
3478 }
3479 }
3480 let ty = self.plan.expr_type(expr).clone();
3481 match self.plan.expr(expr).clone() {
3482 Expr::Column(binding) if binding.table == index => Ok(expr),
3483 Expr::Column(binding) if self.is_window_output(binding) => Ok(expr),
3488 Expr::Column(binding) if self.joined_above.contains(&binding.table) => Ok(expr),
3493 Expr::Column(binding) if self.is_correlation(binding) => Ok(expr),
3499 Expr::Column(binding) => {
3508 let read = self.ungrouped_correlation(binding).unwrap_or(binding);
3509 let name = self.name_of(read, scope);
3510 Err(Error::binder(format!(
3511 "column {name} must appear in the GROUP BY clause or must be part of an aggregate function"
3512 )))
3513 }
3514 Expr::Constant(_)
3515 | Expr::Aggregate { .. }
3516 | Expr::Window { .. }
3517 | Expr::LambdaParam(_) => Ok(expr),
3518 Expr::Lambda { table, params, body } => {
3522 let body = self.over_aggregate(body, scope)?;
3523 Ok(self.plan.add_expr(Expr::Lambda { table, params, body }, ty))
3524 }
3525 Expr::Cast { input, try_cast } => {
3526 let input = self.over_aggregate(input, scope)?;
3527 Ok(self.plan.add_expr(Expr::Cast { input, try_cast }, ty))
3528 }
3529 Expr::Compare { op, left, right } => {
3530 let left = self.over_aggregate(left, scope)?;
3531 let right = self.over_aggregate(right, scope)?;
3532 Ok(self.plan.add_expr(Expr::Compare { op, left, right }, ty))
3533 }
3534 Expr::Conjunction { op, children } => {
3535 let written = self.plan.expr_list(children).to_vec();
3536 let mut rewritten = Vec::with_capacity(written.len());
3537 for child in written {
3538 rewritten.push(self.over_aggregate(child, scope)?);
3539 }
3540 let children = self.plan.add_expr_list(&rewritten);
3541 Ok(self.plan.add_expr(Expr::Conjunction { op, children }, ty))
3542 }
3543 Expr::Function { name, args } => {
3544 let written = self.plan.expr_list(args).to_vec();
3545 let mut rewritten = Vec::with_capacity(written.len());
3546 for arg in written {
3547 rewritten.push(self.over_aggregate(arg, scope)?);
3548 }
3549 let args = self.plan.add_expr_list(&rewritten);
3550 Ok(self.plan.add_expr(Expr::Function { name, args }, ty))
3551 }
3552 Expr::Case { arms, otherwise } => {
3553 let written = self.plan.arm_list(arms).to_vec();
3554 let mut rewritten = Vec::with_capacity(written.len());
3555 for arm in written {
3556 let when = self.over_aggregate(arm.when, scope)?;
3557 let then = self.over_aggregate(arm.then, scope)?;
3558 rewritten.push(rudb_plan::Arm { when, then });
3559 }
3560 let otherwise = match otherwise {
3561 Some(expr) => Some(self.over_aggregate(expr, scope)?),
3562 None => None,
3563 };
3564 let arms = self.plan.add_arms(&rewritten);
3565 Ok(self.plan.add_expr(Expr::Case { arms, otherwise }, ty))
3566 }
3567 }
3568 }
3569
3570 pub(crate) fn same_expr(&self, left: ExprRef, right: ExprRef) -> bool {
3572 same_expr(&self.plan, left, right)
3573 }
3574}
3575
3576#[derive(Debug, Default)]
3586struct Options {
3587 binary_as_string: bool,
3590 all_varchar: bool,
3592 file_row_number: bool,
3597 given: Given,
3599}
3600
3601impl Options {
3602 fn of(written: &[(&'static str, Value, ExprRef)]) -> Result<Self> {
3609 let mut options = Self::default();
3610 for (parameter, value, _) in written {
3611 match (*parameter, value) {
3612 ("binary_as_string", Value::Boolean(on)) => options.binary_as_string = *on,
3613 ("all_varchar", Value::Boolean(on)) => options.all_varchar = *on,
3614 ("file_row_number", Value::Boolean(on)) => options.file_row_number = *on,
3615 _ => {}
3616 }
3617 }
3618 let named: Vec<(&str, Value)> =
3619 written.iter().map(|(parameter, value, _)| (*parameter, value.clone())).collect();
3620 options.given = csv_given(&named)?;
3621 Ok(options)
3622 }
3623}
3624
3625fn mirror_target(paths: &[String]) -> Option<(String, FileStamp)> {
3628 let [path] = paths else { return None };
3629 let canonical = std::fs::canonicalize(path).ok()?;
3630 let stamp = FileStamp::of(&canonical)?;
3631 Some((canonical.to_str()?.to_string(), stamp))
3632}
3633
3634#[derive(Clone, Copy)]
3636struct Operator {
3637 op: SetOp,
3639 quantifier: Quantifier,
3641 by_name: bool,
3643}
3644
3645struct Merged {
3647 name: String,
3649 ty: LogicalType,
3651 left: Option<usize>,
3653 right: Option<usize>,
3655}
3656
3657fn match_by_position(left: &Scope, right: &Scope) -> Result<Vec<Merged>> {
3661 if left.len() != right.len() {
3662 return Err(Error::binder(format!(
3663 "Set operations can only apply to expressions with the same number of result columns, but left side has {} and right side has {}",
3664 left.len(),
3665 right.len()
3666 )));
3667 }
3668 let mut merged = Vec::with_capacity(left.len());
3669 for (at, (held, other)) in left.columns.iter().zip(&right.columns).enumerate() {
3670 merged.push(Merged {
3671 name: held.name.clone(),
3672 ty: meet(&held.ty, &other.ty)?,
3673 left: Some(at),
3674 right: Some(at),
3675 });
3676 }
3677 Ok(merged)
3678}
3679
3680fn match_by_name(left: &Scope, right: &Scope) -> Result<Vec<Merged>> {
3688 named_once(left)?;
3689 named_once(right)?;
3690 let mut merged = Vec::with_capacity(left.len() + right.len());
3691 for (at, held) in left.columns.iter().enumerate() {
3692 let other = right.columns.iter().position(|column| same_name(&column.name, &held.name));
3693 let ty = match other {
3694 Some(other) => meet(&held.ty, &right.columns[other].ty)?,
3695 None => held.ty.clone(),
3696 };
3697 merged.push(Merged { name: held.name.clone(), ty, left: Some(at), right: other });
3698 }
3699 for (at, held) in right.columns.iter().enumerate() {
3700 if left.columns.iter().any(|column| same_name(&column.name, &held.name)) {
3701 continue;
3702 }
3703 merged.push(Merged {
3704 name: held.name.clone(),
3705 ty: held.ty.clone(),
3706 left: None,
3707 right: Some(at),
3708 });
3709 }
3710 Ok(merged)
3711}
3712
3713fn named_once(scope: &Scope) -> Result<()> {
3719 for (at, held) in scope.columns.iter().enumerate() {
3720 if scope.columns[..at].iter().any(|column| same_name(&column.name, &held.name)) {
3721 return Err(Error::binder(format!(
3722 "UNION (ALL) BY NAME operation doesn't support duplicate names in the SELECT list - the name \"\"{}\"\" occurs multiple times",
3723 held.name
3724 )));
3725 }
3726 }
3727 Ok(())
3728}
3729
3730fn meet(left: &LogicalType, right: &LogicalType) -> Result<LogicalType> {
3732 left.promote(right).ok_or_else(|| {
3733 Error::binder(format!(
3734 "Cannot combine a column of type {left} with a column of type {right} in a set operation"
3735 ))
3736 })
3737}
3738
3739fn null_parameter(function: TableFunction, parameter: &str) -> String {
3748 match parameter {
3749 "header" => format!("\"{parameter}\" expects a non-null boolean value (e.g. TRUE or 1)"),
3750 "all_varchar" => format!("{} \"{parameter}\" cannot be NULL", function.name()),
3751 _ => format!("Cannot use NULL as argument to \"{parameter}\""),
3752 }
3753}
3754
3755fn missing_replacement(name: &str, input: &Scope) -> Error {
3760 Error::binder(format!(
3761 "Column \"{name}\" in REPLACE list not found in FROM clause{}",
3762 input.candidates()
3763 ))
3764}
3765
3766fn subtractable(ty: &LogicalType, ordering: bool) -> bool {
3775 if ty.is_numeric() {
3776 return true;
3777 }
3778 match ty {
3779 LogicalType::Date
3780 | LogicalType::Time
3781 | LogicalType::Timestamp
3782 | LogicalType::TimestampS
3783 | LogicalType::TimestampMs
3784 | LogicalType::TimestampNs
3785 | LogicalType::TimestampTz => true,
3786 LogicalType::TimeTz => ordering,
3787 _ => false,
3788 }
3789}
3790
3791fn refuse_fill(
3800 argument: &LogicalType,
3801 order: &[LogicalType],
3802 distinct: bool,
3803 ignore_nulls: bool,
3804) -> Result<()> {
3805 if !subtractable(argument, false) {
3806 return Err(Error::binder("FILL argument must support subtraction"));
3807 }
3808 let [key] = order else {
3809 return Err(Error::binder("FILL functions must have only one ORDER BY expression"));
3810 };
3811 if !subtractable(key, true) {
3812 return Err(Error::binder("FILL ordering must support subtraction"));
3813 }
3814 if distinct {
3815 return Err(Error::binder(
3816 "DISTINCT is not implemented for the window function \"\"fill\"\"",
3817 ));
3818 }
3819 if ignore_nulls {
3820 return Err(Error::binder(
3821 "RESPECT/IGNORE NULLS is not supported for the window function \"fill\"",
3822 ));
3823 }
3824 Ok(())
3825}
3826
3827fn window_signature(name: &str, types: &[LogicalType]) -> Result<Resolved> {
3834 match kind_of(name) {
3835 Some(FunctionKind::Aggregate | FunctionKind::Window) => resolve(name, types),
3836 Some(FunctionKind::Scalar) => {
3837 Err(Error::catalog(format!("{name} is not an aggregate function")))
3838 }
3839 None => Err(Error::catalog(format!("Aggregate Function with name {name} does not exist!"))),
3840 }
3841}
3842
3843fn same_expr(plan: &Plan, left: ExprRef, right: ExprRef) -> bool {
3845 if left == right {
3846 return true;
3847 }
3848 if plan.expr_type(left) != plan.expr_type(right) {
3849 return false;
3850 }
3851 let lists = |left, right| {
3852 let left: &[ExprRef] = plan.expr_list(left);
3853 let right: &[ExprRef] = plan.expr_list(right);
3854 left.len() == right.len()
3855 && left.iter().zip(right).all(|(&left, &right)| same_expr(plan, left, right))
3856 };
3857 match (plan.expr(left), plan.expr(right)) {
3858 (Expr::Column(left), Expr::Column(right)) => left == right,
3859 (Expr::Constant(left), Expr::Constant(right)) => plan.value(*left) == plan.value(*right),
3860 (
3861 Expr::Cast { input: left, try_cast: left_try },
3862 Expr::Cast { input: right, try_cast: right_try },
3863 ) => left_try == right_try && same_expr(plan, *left, *right),
3864 (
3865 Expr::Compare { op: left_op, left: left_a, right: left_b },
3866 Expr::Compare { op: right_op, left: right_a, right: right_b },
3867 ) => {
3868 left_op == right_op
3869 && same_expr(plan, *left_a, *right_a)
3870 && same_expr(plan, *left_b, *right_b)
3871 }
3872 (
3873 Expr::Conjunction { op: left_op, children: left_children },
3874 Expr::Conjunction { op: right_op, children: right_children },
3875 ) => left_op == right_op && lists(*left_children, *right_children),
3876 (
3877 Expr::Function { name: left_name, args: left_args },
3878 Expr::Function { name: right_name, args: right_args },
3879 ) => plan.string(*left_name) == plan.string(*right_name) && lists(*left_args, *right_args),
3880 (
3881 Expr::Aggregate {
3882 name: left_name,
3883 args: left_args,
3884 distinct: left_distinct,
3885 filter: left_filter,
3886 },
3887 Expr::Aggregate {
3888 name: right_name,
3889 args: right_args,
3890 distinct: right_distinct,
3891 filter: right_filter,
3892 },
3893 ) => {
3894 plan.string(*left_name) == plan.string(*right_name)
3895 && left_distinct == right_distinct
3896 && match (left_filter, right_filter) {
3897 (None, None) => true,
3898 (Some(left), Some(right)) => same_expr(plan, *left, *right),
3899 _ => false,
3900 }
3901 && lists(*left_args, *right_args)
3902 }
3903 (
3907 Expr::Window {
3908 name: left_name,
3909 args: left_args,
3910 distinct: left_distinct,
3911 filter: left_filter,
3912 ignore_nulls: left_nulls,
3913 order: left_order,
3914 },
3915 Expr::Window {
3916 name: right_name,
3917 args: right_args,
3918 distinct: right_distinct,
3919 filter: right_filter,
3920 ignore_nulls: right_nulls,
3921 order: right_order,
3922 },
3923 ) => {
3924 let left_keys = plan.sort_key_list(*left_order);
3927 let right_keys = plan.sort_key_list(*right_order);
3928 plan.string(*left_name) == plan.string(*right_name)
3929 && left_distinct == right_distinct
3930 && left_nulls == right_nulls
3931 && left_keys.len() == right_keys.len()
3932 && left_keys.iter().zip(right_keys).all(|(left, right)| {
3933 left.descending == right.descending
3934 && left.nulls_first == right.nulls_first
3935 && same_expr(plan, left.expr, right.expr)
3936 })
3937 && match (left_filter, right_filter) {
3938 (None, None) => true,
3939 (Some(left), Some(right)) => same_expr(plan, *left, *right),
3940 _ => false,
3941 }
3942 && lists(*left_args, *right_args)
3943 }
3944 (
3945 Expr::Case { arms: left_arms, otherwise: left_otherwise },
3946 Expr::Case { arms: right_arms, otherwise: right_otherwise },
3947 ) => {
3948 let left_arms = plan.arm_list(*left_arms);
3949 let right_arms = plan.arm_list(*right_arms);
3950 left_arms.len() == right_arms.len()
3951 && left_arms.iter().zip(right_arms).all(|(left, right)| {
3952 same_expr(plan, left.when, right.when) && same_expr(plan, left.then, right.then)
3953 })
3954 && match (left_otherwise, right_otherwise) {
3955 (None, None) => true,
3956 (Some(left), Some(right)) => same_expr(plan, *left, *right),
3957 _ => false,
3958 }
3959 }
3960 _ => false,
3961 }
3962}