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, TYPES_SET, TableFunction,
24 csv_fields, csv_given, files, is_file, is_pattern, kind_of, parquet_footers, parquet_outline,
25 resolve, 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
131#[derive(Clone, Copy)]
133pub(crate) struct AggregateCall<'a> {
134 name: &'a str,
136 args: &'a [ast::ExprRef],
138 distinct: bool,
140 filter: ast::ExprRef,
142 sorted: &'a [ast::OrderItem],
144}
145
146pub(crate) struct WindowCall<'a> {
151 pub(crate) name: &'a str,
153 pub(crate) args: &'a [ast::ExprRef],
155 pub(crate) distinct: bool,
157 pub(crate) filter: ast::ExprRef,
159 pub(crate) ignore_nulls: bool,
161 pub(crate) order: ast::Slice,
164 pub(crate) spec: ast::WindowRef,
166}
167
168struct WindowParts {
170 args: Vec<ExprRef>,
172 partition: Vec<ExprRef>,
174 order: Vec<SortKey>,
176 inner: Vec<SortKey>,
179 frame: WindowFrame,
181}
182
183#[derive(Debug)]
190struct Read {
191 fields: Vec<Field>,
193 rows: Stat<u64>,
195 distincts: Vec<(String, Stat<u64>)>,
197 zones: Option<Arc<dyn Zones>>,
199}
200
201impl Read {
202 fn uncounted(fields: Vec<Field>) -> Self {
204 Self { fields, rows: Stat::Unknown, distincts: Vec::new(), zones: None }
205 }
206}
207
208#[derive(Debug)]
210struct Materialized {
211 written: u32,
213 cte: u32,
215 name: String,
217 fields: Vec<Field>,
219}
220
221#[derive(Debug)]
222pub(crate) struct PendingSubquery {
223 pub(crate) node: NodeRef,
224 pub(crate) kind: JoinKind,
225 pub(crate) conditions: Vec<ExprRef>,
226 pub(crate) dependent: bool,
227 pub(crate) reads: Vec<ColumnBinding>,
233 pub(crate) index: u32,
239 pub(crate) inside_aggregate: bool,
246}
247
248#[derive(Debug, Clone, Copy, PartialEq, Eq)]
250enum Side {
251 Left,
252 Right,
253}
254
255#[derive(Debug)]
257pub(crate) struct Binder<'a> {
258 catalog: &'a Catalog,
259 pub(crate) parameters: &'a Parameters,
261 pub(crate) session: &'a Session,
263 pub(crate) semantics: Semantics,
265 plan: Plan,
266 next_index: u32,
267 pub(crate) current_span: Span,
269 pub(crate) pinned_span: Option<Span>,
272 pub(crate) aggregation: Option<Aggregation>,
274 want_ascending: bool,
277 pub(crate) upsert: bool,
280 pub(crate) insert_defaults: Option<Vec<(LogicalType, Option<String>)>>,
283 pub(crate) copy_into: Option<Vec<Field>>,
287 pub(crate) default_as_null: bool,
290 pub(crate) in_aggregate: bool,
292 pub(crate) in_filter: bool,
294 pub(crate) windows: Vec<WindowRun>,
296 pub(crate) unnests: Vec<crate::unnest::UnnestCall>,
298 pub(crate) unnest_index: Option<u32>,
300 pub(crate) unnest_here: bool,
303 pub(crate) in_unnest: bool,
305 pub(crate) unnest_root: bool,
308 pub(crate) unnest_struct: Option<crate::unnest::UnnestStruct>,
310 pub(crate) unnest_grouping: Option<bool>,
313 pub(crate) grouped_unnests: Vec<crate::unnest::GroupedUnnest>,
316 pub(crate) sequences: Vec<QualifiedName>,
318 pub(crate) in_window: bool,
320 pub(crate) scalar_subqueries: Vec<PendingSubquery>,
322 pub(crate) joined_above: Vec<u32>,
328 pub(crate) outer_scopes: Vec<Scope>,
329 pub(crate) lateral_scopes: Vec<usize>,
336 pub(crate) correlations: Vec<Vec<ColumnBinding>>,
337 pub(crate) lambda_frames: Vec<crate::lambda::Frame>,
339 pub(crate) trying: bool,
341 pub(crate) clause: &'static str,
343 pub(crate) outlined: bool,
352 expanding: Vec<String>,
354 materialized: Vec<Materialized>,
360 next_cte: u32,
362 started: Option<i64>,
364}
365
366impl<'a> Binder<'a> {
367 pub(crate) fn with(
368 catalog: &'a Catalog,
369 parameters: &'a Parameters,
370 session: &'a Session,
371 ) -> Self {
372 Self {
373 catalog,
374 parameters,
375 session,
376 semantics: session.semantics(),
377 plan: Plan::new(),
378 next_index: 0,
379 current_span: Span::new(0, 0),
380 pinned_span: None,
381 aggregation: None,
382 want_ascending: false,
383 upsert: false,
384 insert_defaults: None,
385 copy_into: None,
386 default_as_null: false,
387 in_aggregate: false,
388 in_filter: false,
389 windows: Vec::new(),
390 unnests: Vec::new(),
391 unnest_index: None,
392 unnest_here: false,
393 in_unnest: false,
394 unnest_root: false,
395 unnest_struct: None,
396 unnest_grouping: None,
397 grouped_unnests: Vec::new(),
398 sequences: Vec::new(),
399 in_window: false,
400 scalar_subqueries: Vec::new(),
401 joined_above: Vec::new(),
402 outer_scopes: Vec::new(),
403 lateral_scopes: Vec::new(),
404 correlations: Vec::new(),
405 lambda_frames: Vec::new(),
406 trying: false,
407 clause: "SELECT clause",
408 outlined: false,
409 expanding: Vec::new(),
410 materialized: Vec::new(),
411 next_cte: 0,
412 started: None,
413 }
414 }
415
416 pub(crate) fn catalog(&self) -> &Catalog {
417 self.catalog
418 }
419
420 pub(crate) fn instant(&mut self) -> i64 {
427 *self.started.get_or_insert_with(crate::context::micros_now)
428 }
429
430 pub(crate) fn plan(&self) -> &Plan {
431 &self.plan
432 }
433
434 pub(crate) fn plan_mut(&mut self) -> &mut Plan {
435 &mut self.plan
436 }
437
438 pub(crate) fn add_expr(&mut self, expr: Expr, ty: LogicalType) -> ExprRef {
439 self.plan.add_expr_at(expr, ty, self.current_span)
440 }
441
442 pub(crate) fn add_constant(&mut self, value: Value) -> ExprRef {
443 let ty = value.logical_type();
444 let reference = self.plan.add_value(value);
445 self.plan.add_expr_at(Expr::Constant(reference), ty, self.current_span)
446 }
447
448 pub(crate) fn add_node(&mut self, node: Node) -> NodeRef {
449 self.plan.add_node_at(node, self.current_span)
450 }
451
452 pub(crate) fn into_plan(self) -> Plan {
453 self.plan
454 }
455
456 pub(crate) fn fresh_index(&mut self) -> u32 {
458 let index = self.next_index;
459 self.next_index += 1;
460 index
461 }
462
463 fn column(&mut self, index: u32, position: usize, ty: LogicalType) -> ExprRef {
465 let binding = ColumnBinding::new(index, position as u32);
466 self.plan.add_expr(Expr::Column(binding), ty)
467 }
468
469 fn attach_scalar_subqueries(&mut self, mut input: NodeRef) -> NodeRef {
471 let subqueries = std::mem::take(&mut self.scalar_subqueries);
472 for pending in subqueries {
473 input = self.attach_subquery(input, pending);
474 }
475 input
476 }
477
478 fn attach_subquery(&mut self, input: NodeRef, pending: PendingSubquery) -> NodeRef {
485 let PendingSubquery {
486 node: mut right,
487 kind,
488 conditions,
489 dependent,
490 reads: _,
491 index: _,
492 inside_aggregate: _,
493 } = pending;
494 if kind == JoinKind::Single && !self.semantics.scalar_subquery_error_on_multiple_rows() {
495 right = self.add_node(Node::Limit {
496 input: right,
497 count: Bound::Rows(1),
498 offset: Bound::Rows(0),
499 });
500 }
501 let conditions = self.plan.add_expr_list(&conditions);
502 if dependent {
503 self.add_node(Node::DependentJoin { left: input, right, kind, conditions })
504 } else {
505 self.add_node(Node::Join {
506 left: input,
507 right,
508 kind,
509 conditions,
510 build: BuildSide::default(),
511 })
512 }
513 }
514
515 pub(crate) fn bind_query(
518 &mut self,
519 ast: &Ast,
520 query: ast::QueryRef,
521 ) -> Result<(NodeRef, Scope)> {
522 let span = ast.query_span(query);
523 let outer = std::mem::replace(&mut self.current_span, span);
524 let result =
525 self.bind_query_inner(ast, query).map_err(|error| error.with_fallback_span(span));
526 self.current_span = outer;
527 result
528 }
529
530 fn bind_query_inner(&mut self, ast: &Ast, query: ast::QueryRef) -> Result<(NodeRef, Scope)> {
531 let written = ast.query(query);
532 if written.ctes.is_empty() {
533 return self.bind_body(ast, &written);
534 }
535 let depth = self.materialized.len();
539 let result = self.bind_materialized(ast, &written);
540 self.materialized.truncate(depth);
541 result
542 }
543
544 fn bind_materialized(&mut self, ast: &Ast, written: &ast::Query) -> Result<(NodeRef, Scope)> {
550 let depth = self.materialized.len();
551 let held = ast.cte_list(written.ctes).to_vec();
552 let mut definitions = Vec::with_capacity(held.len());
553 for &index in &held {
554 definitions.push(self.bind_definition(ast, index)?);
555 }
556 let (mut node, scope) = self.bind_body(ast, written)?;
557 for (at, definition) in definitions.into_iter().enumerate().rev() {
558 let entry = &self.materialized[depth + at];
559 let cte = entry.cte;
560 let name = entry.name.clone();
561 let fields = entry.fields.clone();
562 let name = self.plan.intern(&name);
563 let columns = self.plan.add_fields(&fields);
564 node =
565 self.add_node(Node::MaterializedCte { definition, body: node, name, cte, columns });
566 }
567 Ok((node, scope))
568 }
569
570 fn bind_definition(&mut self, ast: &Ast, index: u32) -> Result<NodeRef> {
580 let held = ast.cte(index);
581 let name = ast.string(held.name).to_string();
582 let (node, mut scope) = self.bind_query(ast, held.query)?;
583 if !held.columns.is_empty() {
584 let names: Vec<&str> = ast.name(held.columns).collect();
585 scope.rename_prefix(&names);
586 }
587 let table = self.fresh_index();
588 let mut exprs = Vec::with_capacity(scope.len());
589 let mut names = Vec::with_capacity(scope.len());
590 for column in &scope.columns {
591 exprs.push(self.plan.add_expr(Expr::Column(column.binding), column.ty.clone()));
592 names.push(self.plan.intern(&column.name));
593 }
594 let exprs = self.plan.add_expr_list(&exprs);
595 let names = self.plan.add_name_list(&names);
596 let node = self.add_node(Node::Project { input: node, index: table, exprs, names });
597 let cte = self.next_cte;
598 self.next_cte += 1;
599 self.materialized.push(Materialized { written: index, cte, name, fields: scope.fields() });
600 Ok(node)
601 }
602
603 fn bind_body(&mut self, ast: &Ast, written: &ast::Query) -> Result<(NodeRef, Scope)> {
604 match written.body {
605 ast::QueryBody::Select(select) => self.bind_select(ast, select, written),
606 ast::QueryBody::SetOp { op, quantifier, by_name, left, right } => {
607 let operator = Operator { op, quantifier, by_name };
608 self.bind_set_op(ast, written, operator, left, right)
609 }
610 ast::QueryBody::Values(rows) => self.bind_values(ast, written, rows),
611 ast::QueryBody::Describe(inner) => self.bind_describe(ast, written, inner),
612 ast::QueryBody::Show { name, relation } => self.bind_show(ast, written, name, relation),
613 }
614 }
615
616 fn bind_show(
618 &mut self,
619 ast: &Ast,
620 query: &ast::Query,
621 name: ast::Slice,
622 relation: ast::QueryRef,
623 ) -> Result<(NodeRef, Scope)> {
624 let text = ast.name_text(name);
625 let parts: Vec<&str> = ast.name(name).collect();
626 let table_exists = self.catalog.resolve(&parts).is_ok();
627 let as_table = match self.semantics.show_behavior() {
628 ShowBehavior::Auto => table_exists,
629 ShowBehavior::Setting => false,
630 ShowBehavior::Table => true,
631 };
632 if as_table {
633 return self.bind_describe(ast, query, relation);
634 }
635 let shown = match self.session.iter().find(|(name, _)| name.eq_ignore_ascii_case(&text)) {
640 Some((_, value)) => value.to_string(),
641 None => match self.beyond(&text)? {
642 Some(Value::Varchar(declared)) => declared,
643 Some(other) => other.to_string(),
644 None => {
645 return Err(Error::catalog(format!(
646 "Setting with name \"{text}\" does not exist"
647 )));
648 }
649 },
650 };
651 let field = Field::new(text, LogicalType::Varchar);
652 let expr = self.plan.add_constant(Value::Varchar(shown));
653 let row = self.plan.add_expr_list(&[expr]);
654 let rows = self.plan.add_rows(&[row]);
655 let columns = self.plan.add_fields(std::slice::from_ref(&field));
656 let index = self.fresh_index();
657 let node = self.add_node(Node::Values { index, columns, rows });
658 let mut scope = Scope::empty();
659 scope.push(Visible {
660 table: String::new(),
661 name: field.name,
662 binding: ColumnBinding::new(index, 0),
663 ty: LogicalType::Varchar,
664 not_null: false,
665 key: None,
666 default: None,
667 qualified: false,
668 also: None,
669 });
670 Ok((node, scope))
671 }
672
673 fn bind_describe(
688 &mut self,
689 ast: &Ast,
690 query: &ast::Query,
691 inner: ast::QueryRef,
692 ) -> Result<(NodeRef, Scope)> {
693 let (_, described) = self.bind_query(ast, inner)?;
694 let fields: Vec<Field> = ["column_name", "column_type", "null", "key", "default", "extra"]
695 .iter()
696 .map(|name| Field::new(*name, LogicalType::Varchar))
697 .collect();
698 let mut slices = Vec::with_capacity(described.columns.len());
699 for column in described.columns.clone() {
700 let written = [
703 column.name.clone(),
704 column.ty.to_string(),
705 if column.not_null { "NO" } else { "YES" }.to_owned(),
706 ];
707 let mut items: Vec<ExprRef> = written
708 .into_iter()
709 .map(|text| self.plan.add_constant(Value::Varchar(text)))
710 .collect();
711 let mark = match column.key {
712 Some(mark) => self.plan.add_constant(Value::Varchar(mark.to_owned())),
713 None => {
714 let empty = self.plan.add_constant(Value::Null);
715 self.cast_to(empty, &LogicalType::Varchar)
716 }
717 };
718 items.push(mark);
719 if let Some(default) = &column.default {
720 items.push(self.plan.add_constant(Value::Varchar(default.clone())));
721 }
722 while items.len() < 6 {
723 let empty = self.plan.add_constant(Value::Null);
724 items.push(self.cast_to(empty, &LogicalType::Varchar));
725 }
726 slices.push(self.plan.add_expr_list(&items));
727 }
728 let rows = self.plan.add_rows(&slices);
729 let columns = self.plan.add_fields(&fields);
730 let index = self.fresh_index();
731 let mut node = self.add_node(Node::Values { index, columns, rows });
732 let mut scope = Scope::empty();
733 for (at, field) in fields.iter().enumerate() {
734 scope.push(Visible {
735 table: String::new(),
736 name: field.name.clone(),
737 binding: ColumnBinding::new(index, at as u32),
738 ty: field.ty.clone(),
739 not_null: false,
740 key: None,
741 default: None,
742 qualified: false,
743 also: None,
744 });
745 }
746 let keys = self.sort_keys(ast, query, &scope, &[])?;
747 if !keys.is_empty() {
748 let keys = self.plan.add_sort_keys(&keys);
749 node = self.add_node(Node::Sort { input: node, keys });
750 }
751 node = self.apply_limit(ast, query, node, &mut scope)?;
752 Ok((node, scope))
753 }
754
755 pub(crate) fn bind_default(&mut self, text: Option<&str>, ty: &LogicalType) -> Result<ExprRef> {
759 let Some(text) = text else {
760 let value = self.plan.add_value(Value::Null);
761 return Ok(self.plan.add_expr(Expr::Constant(value), ty.clone()));
762 };
763 let ast = rudb_parse::parse_ast(&format!("SELECT {text}"))?;
764 let found = match ast.statements.first() {
765 Some(&ast::Statement::Query(query)) => match ast.query(query).body {
766 ast::QueryBody::Select(select) => {
767 ast.target_list(ast.select(select).targets).first().map(|target| target.expr)
768 }
769 _ => None,
770 },
771 _ => None,
772 };
773 let Some(expr) = found else {
774 return Err(Error::internal(format!("a default that is not an expression: {text}")));
775 };
776 let expr = self.bind_expr(&ast, expr, &Scope::empty())?;
777 self.checked_cast_to(expr, ty, false)
778 }
779
780 fn passes_through(&self, expr: ExprRef, input: &Scope) -> bool {
786 let Expr::Column(binding) = *self.plan.expr(expr) else { return false };
787 input.columns.iter().any(|column| column.binding == binding && column.not_null)
788 }
789
790 fn key_through(&self, expr: ExprRef, input: &Scope) -> Option<&'static str> {
794 self.through(expr, input).and_then(|column| column.key)
795 }
796
797 fn through<'s>(&self, expr: ExprRef, input: &'s Scope) -> Option<&'s Visible> {
799 let Expr::Column(binding) = *self.plan.expr(expr) else { return None };
800 input.columns.iter().find(|column| column.binding == binding)
801 }
802
803 fn bind_values(
810 &mut self,
811 ast: &Ast,
812 query: &ast::Query,
813 rows: ast::Slice,
814 ) -> Result<(NodeRef, Scope)> {
815 let written = ast.rows(rows).to_vec();
816 let Some(first) = written.first() else {
817 return Err(Error::binder("VALUES needs at least one row"));
818 };
819 let width = first.len as usize;
820 for (at, row) in written.iter().enumerate() {
821 if row.len as usize != width {
822 return Err(Error::binder(format!(
823 "VALUES lists must all be the same length, expected {width} columns but row {} has {}",
824 at + 1,
825 row.len
826 )));
827 }
828 }
829 let empty = Scope::empty();
831 let defaults = self.insert_defaults.take();
832 let previous = std::mem::replace(&mut self.clause, "VALUES clause");
833 let mut bound: Vec<Vec<ExprRef>> = Vec::with_capacity(written.len());
834 for row in &written {
835 let mut items = Vec::with_capacity(width);
836 for (at, &expr) in ast.expr_list(*row).iter().enumerate() {
837 let column = defaults.as_ref().and_then(|defaults| defaults.get(at));
838 items.push(match (ast.expr(expr), column) {
839 (ast::Expr::Default, Some((ty, default))) => {
840 self.bind_default(default.as_deref(), ty)?
841 }
842 _ => self.bind_expr(ast, expr, &empty)?,
843 });
844 }
845 bound.push(items);
846 }
847 self.clause = previous;
848 let mut types = Vec::with_capacity(width);
849 for at in 0..width {
850 let mut ty = self.plan.expr_type(bound[0][at]).clone();
851 for row in &bound[1..] {
852 let other = self.plan.expr_type(row[at]).clone();
853 ty = ty.promote(&other).ok_or_else(|| {
854 Error::binder(format!(
855 "Cannot combine a value of type {ty} with a value of type {other} in column {} of a VALUES",
856 at + 1
857 ))
858 })?;
859 }
860 types.push(ty);
861 }
862 let mut slices = Vec::with_capacity(bound.len());
863 for row in &bound {
864 let items: Vec<ExprRef> = row
865 .iter()
866 .zip(&types)
867 .map(|(&expr, ty)| self.checked_cast_to(expr, ty, false))
868 .collect::<Result<_>>()?;
869 slices.push(self.plan.add_expr_list(&items));
870 }
871 let rows = self.plan.add_rows(&slices);
872 let fields: Vec<Field> = types
873 .iter()
874 .enumerate()
875 .map(|(at, ty)| Field::new(format!("col{at}"), ty.clone()))
876 .collect();
877 let columns = self.plan.add_fields(&fields);
878 let index = self.fresh_index();
879 let mut node = self.add_node(Node::Values { index, columns, rows });
880 let mut scope = Scope::empty();
881 for (at, field) in fields.iter().enumerate() {
882 scope.push(Visible {
883 table: String::new(),
884 name: field.name.clone(),
885 binding: ColumnBinding::new(index, at as u32),
886 ty: field.ty.clone(),
887 not_null: false,
888 key: None,
889 default: None,
890 qualified: false,
891 also: None,
892 });
893 }
894 let keys = self.sort_keys(ast, query, &scope, &[])?;
895 if !keys.is_empty() {
896 let keys = self.plan.add_sort_keys(&keys);
897 node = self.add_node(Node::Sort { input: node, keys });
898 }
899 node = self.apply_limit(ast, query, node, &mut scope)?;
900 Ok((node, scope))
901 }
902
903 fn bind_set_op(
904 &mut self,
905 ast: &Ast,
906 query: &ast::Query,
907 operator: Operator,
908 left: ast::QueryRef,
909 right: ast::QueryRef,
910 ) -> Result<(NodeRef, Scope)> {
911 let (left_node, left_scope) = self.bind_query(ast, left)?;
912 let (right_node, right_scope) = self.bind_query(ast, right)?;
913 let merged = if operator.by_name {
914 match_by_name(&left_scope, &right_scope)?
915 } else {
916 match_by_position(&left_scope, &right_scope)?
917 };
918 let left_node = self.conform(left_node, &left_scope, &merged, |column| column.left)?;
919 let right_node = self.conform(right_node, &right_scope, &merged, |column| column.right)?;
920 let index = self.fresh_index();
921 let kind = match operator.op {
922 SetOp::Union => SetOpKind::Union,
923 SetOp::Except => SetOpKind::Except,
924 SetOp::Intersect => SetOpKind::Intersect,
925 };
926 let all = operator.quantifier == Quantifier::All;
929 let mut node =
930 self.add_node(Node::SetOp { left: left_node, right: right_node, kind, all, index });
931 let mut scope = Scope::empty();
932 for (at, column) in merged.iter().enumerate() {
933 scope.push(Visible {
934 table: String::new(),
935 name: column.name.clone(),
936 binding: ColumnBinding::new(index, at as u32),
937 ty: column.ty.clone(),
938 not_null: false,
941 key: None,
942 default: None,
943 qualified: false,
944 also: None,
945 });
946 }
947 let keys = self.sort_keys(ast, query, &scope, &[])?;
951 if !keys.is_empty() {
952 let keys = self.plan.add_sort_keys(&keys);
953 node = self.add_node(Node::Sort { input: node, keys });
954 }
955 node = self.apply_limit(ast, query, node, &mut scope)?;
956 Ok((node, scope))
957 }
958
959 fn conform(
965 &mut self,
966 node: NodeRef,
967 scope: &Scope,
968 merged: &[Merged],
969 pick: impl Fn(&Merged) -> Option<usize>,
970 ) -> Result<NodeRef> {
971 let unchanged = merged.len() == scope.len()
972 && merged
973 .iter()
974 .enumerate()
975 .all(|(at, column)| pick(column) == Some(at) && column.ty == scope.columns[at].ty);
976 if unchanged {
977 return Ok(node);
978 }
979 let index = self.fresh_index();
980 let mut exprs = Vec::with_capacity(merged.len());
981 let mut names = Vec::with_capacity(merged.len());
982 for column in merged {
983 let expr = match pick(column) {
984 Some(at) => {
985 let held = &scope.columns[at];
986 self.plan.add_expr(Expr::Column(held.binding), held.ty.clone())
987 }
988 None => self.plan.add_constant(Value::Null),
989 };
990 exprs.push(self.checked_cast_to(expr, &column.ty, false)?);
991 names.push(self.plan.intern(&column.name));
992 }
993 let exprs = self.plan.add_expr_list(&exprs);
994 let names = self.plan.add_name_list(&names);
995 Ok(self.add_node(Node::Project { input: node, index, exprs, names }))
996 }
997
998 fn bind_select(
1001 &mut self,
1002 ast: &Ast,
1003 select: ast::SelectRef,
1004 query: &ast::Query,
1005 ) -> Result<(NodeRef, Scope)> {
1006 let written = ast.select(select);
1007 self.want_ascending |= !written.group_by.is_empty() || written.group_by_all;
1008 let outer_windows = std::mem::take(&mut self.windows);
1012 let outer_unnests = std::mem::take(&mut self.unnests);
1014 let outer_unnest_index = self.unnest_index.take();
1015 let outer_unnest_here = std::mem::replace(&mut self.unnest_here, false);
1016 let outer_in_unnest = std::mem::replace(&mut self.in_unnest, false);
1017 let outer_unnest_grouping = self.unnest_grouping.take();
1018 let outer_grouped_unnests = std::mem::take(&mut self.grouped_unnests);
1019 let outer_joined_above = std::mem::take(&mut self.joined_above);
1024 let (mut node, input) = self.bind_from(ast, written.from)?;
1025 node = self.attach_scalar_subqueries(node);
1026
1027 if written.filter != NONE {
1028 self.clause = "WHERE clause";
1029 let predicate = self.bind_expr(ast, written.filter, &input)?;
1030 let predicate = self.as_boolean(predicate, "WHERE")?;
1031 node = self.attach_scalar_subqueries(node);
1032 node = self.add_node(Node::Filter { input: node, predicate });
1033 }
1034
1035 let targets = ast.target_list(written.targets).to_vec();
1036 if targets.is_empty() {
1037 return Err(Error::binder("a SELECT needs at least one expression to select"));
1038 }
1039
1040 let group_items = self.group_items(ast, &written, &targets)?;
1041 let aggregating = !group_items.is_empty()
1042 || written.having != NONE
1043 || targets.iter().any(|target| has_aggregate(ast, target.expr));
1044 if aggregating {
1045 self.clause = "GROUP BY clause";
1046 self.unnest_here = true;
1050 self.unnest_grouping = Some(written.group_by_all);
1051 let mut groups = Vec::with_capacity(group_items.len());
1052 for item in &group_items {
1053 groups.push(self.bind_expr(ast, *item, &input)?);
1054 }
1055 self.unnest_here = false;
1056 self.unnest_grouping = None;
1057 let unnests = std::mem::take(&mut self.unnests);
1058 if let Some(index) = self.unnest_index.take() {
1059 node = self.plan_unnests(node, index, &unnests)?;
1060 }
1061 let index = self.fresh_index();
1062 self.aggregation = Some(Aggregation { index, groups, aggregates: Vec::new() });
1063 }
1064
1065 let mut above = Vec::new();
1072
1073 self.clause = "SELECT clause";
1074 self.unnest_here = true;
1075 let (mut exprs, mut names) = self.bind_targets(ast, &targets, &input, &mut above)?;
1076 self.unnest_here = false;
1077 let visible = exprs.len();
1078
1079 let mut having = None;
1080 if written.having != NONE {
1081 self.clause = "HAVING clause";
1082 let before = self.scalar_subqueries.len();
1083 let predicate = self.bind_expr(ast, written.having, &input)?;
1084 self.lift_over_aggregate(before, &mut above, &input)?;
1085 let predicate = self.over_aggregate(predicate, &input)?;
1086 having = Some(self.as_boolean(predicate, "HAVING")?);
1087 }
1088
1089 let project = self.fresh_index();
1092 let mut output = Scope::empty();
1093 for (at, (expr, name)) in exprs.iter().zip(&names).enumerate() {
1094 output.push(Visible {
1095 table: String::new(),
1096 name: name.clone(),
1097 binding: ColumnBinding::new(project, at as u32),
1098 ty: self.plan.expr_type(*expr).clone(),
1099 not_null: self.passes_through(*expr, &input),
1100 key: self.key_through(*expr, &input),
1101 default: self.through(*expr, &input).and_then(|column| column.default.clone()),
1102 qualified: false,
1103 also: None,
1104 });
1105 }
1106
1107 self.clause = "ORDER BY clause";
1108 let mut extra = Vec::new();
1109 self.unnest_here = true;
1110 let keys = self.select_sort_keys(
1111 ast, query, &input, &output, project, &mut exprs, &mut names, &mut extra, &mut above,
1112 )?;
1113 self.unnest_here = outer_unnest_here;
1114 self.in_unnest = outer_in_unnest;
1115 self.unnest_grouping = outer_unnest_grouping;
1116 self.grouped_unnests = outer_grouped_unnests;
1117 self.joined_above = outer_joined_above;
1118 if !extra.is_empty() && written.distinct != Distinct::No {
1119 return Err(Error::binder(
1120 "For SELECT DISTINCT, ORDER BY expressions must appear in the select list",
1121 ));
1122 }
1123 let on = self.distinct_on(ast, written.distinct, &output)?;
1124
1125 node = self.attach_scalar_subqueries(node);
1126
1127 if let Some(aggregation) = self.aggregation.take() {
1128 let index = aggregation.index;
1129 let groups = self.plan.add_expr_list(&aggregation.groups);
1130 let aggregates = self.plan.add_expr_list(&aggregation.aggregates);
1131 node = self.add_node(Node::Aggregate { input: node, index, groups, aggregates });
1132 }
1133 if !above.is_empty() {
1134 debug_assert!(self.scalar_subqueries.is_empty(), "a query is waiting to be joined");
1135 self.scalar_subqueries = above;
1136 node = self.attach_scalar_subqueries(node);
1137 }
1138 if let Some(predicate) = having {
1139 node = self.add_node(Node::Filter { input: node, predicate });
1140 }
1141
1142 for run in std::mem::replace(&mut self.windows, outer_windows) {
1146 let partition = self.plan.add_expr_list(&run.partition);
1147 let order = self.plan.add_sort_keys(&run.order);
1148 let expressions = self.plan.add_expr_list(&run.calls);
1149 node = self.add_node(Node::Window {
1150 input: node,
1151 index: run.index,
1152 partition,
1153 order,
1154 frame: run.frame,
1155 expressions,
1156 });
1157 }
1158 let unnests = std::mem::replace(&mut self.unnests, outer_unnests);
1161 if let Some(index) = std::mem::replace(&mut self.unnest_index, outer_unnest_index) {
1162 node = self.plan_unnests(node, index, &unnests)?;
1163 }
1164
1165 let interned: Vec<u32> = names.iter().map(|name| self.plan.intern(name)).collect();
1166 let exprs_slice = self.plan.add_expr_list(&exprs);
1167 let names_slice = self.plan.add_name_list(&interned);
1168 node = self.add_node(Node::Project {
1169 input: node,
1170 index: project,
1171 exprs: exprs_slice,
1172 names: names_slice,
1173 });
1174
1175 if written.distinct != Distinct::No {
1176 let on = self.plan.add_expr_list(&on);
1177 node = self.add_node(Node::Distinct { input: node, on });
1178 }
1179 if !keys.is_empty() {
1180 let keys = self.plan.add_sort_keys(&keys);
1181 node = self.add_node(Node::Sort { input: node, keys });
1182 }
1183 node = self.apply_limit(ast, query, node, &mut output)?;
1184
1185 if extra.is_empty() {
1186 output.columns.truncate(visible);
1187 return Ok((node, output));
1188 }
1189 let index = self.fresh_index();
1192 let mut kept = Vec::with_capacity(visible);
1193 let mut kept_names = Vec::with_capacity(visible);
1194 let mut scope = Scope::empty();
1195 for (at, name) in names.iter().enumerate().take(visible) {
1196 let ty = output.columns[at].ty.clone();
1197 let binding = output.columns[at].binding;
1201 kept.push(self.plan.add_expr(Expr::Column(binding), ty.clone()));
1202 kept_names.push(self.plan.intern(name));
1203 scope.push(Visible {
1204 table: String::new(),
1205 name: name.clone(),
1206 binding: ColumnBinding::new(index, at as u32),
1207 ty,
1208 not_null: output.columns[at].not_null,
1209 key: output.columns[at].key,
1210 default: output.columns[at].default.clone(),
1211 qualified: false,
1212 also: None,
1213 });
1214 }
1215 let exprs = self.plan.add_expr_list(&kept);
1216 let names = self.plan.add_name_list(&kept_names);
1217 node = self.add_node(Node::Project { input: node, index, exprs, names });
1218 Ok((node, scope))
1219 }
1220
1221 fn lift_over_aggregate(
1240 &mut self,
1241 before: usize,
1242 above: &mut Vec<PendingSubquery>,
1243 scope: &Scope,
1244 ) -> Result<()> {
1245 if self.aggregation.is_none() {
1246 return Ok(());
1247 }
1248 let mut lifted = Vec::new();
1249 for mut pending in self.scalar_subqueries.split_off(before) {
1250 let stays = pending.inside_aggregate
1251 || (pending.dependent && !self.lift_correlated(&mut pending));
1252 if stays {
1253 self.scalar_subqueries.push(pending);
1254 } else {
1255 self.joined_above.push(pending.index);
1256 lifted.push(pending);
1257 }
1258 }
1259 for pending in &mut lifted {
1264 let conditions = std::mem::take(&mut pending.conditions);
1265 let mut over = Vec::with_capacity(conditions.len());
1266 for condition in conditions {
1267 over.push(self.over_aggregate(condition, scope)?);
1268 }
1269 pending.conditions = over;
1270 }
1271 above.append(&mut lifted);
1272 Ok(())
1273 }
1274
1275 fn lift_correlated(&mut self, pending: &mut PendingSubquery) -> bool {
1292 let Some(index) = self.aggregation.as_ref().map(|aggregation| aggregation.index) else {
1293 return false;
1294 };
1295 let mut moved = Vec::with_capacity(pending.reads.len());
1296 for read in &pending.reads {
1297 let Some(at) = self.group_of(*read) else {
1298 return false;
1299 };
1300 moved.push((*read, ColumnBinding::new(index, at as u32)));
1301 }
1302 let mut rewrites = Vec::new();
1303 self.plan.subtree_columns(pending.node, &mut |reference, binding| {
1304 if let Some(&(_, to)) = moved.iter().find(|(from, _)| *from == binding) {
1305 rewrites.push((reference, to));
1306 }
1307 });
1308 for (reference, to) in rewrites {
1309 self.plan.rebind(reference, to);
1310 }
1311 pending.reads = moved.into_iter().map(|(_, to)| to).collect();
1312 true
1313 }
1314
1315 fn bind_targets(
1316 &mut self,
1317 ast: &Ast,
1318 targets: &[ast::Target],
1319 input: &Scope,
1320 above: &mut Vec<PendingSubquery>,
1321 ) -> Result<(Vec<ExprRef>, Vec<String>)> {
1322 let mut exprs = Vec::with_capacity(targets.len());
1323 let mut names = Vec::with_capacity(targets.len());
1324 for target in targets {
1325 if let ast::Expr::Star { qualifier, replacements } = ast.expr(target.expr) {
1326 let table = ast.name(qualifier).last().map(str::to_string);
1327 let expanded: Vec<Visible> =
1328 input.star(table.as_deref())?.into_iter().cloned().collect();
1329 let replacements = ast.target_list(replacements).to_vec();
1330 let mut used = vec![false; replacements.len()];
1331 for column in expanded {
1332 let found = replacements.iter().zip(&mut used).find(|(replacement, _)| {
1333 same_name(ast.string(replacement.alias), &column.name)
1334 });
1335 let before = self.scalar_subqueries.len();
1340 let (expr, name) = match found {
1341 Some((replacement, used)) => {
1342 *used = true;
1343 let expr = self.bind_expr(ast, replacement.expr, input)?;
1344 (expr, ast.string(replacement.alias).to_string())
1345 }
1346 None => (
1347 self.plan.add_expr(Expr::Column(column.binding), column.ty),
1348 column.name,
1349 ),
1350 };
1351 self.lift_over_aggregate(before, above, input)?;
1352 exprs.push(self.over_aggregate(expr, input)?);
1353 names.push(name);
1354 }
1355 if let Some((replacement, _)) =
1359 replacements.iter().zip(&used).find(|(_, used)| !**used)
1360 {
1361 return Err(missing_replacement(ast.string(replacement.alias), input));
1362 }
1363 continue;
1364 }
1365 let before = self.scalar_subqueries.len();
1366 self.unnest_root = matches!(ast.expr(target.expr), ast::Expr::Function { name, .. }
1367 if name.len == 1 && same_name(ast.name(name).last().unwrap_or_default(), "unnest"));
1368 let expr = self.bind_expr(ast, target.expr, input);
1369 self.unnest_root = false;
1370 let expr = expr?;
1371 self.lift_over_aggregate(before, above, input)?;
1372 if let Some(taking) = self.unnest_struct.take() {
1373 let expr = self.over_aggregate(expr, input)?;
1375 self.unnest_fields(expr, taking, None, &mut exprs, &mut names)?;
1376 continue;
1377 }
1378 exprs.push(self.over_aggregate(expr, input)?);
1379 names.push(if target.alias == NONE {
1380 self.output_name(ast, target.expr, input)
1381 } else {
1382 ast.string(target.alias).to_string()
1383 });
1384 }
1385 Ok((exprs, names))
1386 }
1387
1388 fn output_name(&self, ast: &Ast, target: ast::ExprRef, input: &Scope) -> String {
1394 if let ast::Expr::Column { name } = ast.expr(target) {
1395 let parts: Vec<&str> = ast.name(name).collect();
1396 if let Ok(found) = input.resolve(&parts) {
1397 let written = parts.last().copied().unwrap_or_default();
1400 if let Some(also) = &found.also
1401 && !same_name(&found.name, written)
1402 && same_name(also, written)
1403 {
1404 return also.clone();
1405 }
1406 return found.name.clone();
1407 }
1408 }
1409 describe(ast, target, self.semantics)
1410 }
1411
1412 fn group_items(
1414 &self,
1415 ast: &Ast,
1416 select: &ast::Select,
1417 targets: &[ast::Target],
1418 ) -> Result<Vec<ast::ExprRef>> {
1419 if select.group_by_all {
1420 return Ok(targets
1423 .iter()
1424 .filter(|target| !has_aggregate(ast, target.expr))
1425 .map(|target| target.expr)
1426 .collect());
1427 }
1428 let mut items = Vec::new();
1429 for &item in ast.expr_list(select.group_by) {
1430 items.push(self.output_reference(ast, item, targets, "GROUP BY")?.unwrap_or(item));
1431 }
1432 Ok(items)
1433 }
1434
1435 fn output_reference(
1437 &self,
1438 ast: &Ast,
1439 item: ast::ExprRef,
1440 targets: &[ast::Target],
1441 clause: &str,
1442 ) -> Result<Option<ast::ExprRef>> {
1443 match ast.expr(item) {
1444 ast::Expr::Literal { kind: LiteralKind::Number, text } => {
1445 let written = ast.string(text);
1446 let position: usize = written.parse().map_err(|_| {
1447 Error::binder(format!("{clause} term {written} is not a column"))
1448 })?;
1449 if position == 0 || position > targets.len() {
1450 return Err(Error::binder(format!(
1451 "{clause} term out of range - should be between 1 and {}",
1452 targets.len()
1453 )));
1454 }
1455 Ok(Some(targets[position - 1].expr))
1456 }
1457 ast::Expr::Column { name } => {
1458 let parts: Vec<&str> = ast.name(name).collect();
1459 let [written] = parts.as_slice() else { return Ok(None) };
1460 let mut found = None;
1461 for target in targets {
1462 if target.alias != NONE && same_name(ast.string(target.alias), written) {
1463 if found.is_some() {
1464 return Ok(None);
1465 }
1466 found = Some(target.expr);
1467 }
1468 }
1469 Ok(found)
1470 }
1471 _ => Ok(None),
1472 }
1473 }
1474
1475 #[allow(clippy::too_many_arguments)]
1479 fn select_sort_keys(
1480 &mut self,
1481 ast: &Ast,
1482 query: &ast::Query,
1483 input: &Scope,
1484 output: &Scope,
1485 project: u32,
1486 exprs: &mut Vec<ExprRef>,
1487 names: &mut Vec<String>,
1488 extra: &mut Vec<usize>,
1489 above: &mut Vec<PendingSubquery>,
1490 ) -> Result<Vec<SortKey>> {
1491 if query.order_by_all {
1492 return Ok(self.every_column(output));
1493 }
1494 let items = ast.order_list(query.order_by).to_vec();
1495 let mut keys = Vec::with_capacity(items.len());
1496 for item in items {
1497 self.check_order_literal(ast, item.expr)?;
1498 let position = match self.output_position(ast, item.expr, output)? {
1499 Some(position) => position,
1500 None => {
1501 let before = self.scalar_subqueries.len();
1502 let bound = self.bind_expr(ast, item.expr, input)?;
1503 self.lift_over_aggregate(before, above, input)?;
1504 let bound = self.over_aggregate(bound, input)?;
1505 match exprs.iter().position(|&held| self.same_expr(held, bound)) {
1506 Some(position) => position,
1507 None => {
1508 exprs.push(bound);
1509 names.push(describe(ast, item.expr, self.semantics));
1510 extra.push(exprs.len() - 1);
1511 exprs.len() - 1
1512 }
1513 }
1514 }
1515 };
1516 let ty = self.plan.expr_type(exprs[position]).clone();
1517 let expr = self.column(project, position, ty);
1518 keys.push(self.sort_key(expr, item));
1519 }
1520 Ok(keys)
1521 }
1522
1523 fn sort_keys(
1525 &mut self,
1526 ast: &Ast,
1527 query: &ast::Query,
1528 output: &Scope,
1529 targets: &[ast::Target],
1530 ) -> Result<Vec<SortKey>> {
1531 if query.order_by_all {
1532 return Ok(self.every_column(output));
1533 }
1534 let items = ast.order_list(query.order_by).to_vec();
1535 let mut keys = Vec::with_capacity(items.len());
1536 for item in items {
1537 self.check_order_literal(ast, item.expr)?;
1538 let expr = match self.output_position(ast, item.expr, output)? {
1539 Some(position) => {
1540 let column = &output.columns[position];
1541 let (binding, ty) = (column.binding, column.ty.clone());
1542 self.plan.add_expr(Expr::Column(binding), ty)
1543 }
1544 None => {
1545 let _ = targets;
1546 self.bind_expr(ast, item.expr, output)?
1547 }
1548 };
1549 keys.push(self.sort_key(expr, item));
1550 }
1551 Ok(keys)
1552 }
1553
1554 fn every_column(&mut self, output: &Scope) -> Vec<SortKey> {
1555 let columns: Vec<(ColumnBinding, LogicalType)> =
1556 output.columns.iter().map(|column| (column.binding, column.ty.clone())).collect();
1557 columns
1558 .into_iter()
1559 .map(|(binding, ty)| {
1560 let expr = self.plan.add_expr(Expr::Column(binding), ty);
1561 let expr = self.by_position(expr);
1562 let descending = self.semantics.default_descending();
1563 SortKey { expr, descending, nulls_first: self.semantics.nulls_first(descending) }
1564 })
1565 .collect()
1566 }
1567
1568 fn sort_key(&mut self, expr: ExprRef, item: ast::OrderItem) -> SortKey {
1570 let expr = self.by_position(expr);
1571 let descending = match item.order {
1572 Order::Unstated => self.semantics.default_descending(),
1573 Order::Ascending => false,
1574 Order::Descending => true,
1575 };
1576 let nulls_first = match item.nulls {
1577 Nulls::First => true,
1578 Nulls::Last => false,
1579 Nulls::Unstated => self.semantics.nulls_first(descending),
1580 };
1581 SortKey { expr, descending, nulls_first }
1582 }
1583
1584 fn output_position(
1586 &self,
1587 ast: &Ast,
1588 item: ast::ExprRef,
1589 output: &Scope,
1590 ) -> Result<Option<usize>> {
1591 match ast.expr(item) {
1592 ast::Expr::Literal { kind: LiteralKind::Number, text } => {
1593 let written = ast.string(text);
1594 if written.contains(['.', 'e', 'E']) {
1595 return Ok(None);
1596 }
1597 let position: usize = written.parse().map_err(|_| {
1598 Error::binder(format!("ORDER BY term {written} is not a column"))
1599 })?;
1600 if position == 0 || position > output.len() {
1601 return Err(Error::binder(format!(
1602 "ORDER BY term out of range - should be between 1 and {}",
1603 output.len()
1604 )));
1605 }
1606 Ok(Some(position - 1))
1607 }
1608 ast::Expr::Column { name } => {
1609 let parts: Vec<&str> = ast.name(name).collect();
1610 let [written] = parts.as_slice() else { return Ok(None) };
1611 Ok(output.position_of(None, written))
1612 }
1613 _ => Ok(None),
1614 }
1615 }
1616
1617 fn check_order_literal(&self, ast: &Ast, item: ast::ExprRef) -> Result<()> {
1619 if !self.semantics.order_by_non_integer_literal()
1620 && matches!(
1621 ast.expr(item),
1622 ast::Expr::Literal { kind, text }
1623 if kind != LiteralKind::Number
1624 || ast.string(text).contains(['.', 'e', 'E'])
1625 )
1626 {
1627 return Err(Error::binder(
1628 "ORDER BY non-integer literal has no effect.\n* SET order_by_non_integer_literal=true to allow this behavior.",
1629 ));
1630 }
1631 Ok(())
1632 }
1633
1634 fn distinct_on(
1636 &mut self,
1637 ast: &Ast,
1638 distinct: Distinct,
1639 output: &Scope,
1640 ) -> Result<Vec<ExprRef>> {
1641 let Distinct::On(items) = distinct else {
1642 return Ok(Vec::new());
1643 };
1644 let items = ast.expr_list(items).to_vec();
1645 let mut on = Vec::with_capacity(items.len());
1646 for item in items {
1647 let Some(position) = self.output_position(ast, item, output)? else {
1648 return Err(Error::not_implemented(
1649 "DISTINCT ON an expression that is not in the select list",
1650 ));
1651 };
1652 let column = &output.columns[position];
1653 let (binding, ty) = (column.binding, column.ty.clone());
1654 on.push(self.plan.add_expr(Expr::Column(binding), ty));
1655 }
1656 Ok(on)
1657 }
1658
1659 fn apply_limit(
1666 &mut self,
1667 ast: &Ast,
1668 query: &ast::Query,
1669 input: NodeRef,
1670 scope: &mut Scope,
1671 ) -> Result<NodeRef> {
1672 let waiting = self.scalar_subqueries.len();
1673 if query.limit_percent {
1674 let percent = self.share(ast, query.limit)?;
1675 let offset = self.skipped(ast, query.offset)?;
1676 let node = |binder: &mut Self, input| match percent {
1677 Some(percent) => binder.add_node(Node::LimitPercent { input, percent, offset }),
1678 None => binder.limited(input, Bound::All, offset),
1681 };
1682 return self.over_subqueries(waiting, input, scope, node);
1683 }
1684 let count = self.count_bound(ast, query.limit, "LIMIT")?;
1685 let offset = self.skipped(ast, query.offset)?;
1686 let node = |binder: &mut Self, input| binder.limited(input, count, offset);
1687 self.over_subqueries(waiting, input, scope, node)
1688 }
1689
1690 fn skipped(&mut self, ast: &Ast, written: ast::ExprRef) -> Result<Bound> {
1695 Ok(match self.count_bound(ast, written, "OFFSET")? {
1696 Bound::All => Bound::Rows(0),
1697 named => named,
1698 })
1699 }
1700
1701 fn over_subqueries(
1708 &mut self,
1709 waiting: usize,
1710 input: NodeRef,
1711 scope: &mut Scope,
1712 node: impl FnOnce(&mut Self, NodeRef) -> NodeRef,
1713 ) -> Result<NodeRef> {
1714 let joined = self.scalar_subqueries.split_off(waiting);
1715 if joined.is_empty() {
1716 return Ok(node(self, input));
1717 }
1718 let mut input = input;
1719 for pending in joined {
1720 input = self.attach_subquery(input, pending);
1721 }
1722 let limit = node(self, input);
1723 Ok(self.reproject(limit, scope))
1724 }
1725
1726 fn limited(&mut self, input: NodeRef, count: Bound, offset: Bound) -> NodeRef {
1729 if count == Bound::All && offset == Bound::Rows(0) {
1730 return input;
1731 }
1732 self.add_node(Node::Limit { input, count, offset })
1733 }
1734
1735 fn reproject(&mut self, node: NodeRef, scope: &mut Scope) -> NodeRef {
1741 let index = self.fresh_index();
1742 let mut exprs = Vec::with_capacity(scope.columns.len());
1743 let mut names = Vec::with_capacity(scope.columns.len());
1744 for column in &scope.columns {
1745 exprs.push(self.plan.add_expr(Expr::Column(column.binding), column.ty.clone()));
1746 names.push(self.plan.intern(&column.name));
1747 }
1748 for (at, column) in scope.columns.iter_mut().enumerate() {
1749 column.binding = ColumnBinding::new(index, at as u32);
1750 }
1751 let exprs = self.plan.add_expr_list(&exprs);
1752 let names = self.plan.add_name_list(&names);
1753 self.add_node(Node::Project { input: node, index, exprs, names })
1754 }
1755
1756 fn share(&mut self, ast: &Ast, written: ast::ExprRef) -> Result<Option<Share>> {
1773 if written == NONE {
1774 return Ok(None);
1775 }
1776 self.clause = "LIMIT clause";
1777 let scope = Scope::empty();
1778 let bound = self.bind_expr(ast, written, &scope)?;
1779 let Some(value) = fold::value_of(&self.plan, bound)? else {
1780 return Ok(Some(Share::Read(bound)));
1781 };
1782 if value.is_null() {
1783 return Ok(None);
1784 }
1785 let percent = percentage(&value)?;
1786 if !(0.0..=100.0).contains(&percent) {
1787 return Err(Error::out_of_range(
1788 "Limit percent out of range, should be between 0% and 100%",
1789 ));
1790 }
1791 Ok(Some(Share::Percent(percent)))
1792 }
1793
1794 fn count_bound(&mut self, ast: &Ast, written: ast::ExprRef, clause: &str) -> Result<Bound> {
1813 if written == NONE {
1814 return Ok(Bound::All);
1815 }
1816 self.clause = "LIMIT clause";
1817 let scope = Scope::empty();
1818 let bound = self.bind_expr(ast, written, &scope)?;
1819 let Some(value) = fold::value_of(&self.plan, bound)? else {
1820 return Ok(Bound::Read(bound));
1821 };
1822 if value.is_null() {
1825 return Ok(Bound::All);
1826 }
1827 row_count(&value, clause).map(Bound::Rows)
1828 }
1829
1830 fn bind_from(&mut self, ast: &Ast, from: ast::Slice) -> Result<(NodeRef, Scope)> {
1833 let sources = ast.source_list(from).to_vec();
1834 let Some((first, rest)) = sources.split_first() else {
1835 return Ok((self.add_node(Node::Dummy), Scope::empty()));
1838 };
1839 let (mut node, mut scope) = self.bind_source(ast, *first)?;
1840 for source in rest {
1841 let (right, right_scope, correlations) = self.bind_lateral(ast, *source, &scope)?;
1842 node = if correlations.is_empty() {
1843 self.add_node(Node::CrossProduct { left: node, right })
1844 } else {
1845 let conditions = self.plan.add_expr_list(&[]);
1846 self.add_node(Node::DependentJoin {
1847 left: node,
1848 right,
1849 kind: JoinKind::Inner,
1850 conditions,
1851 })
1852 };
1853 scope = scope.concat(right_scope);
1854 }
1855 Ok((node, scope))
1856 }
1857
1858 fn bind_lateral(
1870 &mut self,
1871 ast: &Ast,
1872 source: ast::SourceRef,
1873 left: &Scope,
1874 ) -> Result<(NodeRef, Scope, Vec<ColumnBinding>)> {
1875 self.lateral_scopes.push(self.outer_scopes.len());
1876 self.outer_scopes.push(left.clone());
1877 self.correlations.push(Vec::new());
1878 let bound = self.bind_source(ast, source);
1879 let read = self.correlations.pop().expect("correlation frame");
1880 self.outer_scopes.pop();
1881 self.lateral_scopes.pop();
1882 let (node, scope) = bound?;
1883
1884 let mut here = Vec::new();
1885 for binding in read {
1886 if left.columns.iter().any(|column| column.binding == binding) {
1887 here.push(binding);
1888 } else if let Some(enclosing) = self.correlations.last_mut()
1889 && !enclosing.contains(&binding)
1890 {
1891 enclosing.push(binding);
1892 }
1893 }
1894 Ok((node, scope, here))
1904 }
1905
1906 fn bind_source(&mut self, ast: &Ast, source: ast::SourceRef) -> Result<(NodeRef, Scope)> {
1907 match ast.source(source) {
1908 ast::Source::Table { name, alias, columns } => {
1909 self.bind_table(ast, name, alias, columns)
1910 }
1911 ast::Source::Function { name, args, alias, columns, pragma } => {
1912 self.bind_table_function(ast, name, args, alias, columns, pragma)
1913 }
1914 ast::Source::Subquery { query, alias, columns } => {
1915 let (node, mut scope) = self.bind_query(ast, query)?;
1916 let label = if alias == NONE {
1917 "unnamed_subquery".to_string()
1918 } else {
1919 ast.string(alias).to_string()
1920 };
1921 scope.relabel(&label);
1922 if !columns.is_empty() {
1923 let names: Vec<&str> = ast.name(columns).collect();
1924 scope.rename(&names, &label)?;
1925 }
1926 Ok((node, scope))
1927 }
1928 ast::Source::Values { rows, alias, columns } => {
1929 let bare = ast::Query::bare(ast::QueryBody::Values(rows));
1930 let (node, mut scope) = self.bind_values(ast, &bare, rows)?;
1931 let label =
1932 if alias == NONE { String::new() } else { ast.string(alias).to_string() };
1933 scope.relabel(&label);
1934 if !columns.is_empty() {
1935 let names: Vec<&str> = ast.name(columns).collect();
1936 scope.rename(&names, &label)?;
1937 }
1938 Ok((node, scope))
1939 }
1940 ast::Source::Cte { cte, alias, columns } => {
1941 self.bind_cte_scan(ast, cte, alias, columns)
1942 }
1943 ast::Source::Join { left, right, kind, natural, on, using } => {
1944 self.bind_join(ast, left, right, kind, natural, on, using)
1945 }
1946 }
1947 }
1948
1949 fn bind_cte_scan(
1956 &mut self,
1957 ast: &Ast,
1958 written: u32,
1959 alias: ast::StrRef,
1960 columns: ast::Slice,
1961 ) -> Result<(NodeRef, Scope)> {
1962 let Some(held) = self.materialized.iter().rev().find(|held| held.written == written) else {
1963 let name = ast.string(ast.cte(written).name);
1964 return Err(Error::binder(format!("Table with name {name} does not exist!")));
1965 };
1966 let cte = held.cte;
1967 let fields = held.fields.clone();
1968 let text = held.name.clone();
1969 let label = if alias == NONE { text.clone() } else { ast.string(alias).to_string() };
1970 let name = self.plan.intern(&text);
1971 let index = self.fresh_index();
1972 let mut scope = Scope::empty();
1973 for (at, field) in fields.iter().enumerate() {
1974 scope.push(Visible {
1975 table: label.clone(),
1976 name: field.name.clone(),
1977 binding: ColumnBinding::new(index, at as u32),
1978 ty: field.ty.clone(),
1979 not_null: field.not_null,
1980 key: None,
1981 default: None,
1982 qualified: false,
1983 also: None,
1984 });
1985 }
1986 if !columns.is_empty() {
1987 let names: Vec<&str> = ast.name(columns).collect();
1988 scope.rename(&names, &label)?;
1989 }
1990 let columns = self.plan.add_fields(&fields);
1991 let node = self.add_node(Node::CteScan { index, cte, name, columns });
1992 Ok((node, scope))
1993 }
1994
1995 fn bind_table(
1996 &mut self,
1997 ast: &Ast,
1998 name: ast::Slice,
1999 alias: ast::StrRef,
2000 columns: ast::Slice,
2001 ) -> Result<(NodeRef, Scope)> {
2002 let parts: Vec<&str> = ast.name(name).collect();
2003 let catalog = self.catalog;
2004 let resolved = match catalog.resolve(&parts) {
2007 Ok(resolved) => resolved,
2008 Err(missing) => {
2009 return self.bind_replacement_scan(ast, &parts, alias, columns, missing);
2010 }
2011 };
2012 if catalog.entry(&resolved)? == Entry::View {
2013 return self.bind_view(ast, &resolved, alias, columns);
2014 }
2015 let label =
2016 if alias == NONE { resolved.table.clone() } else { ast.string(alias).to_string() };
2017 self.bind_catalog_table(ast, &resolved, label, columns)
2018 }
2019
2020 pub(crate) fn bind_catalog_table(
2023 &mut self,
2024 ast: &Ast,
2025 resolved: &QualifiedName,
2026 label: String,
2027 columns: ast::Slice,
2028 ) -> Result<(NodeRef, Scope)> {
2029 let catalog = self.catalog;
2030 let table = catalog.table(resolved)?;
2031 let fields: &[Field] = table.columns();
2032 let excluded = self.upsert && same_name(&label, "excluded");
2035 let mut marks = vec![None; fields.len()];
2038 for key in table.keys() {
2039 for &column in &key.columns {
2040 if key.primary || marks[column].is_none() {
2041 marks[column] = Some(if key.primary { "PRI" } else { "UNI" });
2042 }
2043 }
2044 }
2045 let index = self.fresh_index();
2046 let mut scope = Scope::empty();
2047 for (at, field) in fields.iter().enumerate() {
2048 scope.push(Visible {
2049 table: label.clone(),
2050 name: field.name.clone(),
2051 binding: ColumnBinding::new(index, at as u32),
2052 ty: field.ty.clone(),
2053 not_null: field.not_null,
2054 key: marks[at],
2055 default: table.default(at).map(str::to_owned),
2056 qualified: excluded,
2057 also: None,
2058 });
2059 }
2060 if !columns.is_empty() {
2061 let names: Vec<&str> = ast.name(columns).collect();
2062 scope.rename(&names, &label)?;
2063 }
2064 let resolved = if excluded { &QualifiedName::excluded() } else { resolved };
2065 let catalog_name = self.plan.intern(&resolved.catalog);
2066 let schema = self.plan.intern(&resolved.schema);
2067 let table_name = self.plan.intern(&resolved.table);
2068 let alias = self.plan.intern(&label);
2069 let columns = self.plan.add_fields(fields);
2070 if let Some(zones) = table.rows().zones().filter(|_| !excluded) {
2075 self.plan.set_zones(index, zones);
2076 }
2077 if let Some(frequencies) = table.frequencies().filter(|_| !excluded) {
2078 self.plan.set_frequencies(index, frequencies);
2079 }
2080 let facts = table.facts().filter(|_| !excluded);
2083 if let Some(facts) = facts {
2084 self.plan.set_facts(index, facts, self.want_ascending);
2085 } else if !excluded {
2086 for (column, distinct) in table.distincts() {
2087 self.plan.measure_distinct(index, &column, distinct);
2088 }
2089 if self.want_ascending {
2090 for column in table.ascending() {
2091 self.plan.mark_ascending(index, &column);
2092 }
2093 }
2094 for (column, bytes) in table.widths() {
2095 self.plan.measure_width(index, &column, bytes);
2096 }
2097 }
2098 let node = self.add_node(Node::Get {
2099 catalog: catalog_name,
2100 schema,
2101 table: table_name,
2102 alias,
2103 index,
2104 columns,
2105 });
2106 Ok((node, scope))
2107 }
2108
2109 fn bind_view(
2121 &mut self,
2122 ast: &Ast,
2123 name: &QualifiedName,
2124 alias: ast::StrRef,
2125 columns: ast::Slice,
2126 ) -> Result<(NodeRef, Scope)> {
2127 let view = self.catalog.view(name)?;
2128 let full = name.to_string();
2129 if self.expanding.contains(&full) {
2130 return Err(Error::binder(format!(
2134 "infinite recursion detected: attempting to recursively bind view \"\"{}\"\"",
2135 name.table
2136 )));
2137 }
2138 let body = parse_ast_with_case(view.sql(), self.semantics.identifier_case())?;
2139 let query = match body.statements.as_slice() {
2140 [ast::Statement::Query(query)] => *query,
2141 _ => return Err(Error::binder(format!("view \"{}\" is not a query", name.table))),
2144 };
2145 self.expanding.push(full);
2146 let bound = self.bind_query(&body, query);
2147 self.expanding.pop();
2148 let (node, mut scope) = bound?;
2149
2150 let aliases: Vec<&str> = view.aliases().iter().map(String::as_str).collect();
2151 if !aliases.is_empty() {
2152 scope.rename(&aliases, "unnamed_subquery")?;
2153 }
2154 view.remember(scope.fields());
2161 let label = if alias == NONE { name.table.clone() } else { ast.string(alias).to_string() };
2162 scope.relabel(&label);
2163 if !columns.is_empty() {
2164 let names: Vec<&str> = ast.name(columns).collect();
2165 scope.rename(&names, &label)?;
2166 }
2167 Ok((node, scope))
2168 }
2169
2170 fn bind_table_function(
2178 &mut self,
2179 ast: &Ast,
2180 name: ast::Slice,
2181 args: ast::Slice,
2182 alias: ast::StrRef,
2183 columns: ast::Slice,
2184 pragma: bool,
2185 ) -> Result<(NodeRef, Scope)> {
2186 let renamed = columns;
2189 let parts: Vec<&str> = ast.name(name).collect();
2190 let function_name = *parts.last().unwrap_or(&"");
2194 if let Some(schema) = parts.iter().rev().nth(1)
2195 && !schema.eq_ignore_ascii_case("main")
2196 && !schema.eq_ignore_ascii_case("system")
2197 {
2198 return Err(Error::catalog(format!(
2199 "Table Function with name {} does not exist!",
2200 parts.join(".")
2201 )));
2202 }
2203 let Some(called) = TableFunction::lookup(function_name) else {
2207 if pragma {
2208 if args.is_empty() && self.catalog.resolve(&parts).is_ok() {
2214 return self.bind_table(ast, name, alias, columns);
2215 }
2216 let spelled = function_name.strip_prefix("pragma_").unwrap_or(function_name);
2217 return Err(Error::catalog(format!(
2218 "Pragma Function with name {spelled} does not exist!"
2219 )));
2220 }
2221 return Err(Error::catalog(format!(
2222 "Table Function with name {function_name} does not exist!"
2223 )));
2224 };
2225 let written = ast.target_list(args).to_vec();
2226 let empty = Scope::empty();
2227 let waiting = self.scalar_subqueries.len();
2228 let previous = std::mem::replace(&mut self.clause, "table function arguments");
2229 let mut bound = Vec::new();
2230 let mut written_options = Vec::new();
2231 for argument in written {
2232 let expr = self.bind_expr(ast, argument.expr, &empty)?;
2233 if argument.alias == NONE {
2234 bound.push(expr);
2235 } else {
2236 let name = ast.string(argument.alias).to_string();
2237 let (parameter, value) = self.named_argument(called, &name, expr)?;
2238 written_options.push((parameter, value, expr));
2239 }
2240 }
2241 self.clause = previous;
2242 let options = Options::of(&written_options)?;
2243
2244 let given: Vec<LogicalType> =
2247 bound.iter().map(|&expr| self.plan.expr_type(expr).clone()).collect();
2248 let resolved = if pragma {
2249 resolve_pragma(function_name, &given)?
2250 } else {
2251 resolve_table(function_name, &given)?
2252 };
2253 let mut cast: Vec<ExprRef> = bound
2254 .iter()
2255 .zip(&resolved.arguments)
2256 .map(|(&expr, ty)| self.checked_cast_to(expr, ty, false))
2257 .collect::<Result<_>>()?;
2258
2259 if resolved.function.answered_when_bound() {
2260 let Columns::Fixed(fields) = resolved.columns else {
2261 return Err(Error::internal("a pragma that resolved to a file"));
2262 };
2263 let [argument] = cast[..] else {
2264 return Err(Error::internal("a pragma that resolved to more than one name"));
2265 };
2266 return self.bind_pragma(ast, resolved.function, &fields, argument, alias, columns);
2267 }
2268 let mut measured = Stat::Unknown;
2271 let mut counted: Vec<(String, Stat<u64>)> = Vec::new();
2272 let mut bounded: Option<Arc<dyn Zones>> = None;
2273 let fields = match resolved.columns {
2274 Columns::Fixed(fields) => fields,
2275 columns => {
2276 let paths = self.file_paths(cast[0], resolved.function.name())?;
2281 let mut mirrorable = None;
2282 if resolved.function == TableFunction::ReadParquet
2283 && !options.file_row_number
2284 && let Some((path, stamp)) = mirror_target(&paths)
2285 {
2286 if let Some(name) = self.catalog.mirror(&path, options.binary_as_string, stamp)
2287 {
2288 let name = name.clone();
2289 let label = if alias == NONE {
2290 resolved.function.name().to_string()
2291 } else {
2292 ast.string(alias).to_string()
2293 };
2294 return self.bind_catalog_table(ast, &name, label, renamed);
2295 }
2296 mirrorable = Some(path);
2297 }
2298 let copy_into = match columns {
2299 Columns::Csv => self.copy_into.take(),
2300 _ => None,
2301 };
2302 let mut fields = match columns {
2303 Columns::Csv if copy_into.is_some() => {
2304 let into = copy_into.unwrap_or_default();
2305 self.copy_fields(&paths, &options.given, &into, &mut written_options)?
2306 }
2307 Columns::Csv => csv_fields(&paths, options.given.clone())?,
2310 _ => {
2311 let footers = self.footers(&paths, mirrorable.as_deref())?;
2312 if let Some(path) = mirrorable.as_deref() {
2313 self.want_mirror(path, options.binary_as_string, &footers.rows);
2314 }
2315 measured = footers.rows;
2316 counted = footers.distincts;
2317 bounded = footers.zones;
2318 footers.fields
2319 }
2320 };
2321 if options.all_varchar {
2322 for field in &mut fields {
2327 field.ty = LogicalType::Varchar;
2328 }
2329 }
2330 if options.binary_as_string {
2331 for field in &mut fields {
2336 if field.ty == LogicalType::Blob {
2337 field.ty = LogicalType::Varchar;
2338 }
2339 }
2340 }
2341 if options.file_row_number {
2342 if fields.iter().any(|field| field.name == FILE_ROW_NUMBER) {
2348 return Err(Error::binder(format!(
2349 "Duplicate column name \"{FILE_ROW_NUMBER}\": the file already has a \
2350 column of that name, so file_row_number cannot add one"
2351 )));
2352 }
2353 fields.push(Field::required(FILE_ROW_NUMBER.to_string(), LogicalType::BigInt));
2354 }
2355 cast = paths.iter().map(|path| self.path_constant(path)).collect();
2356 fields
2357 }
2358 };
2359 let label = if alias == NONE {
2360 resolved.function.name().to_string()
2361 } else {
2362 ast.string(alias).to_string()
2363 };
2364 let names: Vec<&str> = ast.name(columns).collect();
2365 let (node, scope) = self.table_function_source(
2366 resolved.function,
2367 &cast,
2368 &written_options,
2369 Read { fields, rows: measured, distincts: counted, zones: bounded },
2370 &label,
2371 &names,
2372 )?;
2373 Ok((self.lateral_over_subqueries(node, waiting), scope))
2374 }
2375
2376 fn lateral_over_subqueries(&mut self, node: NodeRef, waiting: usize) -> NodeRef {
2383 if self.scalar_subqueries.len() <= waiting {
2384 return node;
2385 }
2386 let Node::TableFunction { index, function, args, options, settings, columns } =
2387 self.plan.node(node).clone()
2388 else {
2389 return node;
2390 };
2391 let series = matches!(
2392 TableFunction::lookup(self.plan.string(function)),
2393 Some(TableFunction::Range | TableFunction::GenerateSeries | TableFunction::Unnest)
2394 );
2395 if !series {
2396 return node;
2397 }
2398 let mut input = self.add_node(Node::Dummy);
2399 for pending in self.scalar_subqueries.split_off(waiting) {
2400 input = self.attach_subquery(input, pending);
2401 }
2402 self.add_node(Node::LateralFunction {
2403 input,
2404 index,
2405 function,
2406 args,
2407 options,
2408 settings,
2409 columns,
2410 })
2411 }
2412
2413 fn copy_fields(
2426 &mut self,
2427 paths: &[String],
2428 given: &Given,
2429 into: &[Field],
2430 written: &mut Vec<(&'static str, Value, ExprRef)>,
2431 ) -> Result<Vec<Field>> {
2432 let sniffed = csv_fields(paths, given.clone())?;
2433 if sniffed.len() != into.len() {
2434 let set: Vec<String> =
2435 into.iter().map(|field| format!("'{}' : '{}'", field.name, field.ty)).collect();
2436 return Err(Error::invalid_input(format!(
2437 "Error when sniffing file \"{}\".\nIt was not possible to automatically detect the \
2438 CSV parsing dialect\n* Columns are set as: \"columns = {{ {}}}\", and they \
2439 contain: {} columns. It does not match the number of columns found by the \
2440 sniffer: {}. Verify the columns parameter is correctly set.",
2441 paths.first().map_or("", String::as_str),
2442 set.join(", "),
2443 into.len(),
2444 sniffed.len()
2445 )));
2446 }
2447 let names: Vec<Value> =
2448 into.iter().map(|field| Value::Varchar(field.name.clone())).collect();
2449 let names = Value::List { element: LogicalType::Varchar, values: names };
2450 let list = LogicalType::List(Box::new(LogicalType::Varchar));
2451 for (parameter, value, ty) in
2452 [("names", names, list), (TYPES_SET, Value::Boolean(true), LogicalType::Boolean)]
2453 {
2454 let reference = self.plan.add_value(value.clone());
2455 let expr = self.plan.add_expr(Expr::Constant(reference), ty);
2456 written.push((parameter, value, expr));
2457 }
2458 Ok(into.to_vec())
2459 }
2460
2461 fn bind_pragma(
2474 &mut self,
2475 ast: &Ast,
2476 function: TableFunction,
2477 fields: &[Field],
2478 argument: ExprRef,
2479 alias: ast::StrRef,
2480 columns: ast::Slice,
2481 ) -> Result<(NodeRef, Scope)> {
2482 let written = self.pragma_name(argument, function)?;
2483 let parts = identifier_parts(&written);
2484 let spelled: Vec<&str> = parts.iter().map(String::as_str).collect();
2485 let name = self.catalog.resolve(&spelled)?;
2486 let described = self.described(ast, &name)?;
2487 let mut rows = Vec::with_capacity(described.len());
2488 for (at, field) in described.iter().enumerate() {
2489 let items = if matches!(function, TableFunction::PragmaShow) {
2490 self.describing(field)
2491 } else {
2492 self.table_info(at, field)
2493 };
2494 rows.push(self.plan.add_expr_list(&items));
2495 }
2496 let rows = self.plan.add_rows(&rows);
2497 let held = self.plan.add_fields(fields);
2498 let index = self.fresh_index();
2499 let node = self.add_node(Node::Values { index, columns: held, rows });
2500 let label =
2501 if alias == NONE { function.name().to_string() } else { ast.string(alias).to_string() };
2502 let mut scope = Scope::empty();
2503 for (at, field) in fields.iter().enumerate() {
2504 scope.push(Visible {
2505 table: label.clone(),
2506 name: field.name.clone(),
2507 binding: ColumnBinding::new(index, at as u32),
2508 ty: field.ty.clone(),
2509 not_null: false,
2510 key: None,
2511 default: None,
2512 qualified: false,
2513 also: None,
2514 });
2515 }
2516 if !columns.is_empty() {
2517 let names: Vec<&str> = ast.name(columns).collect();
2518 scope.rename(&names, &label)?;
2519 }
2520 Ok((node, scope))
2521 }
2522
2523 fn pragma_name(&self, argument: ExprRef, function: TableFunction) -> Result<String> {
2533 let Expr::Constant(reference) = *self.plan.expr(argument) else {
2534 return Err(Error::not_implemented(format!(
2535 "{}() given a name that is not a constant",
2536 function.name()
2537 )));
2538 };
2539 match self.plan.value(reference) {
2540 Value::Varchar(name) => Ok(name.clone()),
2541 Value::Null => Ok("NULL".to_string()),
2542 other => {
2543 Err(Error::internal(format!("a pragma name bound as VARCHAR arrived as {other}")))
2544 }
2545 }
2546 }
2547
2548 fn described(&mut self, ast: &Ast, name: &QualifiedName) -> Result<Vec<Field>> {
2559 if self.catalog.entry(name)? == Entry::Table {
2560 return Ok(self.catalog.table(name)?.columns().to_vec());
2561 }
2562 let (_, scope) = self.bind_view(ast, name, NONE, ast::Slice::default())?;
2563 Ok(scope.fields())
2564 }
2565
2566 fn describing(&mut self, field: &Field) -> Vec<ExprRef> {
2568 let written = [
2569 field.name.clone(),
2570 field.ty.to_string(),
2571 if field.not_null { "NO" } else { "YES" }.to_owned(),
2572 ];
2573 let mut items: Vec<ExprRef> =
2574 written.into_iter().map(|text| self.plan.add_constant(Value::Varchar(text))).collect();
2575 for _ in 0..3 {
2576 let empty = self.plan.add_constant(Value::Null);
2577 items.push(self.cast_to(empty, &LogicalType::Varchar));
2578 }
2579 items
2580 }
2581
2582 fn table_info(&mut self, at: usize, field: &Field) -> Vec<ExprRef> {
2588 let cid = self.plan.add_constant(Value::Integer(i32::try_from(at).unwrap_or(i32::MAX)));
2589 let name = self.plan.add_constant(Value::Varchar(field.name.clone()));
2590 let ty = self.plan.add_constant(Value::Varchar(field.ty.to_string()));
2591 let not_null = self.plan.add_constant(Value::Boolean(field.not_null));
2592 let default = self.plan.add_constant(Value::Null);
2593 let default = self.cast_to(default, &LogicalType::Varchar);
2594 let key = self.plan.add_constant(Value::Boolean(false));
2595 vec![cid, name, ty, not_null, default, key]
2596 }
2597
2598 fn named_argument(
2612 &mut self,
2613 function: TableFunction,
2614 name: &str,
2615 expr: ExprRef,
2616 ) -> Result<(&'static str, Value)> {
2617 let known = function
2618 .parameters()
2619 .iter()
2620 .find(|(parameter, _)| parameter.eq_ignore_ascii_case(name));
2621 let Some((parameter, wanted)) = known else {
2622 let candidates: Vec<String> = function
2623 .parameters()
2624 .iter()
2625 .map(|(parameter, ty)| format!(" {parameter} {ty}"))
2626 .collect();
2627 if candidates.is_empty() {
2629 return Err(Error::binder(format!(
2630 "Invalid named parameter \"{name}\" for function {}\nFunction does not \
2631 accept any named parameters.",
2632 function.name()
2633 )));
2634 }
2635 return Err(Error::binder(format!(
2636 "Invalid named parameter \"{name}\" for function {}\nCandidates:\n{}\n",
2637 function.name(),
2638 candidates.join("\n")
2639 )));
2640 };
2641 let Some(value) = fold::value_of(&self.plan, expr)? else {
2644 return Err(Error::not_implemented(format!(
2645 "the named parameter {parameter} with a value that is not a constant"
2646 )));
2647 };
2648 if value == Value::Null {
2649 return Err(Error::binder(null_parameter(function, parameter)));
2650 }
2651 let given = self.plan.expr_type(expr).clone();
2652 let listed = *parameter == "nullstr" && given == LogicalType::list(LogicalType::Varchar);
2655 if *parameter == "nullstr" && given != *wanted && !listed {
2656 return Err(Error::binder(
2657 "CSV Reader function option \"nullstr\" requires a string or a list as input",
2658 ));
2659 }
2660 if given != *wanted && !listed {
2661 return Err(Error::not_implemented(format!(
2662 "the named parameter {parameter} given a {given} where a {wanted} was wanted"
2663 )));
2664 }
2665 Ok((parameter, value))
2666 }
2667
2668 fn bind_replacement_scan(
2679 &mut self,
2680 ast: &Ast,
2681 parts: &[&str],
2682 alias: ast::StrRef,
2683 columns: ast::Slice,
2684 missing: Error,
2685 ) -> Result<(NodeRef, Scope)> {
2686 let [path] = parts else { return Err(missing) };
2687 let path = *path;
2688 let extension = path.rsplit_once('.').map(|(_, after)| after).unwrap_or_default();
2689 let Some(function) = Self::reader_for(extension) else {
2690 if is_file(path) {
2691 return Err(Error::binder(format!(
2696 "No extension found that is capable of reading the file \"{path}\"\n* If this \
2697 file is a supported file format you can explicitly use the reader functions, \
2698 such as read_csv, read_json or read_parquet"
2699 )));
2700 }
2701 return Err(missing);
2702 };
2703 let paths = files(path)?;
2708 let label = if alias == NONE {
2714 if is_pattern(path) {
2715 path.to_string()
2716 } else {
2717 let file = path.rsplit_once('/').map_or(path, |(_, file)| file);
2718 file.rsplit_once('.').map_or(file, |(stem, _)| stem).to_string()
2719 }
2720 } else {
2721 ast.string(alias).to_string()
2722 };
2723 let mut mirrorable = None;
2724 if function == TableFunction::ReadParquet
2725 && let Some((canonical, stamp)) = mirror_target(&paths)
2726 {
2727 if let Some(name) = self.catalog.mirror(&canonical, false, stamp) {
2728 let name = name.clone();
2729 return self.bind_catalog_table(ast, &name, label, columns);
2730 }
2731 mirrorable = Some(canonical);
2732 }
2733 let read = match function {
2734 TableFunction::ReadParquet => {
2735 let footers = self.footers(&paths, mirrorable.as_deref())?;
2736 if let Some(canonical) = mirrorable.as_deref() {
2737 self.want_mirror(canonical, false, &footers.rows);
2738 }
2739 Read {
2740 fields: footers.fields,
2741 rows: footers.rows,
2742 distincts: footers.distincts,
2743 zones: footers.zones,
2744 }
2745 }
2746 _ => Read::uncounted(csv_fields(&paths, Given::default())?),
2747 };
2748 let arguments: Vec<ExprRef> = paths.iter().map(|path| self.path_constant(path)).collect();
2749 let names: Vec<&str> = ast.name(columns).collect();
2750 self.table_function_source(function, &arguments, &[], read, &label, &names)
2751 }
2752
2753 fn footers(&self, paths: &[String], mirrorable: Option<&str>) -> Result<Footers> {
2759 if let Some(path) = mirrorable.filter(|_| self.outlined) {
2760 let outline = parquet_outline(path)?;
2761 if outline.rows.value().is_some() {
2762 return Ok(outline);
2763 }
2764 }
2765 parquet_footers(paths)
2766 }
2767
2768 fn want_mirror(&mut self, path: &str, binary_as_string: bool, rows: &Stat<u64>) {
2772 if let Some(&rows) = rows.value() {
2773 self.plan.want_mirror(path, binary_as_string, rows);
2774 }
2775 }
2776
2777 fn path_constant(&mut self, path: &str) -> ExprRef {
2779 let value = self.plan.add_value(Value::Varchar(path.to_string()));
2780 self.plan.add_expr(Expr::Constant(value), LogicalType::Varchar)
2781 }
2782
2783 fn reader_for(extension: &str) -> Option<TableFunction> {
2790 if extension.eq_ignore_ascii_case("parquet") {
2791 return Some(TableFunction::ReadParquet);
2792 }
2793 if extension.eq_ignore_ascii_case("csv") || extension.eq_ignore_ascii_case("tsv") {
2794 return Some(TableFunction::ReadCsv);
2795 }
2796 None
2797 }
2798
2799 fn table_function_source(
2809 &mut self,
2810 function: TableFunction,
2811 args: &[ExprRef],
2812 written: &[(&'static str, Value, ExprRef)],
2813 read: Read,
2814 label: &str,
2815 names: &[&str],
2816 ) -> Result<(NodeRef, Scope)> {
2817 let Read { fields, rows, distincts, zones } = read;
2818 let index = self.fresh_index();
2819 if rows.is_known() {
2824 self.plan.measure(index, rows);
2825 }
2826 for (column, distinct) in distincts {
2827 self.plan.measure_distinct(index, &column, distinct);
2828 }
2829 if let Some(zones) = zones {
2830 self.plan.set_zones(index, zones);
2831 }
2832 let mut scope = Scope::empty();
2833 for (at, field) in fields.iter().enumerate() {
2834 scope.push(Visible {
2835 table: label.to_string(),
2836 name: field.name.clone(),
2837 binding: ColumnBinding::new(index, at as u32),
2838 ty: field.ty.clone(),
2839 not_null: false,
2842 key: None,
2843 default: None,
2844 qualified: false,
2845 also: None,
2846 });
2847 }
2848 if !names.is_empty() {
2849 scope.rename(names, label)?;
2850 } else if matches!(
2851 function,
2852 TableFunction::Range | TableFunction::GenerateSeries | TableFunction::Unnest
2853 ) {
2854 for column in &mut scope.columns {
2857 column.also = Some(std::mem::replace(&mut column.name, label.to_string()));
2858 }
2859 }
2860 let function = self.plan.intern(function.name());
2861 let args = self.plan.add_expr_list(args);
2862 let named: Vec<u32> =
2863 written.iter().map(|(parameter, _, _)| self.plan.intern(parameter)).collect();
2864 let settings: Vec<ExprRef> = written.iter().map(|(_, _, expr)| *expr).collect();
2865 let options = self.plan.add_name_list(&named);
2866 let settings = self.plan.add_expr_list(&settings);
2867 let columns = self.plan.add_fields(&fields);
2868 let node = self.add_node(Node::TableFunction {
2869 index,
2870 function,
2871 args,
2872 options,
2873 settings,
2874 columns,
2875 });
2876 Ok((node, scope))
2877 }
2878
2879 fn file_paths(&self, expr: ExprRef, name: &str) -> Result<Vec<String>> {
2886 let mut paths = Vec::new();
2887 for pattern in self.file_patterns(expr, name)? {
2888 paths.extend(files(&pattern)?);
2889 }
2890 Ok(paths)
2891 }
2892
2893 fn file_patterns(&self, expr: ExprRef, name: &str) -> Result<Vec<String>> {
2912 let Some(value) = fold::value_of(&self.plan, expr)? else {
2913 return Err(Error::not_implemented(
2914 "a table function file name that is not a constant",
2915 ));
2916 };
2917 match value {
2918 Value::Varchar(path) => Ok(vec![path]),
2919 Value::Null => Err(Error::parser(format!("{name} cannot take NULL list as parameter"))),
2921 Value::List { values, .. } if values.is_empty() => {
2926 Err(Error::io(format!("\"{name}\" needs at least one file to read")))
2927 }
2928 Value::List { values, .. } => values
2929 .iter()
2930 .map(|value| match value {
2931 Value::Varchar(path) => Ok(path.clone()),
2932 _ => Err(Error::parser(format!(
2933 "{name} reader cannot take NULL input as parameter"
2934 ))),
2935 })
2936 .collect(),
2937 other => {
2938 Err(Error::internal(format!("a file name bound as VARCHAR arrived as {other}")))
2939 }
2940 }
2941 }
2942
2943 fn side_of(
2961 &self,
2962 pending: &PendingSubquery,
2963 left_tables: &[u32],
2964 right_tables: &[u32],
2965 ) -> Option<Side> {
2966 let mut needs_left = false;
2967 let mut needs_right = false;
2968 let mut note = |binding: ColumnBinding| {
2969 needs_left |= left_tables.contains(&binding.table);
2970 needs_right |= right_tables.contains(&binding.table);
2971 };
2972 for &binding in &pending.reads {
2973 note(binding);
2974 }
2975 for &condition in &pending.conditions {
2980 self.plan.read_columns(condition, &mut |_, binding| note(binding));
2981 }
2982 match (needs_left, needs_right) {
2983 (true, true) => None,
2984 (_, true) => Some(Side::Right),
2985 _ => Some(Side::Left),
2986 }
2987 }
2988
2989 #[allow(clippy::too_many_arguments)]
3010 fn bind_pair_dependent_join(
3011 &mut self,
3012 kind: ast::JoinKind,
3013 independent: bool,
3014 left: NodeRef,
3015 right: NodeRef,
3016 pair: Vec<PendingSubquery>,
3017 conditions: Vec<ExprRef>,
3018 scope: Scope,
3019 ) -> Result<(NodeRef, Scope)> {
3020 if kind != ast::JoinKind::Inner {
3021 return Err(Error::not_implemented(
3022 "a subquery that reads both sides of that join, written in the condition of a join \
3023 that is not an inner join"
3024 .to_string(),
3025 ));
3026 }
3027 if !independent {
3030 return Err(Error::not_implemented(
3031 "a subquery that reads both sides of that join, written in the condition of a join \
3032 whose right side is lateral"
3033 .to_string(),
3034 ));
3035 }
3036 let mut node = self.add_node(Node::CrossProduct { left, right });
3037 for pending in pair {
3038 node = self.attach_subquery(node, pending);
3039 }
3040 let mut conditions = conditions.into_iter();
3044 let mut predicate = conditions.next().expect("a join condition was bound");
3045 for next in conditions {
3046 let children = self.plan.add_expr_list(&[predicate, next]);
3047 let conjunction = Expr::Conjunction { op: ConjunctionOp::And, children };
3048 predicate = self.plan.add_expr(conjunction, LogicalType::Boolean);
3049 }
3050 let node = self.add_node(Node::Filter { input: node, predicate });
3051 Ok((node, scope))
3052 }
3053
3054 #[allow(clippy::too_many_arguments)]
3055 fn bind_join(
3056 &mut self,
3057 ast: &Ast,
3058 left: ast::SourceRef,
3059 right: ast::SourceRef,
3060 kind: ast::JoinKind,
3061 natural: bool,
3062 on: ast::ExprRef,
3063 using: ast::Slice,
3064 ) -> Result<(NodeRef, Scope)> {
3065 let (left_node, left_scope) = self.bind_source(ast, left)?;
3066 let (right_node, right_scope, correlated) = self.bind_lateral(ast, right, &left_scope)?;
3067 if !correlated.is_empty()
3071 && !matches!(kind, ast::JoinKind::Inner | ast::JoinKind::Cross | ast::JoinKind::Left)
3072 {
3073 return Err(Error::binder(
3074 "The combining JOIN type must be INNER or LEFT for a LATERAL reference",
3075 ));
3076 }
3077 let split = left_scope.len();
3078 let left_tables: Vec<u32> =
3084 left_scope.columns.iter().map(|column| column.binding.table).collect();
3085 let right_tables: Vec<u32> =
3086 right_scope.columns.iter().map(|column| column.binding.table).collect();
3087 let mut scope = left_scope.concat(right_scope);
3088
3089 let merged: Vec<String> = if natural {
3092 let mut names = Vec::new();
3093 for (at, column) in scope.columns.iter().enumerate().take(split) {
3094 if scope.columns[split..].iter().any(|right| same_name(&right.name, &column.name))
3095 && !names.iter().any(|held: &String| same_name(held, &column.name))
3096 {
3097 let _ = at;
3098 names.push(column.name.clone());
3099 }
3100 }
3101 names
3102 } else {
3103 let mut names: Vec<String> = Vec::new();
3109 for name in ast.name(using) {
3110 if !names.iter().any(|held| same_name(held, name)) {
3111 names.push(name.to_string());
3112 }
3113 }
3114 names
3115 };
3116
3117 let mut conditions = Vec::new();
3118 let mut dropped = Vec::new();
3119 for name in &merged {
3120 let left_at = scope.columns[..split]
3121 .iter()
3122 .position(|column| same_name(&column.name, name))
3123 .ok_or_else(|| {
3124 Error::binder(format!(
3125 "column \"{name}\" specified in USING clause does not exist in left table"
3126 ))
3127 })?;
3128 let right_at = scope.columns[split..]
3129 .iter()
3130 .position(|column| same_name(&column.name, name))
3131 .map(|at| at + split)
3132 .ok_or_else(|| {
3133 Error::binder(format!(
3134 "column \"{name}\" specified in USING clause does not exist in right table"
3135 ))
3136 })?;
3137 let left_column = &scope.columns[left_at];
3138 let (left_binding, left_type) = (left_column.binding, left_column.ty.clone());
3139 let right_column = &scope.columns[right_at];
3140 let (right_binding, right_type) = (right_column.binding, right_column.ty.clone());
3141 let left_expr = self.plan.add_expr(Expr::Column(left_binding), left_type);
3142 let right_expr = self.plan.add_expr(Expr::Column(right_binding), right_type);
3143 conditions.push(self.compare(rudb_plan::CompareOp::Equal, left_expr, right_expr)?);
3144 dropped.push(right_at);
3145 }
3146 dropped.sort_unstable();
3149 for at in dropped.into_iter().rev() {
3150 scope.remove(at);
3151 }
3152
3153 let mut left_node = left_node;
3154 let mut right_node = right_node;
3155 let mut pair = Vec::new();
3156 if on != NONE {
3157 if !merged.is_empty() {
3158 return Err(Error::binder("a join cannot have both ON and USING"));
3159 }
3160 self.clause = "JOIN condition";
3161 let waiting = self.scalar_subqueries.len();
3162 let predicate = self.bind_expr(ast, on, &scope)?;
3163 conditions.push(self.as_boolean(predicate, "JOIN")?);
3164 for pending in self.scalar_subqueries.split_off(waiting) {
3165 match self.side_of(&pending, &left_tables, &right_tables) {
3166 Some(Side::Right) => right_node = self.attach_subquery(right_node, pending),
3167 Some(Side::Left) => left_node = self.attach_subquery(left_node, pending),
3168 None => pair.push(pending),
3169 }
3170 }
3171 }
3172
3173 if kind == ast::JoinKind::Cross && !conditions.is_empty() {
3174 return Err(Error::binder("a CROSS JOIN cannot have a condition"));
3175 }
3176 if !pair.is_empty() {
3177 return self.bind_pair_dependent_join(
3178 kind,
3179 correlated.is_empty(),
3180 left_node,
3181 right_node,
3182 pair,
3183 conditions,
3184 scope,
3185 );
3186 }
3187 if correlated.is_empty()
3191 && conditions.is_empty()
3192 && matches!(kind, ast::JoinKind::Cross | ast::JoinKind::Inner)
3193 {
3194 let node = self.add_node(Node::CrossProduct { left: left_node, right: right_node });
3195 return Ok((node, scope));
3196 }
3197 if matches!(kind, ast::JoinKind::Semi | ast::JoinKind::Anti) {
3206 scope.truncate(split);
3207 }
3208 let kind = match kind {
3209 ast::JoinKind::Inner | ast::JoinKind::Cross => JoinKind::Inner,
3210 ast::JoinKind::Left => JoinKind::Left,
3211 ast::JoinKind::Right => JoinKind::Right,
3212 ast::JoinKind::Full => JoinKind::Full,
3213 ast::JoinKind::Semi => JoinKind::Semi,
3214 ast::JoinKind::Anti => JoinKind::Anti,
3215 ast::JoinKind::Positional => JoinKind::Positional,
3216 };
3217 let conditions = self.plan.add_expr_list(&conditions);
3218 let node = if correlated.is_empty() {
3219 self.add_node(Node::Join {
3220 left: left_node,
3221 right: right_node,
3222 kind,
3223 conditions,
3224 build: BuildSide::default(),
3225 })
3226 } else {
3227 self.add_node(Node::DependentJoin {
3228 left: left_node,
3229 right: right_node,
3230 kind,
3231 conditions,
3232 })
3233 };
3234 Ok((node, scope))
3235 }
3236
3237 fn bind_filter(
3245 &mut self,
3246 ast: &Ast,
3247 filter: ast::ExprRef,
3248 scope: &Scope,
3249 ) -> Result<Option<ExprRef>> {
3250 if filter == NONE {
3251 return Ok(None);
3252 }
3253 let bound = self.bind_expr(ast, filter, scope)?;
3254 Ok(Some(self.checked_cast_to(bound, &LogicalType::Boolean, false)?))
3255 }
3256
3257 #[allow(clippy::too_many_arguments)]
3262 pub(crate) fn bind_aggregate(
3263 &mut self,
3264 ast: &Ast,
3265 name: &str,
3266 args: &[ast::ExprRef],
3267 distinct: bool,
3268 filter: ast::ExprRef,
3269 sorted: &[ast::OrderItem],
3270 scope: &Scope,
3271 ) -> Result<ExprRef> {
3272 if self.trying {
3273 return Err(Error::binder("aggregates are not allowed inside the TRY expression"));
3274 }
3275 let frames = std::mem::take(&mut self.lambda_frames);
3276 let call = AggregateCall { name, args, distinct, filter, sorted };
3277 let bound = self.bind_aggregate_over_rows(ast, &call, scope);
3278 self.lambda_frames = frames;
3279 bound
3280 }
3281
3282 fn bind_aggregate_over_rows(
3283 &mut self,
3284 ast: &Ast,
3285 written: &AggregateCall<'_>,
3286 scope: &Scope,
3287 ) -> Result<ExprRef> {
3288 let AggregateCall { name, args, distinct, filter, sorted } = *written;
3289 if self.in_filter {
3290 return Err(Error::binder("aggregate functions are not allowed in FILTER"));
3291 }
3292 if self.in_aggregate {
3293 return Err(Error::binder(format!(
3294 "aggregate function calls cannot be nested, and {name}() is inside one"
3295 )));
3296 }
3297 if self.aggregation.is_none() {
3298 let clause = if self.clause == "JOIN condition" { "WHERE clause" } else { self.clause };
3300 return Err(Error::binder(format!("{clause} cannot contain aggregates!")));
3301 }
3302 self.in_aggregate = true;
3307 self.in_filter = true;
3308 let filter = self.bind_filter(ast, filter, scope);
3309 self.in_filter = false;
3310 self.in_aggregate = false;
3311 let filter = filter?;
3312
3313 let (ordered_set, taken) = ordered_set(name, args.len(), sorted);
3317 let injected = sorted.iter().map(|item| item.expr).take(usize::from(taken));
3318 let args: Vec<ast::ExprRef> = injected.chain(args.iter().copied()).collect();
3319 let from_top = ordered_set
3320 && sorted.len() == 1
3321 && match sorted[0].order {
3322 Order::Unstated => self.semantics.default_descending(),
3323 Order::Ascending => false,
3324 Order::Descending => true,
3325 };
3326
3327 self.in_aggregate = true;
3328 let mut bound = Vec::with_capacity(args.len());
3329 let mut failure = None;
3330 let written_keys = sorted.iter().map(|item| item.expr);
3331 for arg in args.iter().copied().chain(written_keys) {
3332 match self.bind_expr(ast, arg, scope) {
3333 Ok(expr) => bound.push(expr),
3334 Err(error) => {
3335 failure = Some(error);
3336 break;
3337 }
3338 }
3339 }
3340 self.in_aggregate = false;
3341 if let Some(error) = failure {
3342 return Err(error);
3343 }
3344 let keys = bound.split_off(args.len());
3345 if distinct && !keys.iter().all(|&key| bound.iter().any(|&arg| self.same_expr(arg, key))) {
3348 return Err(Error::binder(
3349 "In a DISTINCT aggregate, ORDER BY expressions must appear in the argument list",
3350 ));
3351 }
3352
3353 let types: Vec<LogicalType> =
3354 bound.iter().map(|&arg| self.plan.expr_type(arg).clone()).collect();
3355 let resolved = resolve(name, &types)?;
3356 if resolved.name == "string_agg"
3359 && bound.len() == 2
3360 && !matches!(fold::value_of(&self.plan, bound[1]), Ok(Some(_)))
3361 {
3362 return Err(Error::binder(
3363 "The \"separator\" argument in function \"string_agg\" must be a constant expression",
3364 ));
3365 }
3366 if matches!(resolved.name, "quantile_cont" | "quantile_disc") {
3367 let ordered = ordered_set && sorted.len() == 1;
3368 bound[1] = self.quantile_fraction(resolved.name, bound[1], ordered, from_top)?;
3369 }
3370 if resolved.name == "approx_quantile" {
3371 self.digest_arguments(&bound)?;
3372 }
3373 if resolved.name == "reservoir_quantile" {
3374 self.reservoir_arguments(&bound)?;
3375 }
3376 if resolved.name == "approx_top_k" {
3377 self.top_k_argument(&bound)?;
3378 }
3379 let mut cast = Vec::with_capacity(bound.len());
3380 for (arg, wanted) in bound.iter().zip(&resolved.arguments) {
3381 cast.push(self.checked_cast_to(*arg, wanted, false)?);
3382 }
3383 let name = self.ordered_aggregate(resolved.name, sorted, &keys, &mut cast);
3384 let args = self.plan.add_expr_list(&cast);
3385 let name = self.plan.intern(&name);
3386 let ty = resolved.returns;
3387 let call = self.plan.add_expr(Expr::Aggregate { name, args, distinct, filter }, ty.clone());
3388
3389 let existing = self.aggregation.as_ref().map(|held| held.aggregates.clone());
3392 let existing = existing.unwrap_or_default();
3393 let at = match existing.iter().position(|&held| self.same_expr(held, call)) {
3394 Some(at) => at,
3395 None => {
3396 let aggregation = self.aggregation.as_mut().expect("checked above");
3397 aggregation.aggregates.push(call);
3398 aggregation.aggregates.len() - 1
3399 }
3400 };
3401 let aggregation = self.aggregation.as_ref().expect("checked above");
3402 let (index, groups) = (aggregation.index, aggregation.groups.len());
3403 Ok(self.column(index, groups + at, ty))
3404 }
3405
3406 fn quantile_fraction(
3412 &mut self,
3413 name: &str,
3414 fraction: ExprRef,
3415 ordered: bool,
3416 from_top: bool,
3417 ) -> Result<ExprRef> {
3418 let Ok(Some(value)) = fold::value_of(&self.plan, fraction) else {
3419 return Err(Error::binder(format!(
3420 "The \"quantile\" argument in function \"{name}\" must be a constant expression"
3421 )));
3422 };
3423 if value.is_null() {
3424 return Err(Error::binder(format!(
3425 "The \"quantile\" argument in function '\"{name}\"' must not be NULL"
3426 )));
3427 }
3428 let each = match &value {
3429 Value::List { values, .. } => values.as_slice(),
3430 one => std::slice::from_ref(one),
3431 };
3432 let mut signs = (false, false);
3433 for one in each {
3434 if one.is_null() {
3435 return Err(Error::binder("QUANTILE parameter cannot be NULL"));
3436 }
3437 let share = share(one).unwrap_or(f64::NAN);
3438 if !(-1.0..=1.0).contains(&share) {
3439 return Err(Error::binder(
3440 "QUANTILE can only take parameters in the range [-1, 1]",
3441 ));
3442 }
3443 if share < 0.0 {
3444 signs.0 = true;
3445 } else {
3446 signs.1 = true;
3447 }
3448 }
3449 if ordered && signs.0 {
3450 return Err(Error::binder("PERCENTILEs can only take parameters in the range [0, 1]"));
3451 }
3452 if signs.0 && signs.1 {
3453 return Err(Error::binder("QUANTILE parameters must have consistent signs"));
3454 }
3455 if !from_top {
3456 return Ok(fraction);
3457 }
3458 let negated = match value {
3459 Value::List { element, values } => {
3460 Value::List { element, values: values.iter().map(negated).collect() }
3461 }
3462 one => negated(&one),
3463 };
3464 Ok(self.add_constant(negated))
3465 }
3466
3467 fn top_k_argument(&self, bound: &[ExprRef]) -> Result<()> {
3471 if matches!(fold::value_of(&self.plan, bound[1]), Ok(Some(_))) {
3472 return Ok(());
3473 }
3474 Err(Error::binder(
3475 "The \"col1\" argument in function \"approx_top_k\" must be a constant expression",
3476 ))
3477 }
3478
3479 fn reservoir_arguments(&self, bound: &[ExprRef]) -> Result<()> {
3482 let constant = |arg: ExprRef, parameter: &str| match fold::value_of(&self.plan, arg) {
3483 Ok(Some(value)) => Ok(value),
3484 _ => Err(Error::binder(format!(
3485 "The \"{parameter}\" argument in function \"reservoir_quantile\" must be a constant \
3486 expression"
3487 ))),
3488 };
3489 let fraction = constant(bound[1], "quantile")?;
3490 let each = match &fraction {
3491 Value::List { values, .. } => values.as_slice(),
3492 one => std::slice::from_ref(one),
3493 };
3494 for one in each {
3495 if one.is_null() {
3496 return Err(Error::binder("RESERVOIR_QUANTILE QUANTILE parameter cannot be NULL"));
3497 }
3498 if !(0.0..=1.0).contains(&share(one).unwrap_or(f64::NAN)) {
3499 return Err(Error::binder(
3500 "RESERVOIR_QUANTILE can only take parameters in the range [0, 1]",
3501 ));
3502 }
3503 }
3504 let Some(&size) = bound.get(2) else {
3505 return Ok(());
3506 };
3507 let size = constant(size, "sample_size")?;
3508 if size.is_null() {
3509 return Err(Error::binder(
3510 "The \"sample_size\" argument in function '\"reservoir_quantile\"' must not be NULL",
3511 ));
3512 }
3513 if share(&size).is_none_or(|n| n <= 0.0) {
3514 return Err(Error::binder(
3515 "Size of the RESERVOIR_QUANTILE sample must be bigger than 0",
3516 ));
3517 }
3518 Ok(())
3519 }
3520
3521 fn digest_arguments(&self, bound: &[ExprRef]) -> Result<()> {
3524 let Ok(Some(fraction)) = fold::value_of(&self.plan, bound[1]) else {
3525 return Err(Error::binder(
3526 "The \"quantile\" argument in function \"approx_quantile\" must be a constant \
3527 expression",
3528 ));
3529 };
3530 if fraction.is_null() {
3531 return Err(Error::binder(
3532 "The \"quantile\" argument in function '\"approx_quantile\"' must not be NULL",
3533 ));
3534 }
3535 let each = match &fraction {
3536 Value::List { values, .. } => values.as_slice(),
3537 one => std::slice::from_ref(one),
3538 };
3539 for one in each {
3540 if one.is_null() {
3541 return Err(Error::binder("APPROXIMATE QUANTILE parameter cannot be NULL"));
3542 }
3543 if !(0.0..=1.0).contains(&share(one).unwrap_or(f64::NAN)) {
3544 return Err(Error::binder(
3545 "APPROXIMATE QUANTILE can only take parameters in range [0, 1]",
3546 ));
3547 }
3548 }
3549 Ok(())
3550 }
3551
3552 fn ordered_aggregate(
3561 &mut self,
3562 name: &str,
3563 sorted: &[ast::OrderItem],
3564 keys: &[ExprRef],
3565 args: &mut Vec<ExprRef>,
3566 ) -> String {
3567 const DEPENDS_ON_ORDER: &[&str] = &["list", "first", "last", "any_value", "string_agg"];
3568 if !DEPENDS_ON_ORDER.contains(&name) {
3569 return name.to_string();
3570 }
3571 let mut flags = Vec::new();
3572 for (&key, item) in keys.iter().zip(sorted) {
3573 if matches!(fold::value_of(&self.plan, key), Ok(Some(_))) {
3574 continue;
3575 }
3576 let descending = match item.order {
3577 Order::Unstated => self.semantics.default_descending(),
3578 Order::Ascending => false,
3579 Order::Descending => true,
3580 };
3581 let nulls_first = match item.nulls {
3582 Nulls::First => true,
3583 Nulls::Last => false,
3584 Nulls::Unstated => self.semantics.nulls_first(descending),
3585 };
3586 flags.push((descending, nulls_first));
3587 args.push(key);
3588 }
3589 if flags.is_empty() {
3590 return name.to_string();
3591 }
3592 rudb_kernels::ordered_name(name, &flags)
3593 }
3594
3595 pub(crate) fn bind_window(
3606 &mut self,
3607 ast: &Ast,
3608 written: &WindowCall<'_>,
3609 scope: &Scope,
3610 ) -> Result<ExprRef> {
3611 if self.trying {
3612 return Err(Error::binder("window functions are not allowed in try"));
3613 }
3614 let frames = std::mem::take(&mut self.lambda_frames);
3615 let bound = self.bind_window_over_rows(ast, written, scope);
3616 self.lambda_frames = frames;
3617 bound
3618 }
3619
3620 fn bind_window_over_rows(
3621 &mut self,
3622 ast: &Ast,
3623 written: &WindowCall<'_>,
3624 scope: &Scope,
3625 ) -> Result<ExprRef> {
3626 let WindowCall { name, args, distinct, filter, ignore_nulls, spec, .. } = *written;
3627 if self.in_aggregate {
3628 return Err(Error::binder(
3629 "aggregate function calls cannot contain window function calls",
3630 ));
3631 }
3632 if self.in_window {
3633 return Err(Error::binder("window function calls cannot be nested"));
3634 }
3635 let clause = if self.clause == "JOIN condition" { "WHERE clause" } else { self.clause };
3639 if clause != "SELECT clause" && clause != "ORDER BY clause" {
3640 return Err(Error::binder(format!("{clause} cannot contain window functions!")));
3641 }
3642
3643 let starred = args.iter().any(|&arg| {
3647 matches!(ast.expr(arg), ast::Expr::Star { qualifier, replacements }
3648 if qualifier.is_empty() && replacements.is_empty())
3649 });
3650 let (name, args): (&str, &[ast::ExprRef]) = if starred {
3651 if !same_name(name, "count") || args.len() != 1 {
3652 return Err(Error::binder(format!("* is not allowed in {name}()")));
3653 }
3654 ("count_star", &[])
3655 } else if same_name(name, "count") && args.is_empty() {
3656 ("count_star", &[])
3659 } else {
3660 (name, args)
3661 };
3662
3663 let held = ast.window(spec);
3664 self.in_window = true;
3665 let parts = self.window_parts(ast, written, args, held, scope);
3666 let filter = if parts.is_ok() { self.bind_filter(ast, filter, scope) } else { Ok(None) };
3671 self.in_window = false;
3672 let parts = parts?;
3673 let filter = filter?;
3674 let offsets = [parts.frame.start, parts.frame.end]
3677 .iter()
3678 .any(|end| matches!(end, WindowBound::Preceding(_) | WindowBound::Following(_)));
3679 if parts.frame.unit == WindowUnit::Range && offsets && parts.order.len() != 1 {
3680 return Err(Error::binder("RANGE frames must have only one ORDER BY expression"));
3681 }
3682
3683 let types: Vec<LogicalType> =
3684 parts.args.iter().map(|&arg| self.plan.expr_type(arg).clone()).collect();
3685 let resolved = window_signature(name, &types)?;
3686 if resolved.name == "fill" {
3689 let keys: Vec<LogicalType> =
3690 parts.order.iter().map(|key| self.plan.expr_type(key.expr).clone()).collect();
3691 refuse_fill(&types[0], &keys, distinct, ignore_nulls)?;
3692 }
3693 if distinct && kind_of(resolved.name) == Some(FunctionKind::Window) {
3697 return Err(Error::binder(format!(
3698 "DISTINCT is not implemented for the window function \"\"{name}\"\""
3699 )));
3700 }
3701 if filter.is_some() && kind_of(resolved.name) == Some(FunctionKind::Window) {
3704 return Err(Error::binder(format!(
3705 "FILTER is not implemented for the window function \"\"{name}\"\""
3706 )));
3707 }
3708 if !parts.inner.is_empty() && kind_of(resolved.name) == Some(FunctionKind::Window) {
3716 let counts = matches!(resolved.name, "first_value" | "last_value" | "nth_value");
3717 if !counts {
3718 if parts.frame.exclude != WindowExclude::NoOthers {
3719 return Err(Error::binder(format!(
3720 "EXCLUDE is not supported for the window function \"\"{}\"\"",
3721 resolved.name
3722 )));
3723 }
3724 return Err(Error::not_implemented(format!(
3725 "ORDER BY inside the arguments of the window function \"{}\"",
3726 resolved.name
3727 )));
3728 }
3729 }
3730 if resolved.name == "approx_top_k" {
3731 self.top_k_argument(&parts.args)?;
3732 }
3733 let mut cast = Vec::with_capacity(parts.args.len());
3734 for (arg, wanted) in parts.args.iter().zip(&resolved.arguments) {
3735 cast.push(self.checked_cast_to(*arg, wanted, false)?);
3736 }
3737 let args = self.plan.add_expr_list(&cast);
3738 let order = self.plan.add_sort_keys(&parts.inner);
3739 let name = self.plan.intern(resolved.name);
3740 let ty = resolved.returns;
3741 let call = self.plan.add_expr(
3742 Expr::Window { name, args, distinct, filter, ignore_nulls, order },
3743 ty.clone(),
3744 );
3745
3746 let at = self.window_run(parts.partition, parts.order, parts.frame, call);
3747 let index = self.windows.last().expect("the run was just filed").index;
3748 Ok(self.column(index, at, ty))
3749 }
3750
3751 fn window_run(
3758 &mut self,
3759 partition: Vec<ExprRef>,
3760 order: Vec<SortKey>,
3761 frame: WindowFrame,
3762 call: ExprRef,
3763 ) -> usize {
3764 let matches = self.windows.last().is_some_and(|run| {
3765 run.frame == frame
3766 && run.partition.len() == partition.len()
3767 && run.order.len() == order.len()
3768 && run.partition.iter().zip(&partition).all(|(&l, &r)| self.same_expr(l, r))
3769 && run.order.iter().zip(&order).all(|(l, r)| {
3770 l.descending == r.descending
3771 && l.nulls_first == r.nulls_first
3772 && self.same_expr(l.expr, r.expr)
3773 })
3774 });
3775 if !matches {
3776 let index = self.fresh_index();
3777 self.windows.push(WindowRun { index, partition, order, frame, calls: Vec::new() });
3778 }
3779 let calls = self.windows.last().expect("a run is open").calls.clone();
3782 if let Some(at) = calls.iter().position(|&held| self.same_expr(held, call)) {
3783 return at;
3784 }
3785 let run = self.windows.last_mut().expect("a run is open");
3786 run.calls.push(call);
3787 run.calls.len() - 1
3788 }
3789
3790 fn window_parts(
3796 &mut self,
3797 ast: &Ast,
3798 written: &WindowCall<'_>,
3799 args: &[ast::ExprRef],
3800 held: ast::WindowSpec,
3801 scope: &Scope,
3802 ) -> Result<WindowParts> {
3803 let mut bound = Vec::with_capacity(args.len());
3804 for &arg in args {
3805 let expr = self.bind_expr(ast, arg, scope)?;
3806 bound.push(self.over_aggregate(expr, scope)?);
3807 }
3808 let mut inner = Vec::new();
3812 for item in ast.order_list(written.order).to_vec() {
3813 let expr = self.bind_expr(ast, item.expr, scope)?;
3814 let expr = self.over_aggregate(expr, scope)?;
3815 inner.push(self.sort_key(expr, item));
3816 }
3817 let mut partition = Vec::new();
3818 for &key in ast.expr_list(held.partition) {
3819 let expr = self.bind_expr(ast, key, scope)?;
3820 partition.push(self.over_aggregate(expr, scope)?);
3821 }
3822 let mut order = Vec::new();
3823 for item in ast.order_list(held.order).to_vec() {
3824 let expr = self.bind_expr(ast, item.expr, scope)?;
3825 let expr = self.over_aggregate(expr, scope)?;
3826 order.push(self.sort_key(expr, item));
3827 }
3828 let frame = WindowFrame {
3829 unit: match held.unit {
3830 ast::WindowUnit::Rows => WindowUnit::Rows,
3831 ast::WindowUnit::Range => WindowUnit::Range,
3832 ast::WindowUnit::Groups => WindowUnit::Groups,
3833 },
3834 start: self.window_bound(ast, held.start, scope)?,
3835 end: self.window_bound(ast, held.end, scope)?,
3836 exclude: match held.exclude {
3837 ast::WindowExclude::NoOthers => WindowExclude::NoOthers,
3838 ast::WindowExclude::CurrentRow => WindowExclude::CurrentRow,
3839 ast::WindowExclude::Group => WindowExclude::Group,
3840 ast::WindowExclude::Ties => WindowExclude::Ties,
3841 },
3842 };
3843 Ok(WindowParts { args: bound, partition, order, inner, frame })
3844 }
3845
3846 fn window_bound(
3848 &mut self,
3849 ast: &Ast,
3850 bound: ast::WindowBound,
3851 scope: &Scope,
3852 ) -> Result<WindowBound> {
3853 let offset = |binder: &mut Self, written| {
3854 let expr = binder.bind_expr(ast, written, scope)?;
3855 binder.over_aggregate(expr, scope)
3856 };
3857 Ok(match bound {
3858 ast::WindowBound::UnboundedPreceding => WindowBound::UnboundedPreceding,
3859 ast::WindowBound::CurrentRow => WindowBound::CurrentRow,
3860 ast::WindowBound::UnboundedFollowing => WindowBound::UnboundedFollowing,
3861 ast::WindowBound::Preceding(written) => WindowBound::Preceding(offset(self, written)?),
3862 ast::WindowBound::Following(written) => WindowBound::Following(offset(self, written)?),
3863 })
3864 }
3865
3866 fn group_of(&self, read: ColumnBinding) -> Option<usize> {
3872 self.aggregation.as_ref()?.groups.iter().position(
3873 |group| matches!(*self.plan.expr(*group), Expr::Column(binding) if binding == read),
3874 )
3875 }
3876
3877 fn ungrouped_correlation(&self, binding: ColumnBinding) -> Option<ColumnBinding> {
3883 let pending =
3884 self.scalar_subqueries.iter().find(|pending| pending.index == binding.table)?;
3885 pending.reads.iter().copied().find(|read| self.group_of(*read).is_none())
3886 }
3887
3888 fn is_window_output(&self, binding: ColumnBinding) -> bool {
3890 self.windows.iter().any(|run| run.index == binding.table)
3891 }
3892
3893 fn is_correlation(&self, binding: ColumnBinding) -> bool {
3899 self.correlations.last().is_some_and(|frame| frame.contains(&binding))
3900 }
3901
3902 fn name_of(&self, binding: ColumnBinding, scope: &Scope) -> String {
3908 std::iter::once(scope)
3909 .chain(self.outer_scopes.iter().rev())
3910 .flat_map(|visible| visible.columns.iter())
3911 .find(|column| column.binding == binding)
3912 .map_or_else(|| "a column".to_string(), |column| format!("\"{}\"", column.name))
3913 }
3914
3915 pub(crate) fn over_aggregate(&mut self, expr: ExprRef, scope: &Scope) -> Result<ExprRef> {
3921 let Some(aggregation) = self.aggregation.as_ref() else {
3922 return Ok(expr);
3923 };
3924 let index = aggregation.index;
3925 let groups = aggregation.groups.clone();
3926 for (at, group) in groups.iter().enumerate() {
3927 if self.same_expr(expr, *group) {
3928 let ty = self.plan.expr_type(*group).clone();
3929 return Ok(self.column(index, at, ty));
3930 }
3931 }
3932 let ty = self.plan.expr_type(expr).clone();
3933 match self.plan.expr(expr).clone() {
3934 Expr::Column(binding) if binding.table == index => Ok(expr),
3935 Expr::Column(binding) if self.is_window_output(binding) => Ok(expr),
3940 Expr::Column(binding) if self.is_unnest_output(binding) => Ok(expr),
3943 Expr::Column(binding) if self.joined_above.contains(&binding.table) => Ok(expr),
3948 Expr::Column(binding) if self.is_correlation(binding) => Ok(expr),
3954 Expr::Column(binding) => {
3963 let read = self.ungrouped_correlation(binding).unwrap_or(binding);
3964 let name = self.name_of(read, scope);
3965 Err(Error::binder(format!(
3966 "column {name} must appear in the GROUP BY clause or must be part of an aggregate function"
3967 )))
3968 }
3969 Expr::Constant(_)
3970 | Expr::Aggregate { .. }
3971 | Expr::Window { .. }
3972 | Expr::LambdaParam(_) => Ok(expr),
3973 Expr::Lambda { table, params, body } => {
3977 let body = self.over_aggregate(body, scope)?;
3978 Ok(self.plan.add_expr(Expr::Lambda { table, params, body }, ty))
3979 }
3980 Expr::Cast { input, try_cast } => {
3981 let input = self.over_aggregate(input, scope)?;
3982 Ok(self.plan.add_expr(Expr::Cast { input, try_cast }, ty))
3983 }
3984 Expr::Compare { op, left, right } => {
3985 let left = self.over_aggregate(left, scope)?;
3986 let right = self.over_aggregate(right, scope)?;
3987 Ok(self.plan.add_expr(Expr::Compare { op, left, right }, ty))
3988 }
3989 Expr::Conjunction { op, children } => {
3990 let written = self.plan.expr_list(children).to_vec();
3991 let mut rewritten = Vec::with_capacity(written.len());
3992 for child in written {
3993 rewritten.push(self.over_aggregate(child, scope)?);
3994 }
3995 let children = self.plan.add_expr_list(&rewritten);
3996 Ok(self.plan.add_expr(Expr::Conjunction { op, children }, ty))
3997 }
3998 Expr::Function { name, args } => {
3999 let written = self.plan.expr_list(args).to_vec();
4000 let mut rewritten = Vec::with_capacity(written.len());
4001 for arg in written {
4002 rewritten.push(self.over_aggregate(arg, scope)?);
4003 }
4004 let args = self.plan.add_expr_list(&rewritten);
4005 Ok(self.plan.add_expr(Expr::Function { name, args }, ty))
4006 }
4007 Expr::Case { arms, otherwise } => {
4008 let written = self.plan.arm_list(arms).to_vec();
4009 let mut rewritten = Vec::with_capacity(written.len());
4010 for arm in written {
4011 let when = self.over_aggregate(arm.when, scope)?;
4012 let then = self.over_aggregate(arm.then, scope)?;
4013 rewritten.push(rudb_plan::Arm { when, then });
4014 }
4015 let otherwise = match otherwise {
4016 Some(expr) => Some(self.over_aggregate(expr, scope)?),
4017 None => None,
4018 };
4019 let arms = self.plan.add_arms(&rewritten);
4020 Ok(self.plan.add_expr(Expr::Case { arms, otherwise }, ty))
4021 }
4022 }
4023 }
4024
4025 pub(crate) fn same_expr(&self, left: ExprRef, right: ExprRef) -> bool {
4027 same_expr(&self.plan, left, right)
4028 }
4029}
4030
4031#[derive(Debug, Default)]
4041struct Options {
4042 binary_as_string: bool,
4045 all_varchar: bool,
4047 file_row_number: bool,
4052 given: Given,
4054}
4055
4056impl Options {
4057 fn of(written: &[(&'static str, Value, ExprRef)]) -> Result<Self> {
4064 let mut options = Self::default();
4065 for (parameter, value, _) in written {
4066 match (*parameter, value) {
4067 ("binary_as_string", Value::Boolean(on)) => options.binary_as_string = *on,
4068 ("all_varchar", Value::Boolean(on)) => options.all_varchar = *on,
4069 ("file_row_number", Value::Boolean(on)) => options.file_row_number = *on,
4070 _ => {}
4071 }
4072 }
4073 let named: Vec<(&str, Value)> =
4074 written.iter().map(|(parameter, value, _)| (*parameter, value.clone())).collect();
4075 options.given = csv_given(&named)?;
4076 Ok(options)
4077 }
4078}
4079
4080fn mirror_target(paths: &[String]) -> Option<(String, FileStamp)> {
4083 let [path] = paths else { return None };
4084 let canonical = std::fs::canonicalize(path).ok()?;
4085 let stamp = FileStamp::of(&canonical)?;
4086 Some((canonical.to_str()?.to_string(), stamp))
4087}
4088
4089#[derive(Clone, Copy)]
4091struct Operator {
4092 op: SetOp,
4094 quantifier: Quantifier,
4096 by_name: bool,
4098}
4099
4100struct Merged {
4102 name: String,
4104 ty: LogicalType,
4106 left: Option<usize>,
4108 right: Option<usize>,
4110}
4111
4112fn match_by_position(left: &Scope, right: &Scope) -> Result<Vec<Merged>> {
4116 if left.len() != right.len() {
4117 return Err(Error::binder(format!(
4118 "Set operations can only apply to expressions with the same number of result columns, but left side has {} and right side has {}",
4119 left.len(),
4120 right.len()
4121 )));
4122 }
4123 let mut merged = Vec::with_capacity(left.len());
4124 for (at, (held, other)) in left.columns.iter().zip(&right.columns).enumerate() {
4125 merged.push(Merged {
4126 name: held.name.clone(),
4127 ty: meet(&held.ty, &other.ty)?,
4128 left: Some(at),
4129 right: Some(at),
4130 });
4131 }
4132 Ok(merged)
4133}
4134
4135fn match_by_name(left: &Scope, right: &Scope) -> Result<Vec<Merged>> {
4143 named_once(left)?;
4144 named_once(right)?;
4145 let mut merged = Vec::with_capacity(left.len() + right.len());
4146 for (at, held) in left.columns.iter().enumerate() {
4147 let other = right.columns.iter().position(|column| same_name(&column.name, &held.name));
4148 let ty = match other {
4149 Some(other) => meet(&held.ty, &right.columns[other].ty)?,
4150 None => held.ty.clone(),
4151 };
4152 merged.push(Merged { name: held.name.clone(), ty, left: Some(at), right: other });
4153 }
4154 for (at, held) in right.columns.iter().enumerate() {
4155 if left.columns.iter().any(|column| same_name(&column.name, &held.name)) {
4156 continue;
4157 }
4158 merged.push(Merged {
4159 name: held.name.clone(),
4160 ty: held.ty.clone(),
4161 left: None,
4162 right: Some(at),
4163 });
4164 }
4165 Ok(merged)
4166}
4167
4168fn named_once(scope: &Scope) -> Result<()> {
4174 for (at, held) in scope.columns.iter().enumerate() {
4175 if scope.columns[..at].iter().any(|column| same_name(&column.name, &held.name)) {
4176 return Err(Error::binder(format!(
4177 "UNION (ALL) BY NAME operation doesn't support duplicate names in the SELECT list - the name \"\"{}\"\" occurs multiple times",
4178 held.name
4179 )));
4180 }
4181 }
4182 Ok(())
4183}
4184
4185fn meet(left: &LogicalType, right: &LogicalType) -> Result<LogicalType> {
4187 left.promote(right).ok_or_else(|| {
4188 Error::binder(format!(
4189 "Cannot combine a column of type {left} with a column of type {right} in a set operation"
4190 ))
4191 })
4192}
4193
4194fn null_parameter(function: TableFunction, parameter: &str) -> String {
4203 match parameter {
4204 "header" => format!("\"{parameter}\" expects a non-null boolean value (e.g. TRUE or 1)"),
4205 "all_varchar" => format!("{} \"{parameter}\" cannot be NULL", function.name()),
4206 _ => format!("Cannot use NULL as argument to \"{parameter}\""),
4207 }
4208}
4209
4210fn missing_replacement(name: &str, input: &Scope) -> Error {
4215 Error::binder(format!(
4216 "Column \"{name}\" in REPLACE list not found in FROM clause{}",
4217 input.candidates()
4218 ))
4219}
4220
4221fn subtractable(ty: &LogicalType, ordering: bool) -> bool {
4230 if ty.is_numeric() {
4231 return true;
4232 }
4233 match ty {
4234 LogicalType::Date
4235 | LogicalType::Time
4236 | LogicalType::Timestamp
4237 | LogicalType::TimestampS
4238 | LogicalType::TimestampMs
4239 | LogicalType::TimestampNs
4240 | LogicalType::TimestampTz => true,
4241 LogicalType::TimeTz => ordering,
4242 _ => false,
4243 }
4244}
4245
4246fn refuse_fill(
4255 argument: &LogicalType,
4256 order: &[LogicalType],
4257 distinct: bool,
4258 ignore_nulls: bool,
4259) -> Result<()> {
4260 if !subtractable(argument, false) {
4261 return Err(Error::binder("FILL argument must support subtraction"));
4262 }
4263 let [key] = order else {
4264 return Err(Error::binder("FILL functions must have only one ORDER BY expression"));
4265 };
4266 if !subtractable(key, true) {
4267 return Err(Error::binder("FILL ordering must support subtraction"));
4268 }
4269 if distinct {
4270 return Err(Error::binder(
4271 "DISTINCT is not implemented for the window function \"\"fill\"\"",
4272 ));
4273 }
4274 if ignore_nulls {
4275 return Err(Error::binder(
4276 "RESPECT/IGNORE NULLS is not supported for the window function \"fill\"",
4277 ));
4278 }
4279 Ok(())
4280}
4281
4282fn window_signature(name: &str, types: &[LogicalType]) -> Result<Resolved> {
4289 match kind_of(name) {
4290 Some(FunctionKind::Aggregate | FunctionKind::Window) => resolve(name, types),
4291 Some(FunctionKind::Scalar) => {
4292 Err(Error::catalog(format!("{name} is not an aggregate function")))
4293 }
4294 None => Err(Error::catalog(format!("Aggregate Function with name {name} does not exist!"))),
4295 }
4296}
4297
4298fn same_expr(plan: &Plan, left: ExprRef, right: ExprRef) -> bool {
4300 if left == right {
4301 return true;
4302 }
4303 if plan.expr_type(left) != plan.expr_type(right) {
4304 return false;
4305 }
4306 let lists = |left, right| {
4307 let left: &[ExprRef] = plan.expr_list(left);
4308 let right: &[ExprRef] = plan.expr_list(right);
4309 left.len() == right.len()
4310 && left.iter().zip(right).all(|(&left, &right)| same_expr(plan, left, right))
4311 };
4312 match (plan.expr(left), plan.expr(right)) {
4313 (Expr::Column(left), Expr::Column(right)) => left == right,
4314 (Expr::Constant(left), Expr::Constant(right)) => plan.value(*left) == plan.value(*right),
4315 (
4316 Expr::Cast { input: left, try_cast: left_try },
4317 Expr::Cast { input: right, try_cast: right_try },
4318 ) => left_try == right_try && same_expr(plan, *left, *right),
4319 (
4320 Expr::Compare { op: left_op, left: left_a, right: left_b },
4321 Expr::Compare { op: right_op, left: right_a, right: right_b },
4322 ) => {
4323 left_op == right_op
4324 && same_expr(plan, *left_a, *right_a)
4325 && same_expr(plan, *left_b, *right_b)
4326 }
4327 (
4328 Expr::Conjunction { op: left_op, children: left_children },
4329 Expr::Conjunction { op: right_op, children: right_children },
4330 ) => left_op == right_op && lists(*left_children, *right_children),
4331 (
4332 Expr::Function { name: left_name, args: left_args },
4333 Expr::Function { name: right_name, args: right_args },
4334 ) => plan.string(*left_name) == plan.string(*right_name) && lists(*left_args, *right_args),
4335 (
4336 Expr::Aggregate {
4337 name: left_name,
4338 args: left_args,
4339 distinct: left_distinct,
4340 filter: left_filter,
4341 },
4342 Expr::Aggregate {
4343 name: right_name,
4344 args: right_args,
4345 distinct: right_distinct,
4346 filter: right_filter,
4347 },
4348 ) => {
4349 plan.string(*left_name) == plan.string(*right_name)
4350 && left_distinct == right_distinct
4351 && match (left_filter, right_filter) {
4352 (None, None) => true,
4353 (Some(left), Some(right)) => same_expr(plan, *left, *right),
4354 _ => false,
4355 }
4356 && lists(*left_args, *right_args)
4357 }
4358 (
4362 Expr::Window {
4363 name: left_name,
4364 args: left_args,
4365 distinct: left_distinct,
4366 filter: left_filter,
4367 ignore_nulls: left_nulls,
4368 order: left_order,
4369 },
4370 Expr::Window {
4371 name: right_name,
4372 args: right_args,
4373 distinct: right_distinct,
4374 filter: right_filter,
4375 ignore_nulls: right_nulls,
4376 order: right_order,
4377 },
4378 ) => {
4379 let left_keys = plan.sort_key_list(*left_order);
4382 let right_keys = plan.sort_key_list(*right_order);
4383 plan.string(*left_name) == plan.string(*right_name)
4384 && left_distinct == right_distinct
4385 && left_nulls == right_nulls
4386 && left_keys.len() == right_keys.len()
4387 && left_keys.iter().zip(right_keys).all(|(left, right)| {
4388 left.descending == right.descending
4389 && left.nulls_first == right.nulls_first
4390 && same_expr(plan, left.expr, right.expr)
4391 })
4392 && match (left_filter, right_filter) {
4393 (None, None) => true,
4394 (Some(left), Some(right)) => same_expr(plan, *left, *right),
4395 _ => false,
4396 }
4397 && lists(*left_args, *right_args)
4398 }
4399 (
4400 Expr::Case { arms: left_arms, otherwise: left_otherwise },
4401 Expr::Case { arms: right_arms, otherwise: right_otherwise },
4402 ) => {
4403 let left_arms = plan.arm_list(*left_arms);
4404 let right_arms = plan.arm_list(*right_arms);
4405 left_arms.len() == right_arms.len()
4406 && left_arms.iter().zip(right_arms).all(|(left, right)| {
4407 same_expr(plan, left.when, right.when) && same_expr(plan, left.then, right.then)
4408 })
4409 && match (left_otherwise, right_otherwise) {
4410 (None, None) => true,
4411 (Some(left), Some(right)) => same_expr(plan, *left, *right),
4412 _ => false,
4413 }
4414 }
4415 _ => false,
4416 }
4417}
4418
4419fn ordered_set(name: &str, written: usize, sorted: &[ast::OrderItem]) -> (bool, bool) {
4422 let name = name.to_ascii_lowercase();
4423 let wants = match name.as_str() {
4424 "quantile_cont" | "quantile_disc" | "quantile" => 1,
4425 "mode" => 0,
4426 _ => return (false, false),
4427 };
4428 (true, sorted.len() == 1 && written == wants)
4429}
4430
4431fn negated(value: &Value) -> Value {
4433 match *value {
4434 Value::Decimal { unscaled, width, scale } => {
4435 Value::Decimal { unscaled: -unscaled, width, scale }
4436 }
4437 Value::Double(share) => Value::Double(-share),
4438 Value::Float(share) => Value::Float(-share),
4439 ref whole => match share(whole) {
4440 Some(share) => Value::Double(-share),
4441 None => whole.clone(),
4442 },
4443 }
4444}
4445
4446#[expect(
4448 clippy::cast_precision_loss,
4449 reason = "a fraction is compared with -1 and 1, which a double holds exactly"
4450)]
4451fn share(value: &Value) -> Option<f64> {
4452 Some(match *value {
4453 Value::TinyInt(v) => f64::from(v),
4454 Value::SmallInt(v) => f64::from(v),
4455 Value::Integer(v) => f64::from(v),
4456 Value::BigInt(v) => v as f64,
4457 Value::HugeInt(v) => v as f64,
4458 Value::UTinyInt(v) => f64::from(v),
4459 Value::USmallInt(v) => f64::from(v),
4460 Value::UInteger(v) => f64::from(v),
4461 Value::UBigInt(v) => v as f64,
4462 Value::UHugeInt(v) => v as f64,
4463 Value::Float(v) => f64::from(v),
4464 Value::Double(v) => v,
4465 Value::Decimal { unscaled, scale, .. } => unscaled as f64 / 10f64.powi(i32::from(scale)),
4466 _ => return None,
4467 })
4468}