1use std::sync::Arc;
16
17use rudb_catalog::{Catalog, Entry, FileStamp, QualifiedName, same_name};
18use rudb_common::bounds::Zones;
19use rudb_common::{
20 Error, Field, LogicalType, Result, Semantics, Session, ShowBehavior, Span, Stat, Value,
21};
22use rudb_functions::{
23 Columns, FILE_ROW_NUMBER, Footers, FunctionKind, Given, Resolved, TableFunction, csv_fields,
24 csv_given, files, is_file, is_pattern, kind_of, parquet_footers, parquet_outline, resolve,
25 resolve_pragma, resolve_table,
26};
27use rudb_kernels::{percentage, row_count};
28use rudb_parse::ast::{self, Ast, Distinct, LiteralKind, Nulls, Order, Quantifier, SetOp};
29use rudb_parse::{NONE, identifier_parts, parse_ast_with_case};
30use rudb_plan::{
31 Bound, BuildSide, ColumnBinding, ConjunctionOp, Expr, ExprRef, JoinKind, Node, NodeRef, Plan,
32 SetOpKind, Share, SortKey, WindowBound, WindowExclude, WindowFrame, WindowUnit,
33};
34
35use crate::expr::{describe, has_aggregate};
36use crate::fold;
37use crate::parameters::Parameters;
38use crate::scope::{Scope, Visible};
39
40pub fn bind(ast: &Ast, catalog: &Catalog) -> Result<Plan> {
47 bind_with(ast, catalog, &Parameters::new(), &Session::new())
48}
49
50pub fn bind_with(
59 ast: &Ast,
60 catalog: &Catalog,
61 parameters: &Parameters,
62 session: &Session,
63) -> Result<Plan> {
64 let query = match ast.statements.as_slice() {
65 [ast::Statement::Query(query)] => *query,
66 [] => return Err(Error::binder("no statement to bind")),
67 [_] => return Err(Error::not_implemented("a statement that is not a query")),
70 _ => return Err(Error::not_implemented("a script of more than one statement")),
71 };
72 let mut binder = Binder::with(catalog, parameters, session);
73 let (root, _) = binder.bind_query(ast, query)?;
74 let mut plan = binder.into_plan();
75 plan.set_root(root);
76 plan.validate()?;
77 Ok(plan)
78}
79
80pub fn bind_sql(query: &str, catalog: &Catalog) -> Result<Plan> {
86 bind_sql_with(query, catalog, &Session::new())
87}
88
89pub fn bind_sql_with(query: &str, catalog: &Catalog, session: &Session) -> Result<Plan> {
95 let ast = parse_ast_with_case(query, session.semantics().identifier_case())?;
96 bind_with(&ast, catalog, &Parameters::new(), session)
97}
98
99#[derive(Debug)]
101pub(crate) struct Aggregation {
102 pub(crate) index: u32,
104 pub(crate) groups: Vec<ExprRef>,
106 pub(crate) aggregates: Vec<ExprRef>,
108}
109
110#[derive(Debug)]
118pub(crate) struct WindowRun {
119 index: u32,
121 partition: Vec<ExprRef>,
123 order: Vec<SortKey>,
125 frame: WindowFrame,
127 calls: Vec<ExprRef>,
129}
130
131pub(crate) struct WindowCall<'a> {
136 pub(crate) name: &'a str,
138 pub(crate) args: &'a [ast::ExprRef],
140 pub(crate) distinct: bool,
142 pub(crate) filter: ast::ExprRef,
144 pub(crate) ignore_nulls: bool,
146 pub(crate) order: ast::Slice,
149 pub(crate) spec: ast::WindowRef,
151}
152
153struct WindowParts {
155 args: Vec<ExprRef>,
157 partition: Vec<ExprRef>,
159 order: Vec<SortKey>,
161 inner: Vec<SortKey>,
164 frame: WindowFrame,
166}
167
168#[derive(Debug)]
175struct Read {
176 fields: Vec<Field>,
178 rows: Stat<u64>,
180 distincts: Vec<(String, Stat<u64>)>,
182 zones: Option<Arc<dyn Zones>>,
184}
185
186impl Read {
187 fn uncounted(fields: Vec<Field>) -> Self {
189 Self { fields, rows: Stat::Unknown, distincts: Vec::new(), zones: None }
190 }
191}
192
193#[derive(Debug)]
195struct Materialized {
196 written: u32,
198 cte: u32,
200 name: String,
202 fields: Vec<Field>,
204}
205
206#[derive(Debug)]
207pub(crate) struct PendingSubquery {
208 pub(crate) node: NodeRef,
209 pub(crate) kind: JoinKind,
210 pub(crate) conditions: Vec<ExprRef>,
211 pub(crate) dependent: bool,
212 pub(crate) reads: Vec<ColumnBinding>,
218 pub(crate) index: u32,
224 pub(crate) inside_aggregate: bool,
231}
232
233#[derive(Debug, Clone, Copy, PartialEq, Eq)]
235enum Side {
236 Left,
237 Right,
238}
239
240#[derive(Debug)]
242pub(crate) struct Binder<'a> {
243 catalog: &'a Catalog,
244 pub(crate) parameters: &'a Parameters,
246 pub(crate) session: &'a Session,
248 pub(crate) semantics: Semantics,
250 plan: Plan,
251 next_index: u32,
252 pub(crate) current_span: Span,
254 pub(crate) aggregation: Option<Aggregation>,
256 want_ascending: bool,
259 pub(crate) upsert: bool,
262 pub(crate) insert_defaults: Option<Vec<(LogicalType, Option<String>)>>,
265 pub(crate) default_as_null: bool,
268 pub(crate) in_aggregate: bool,
270 pub(crate) in_filter: bool,
272 pub(crate) windows: Vec<WindowRun>,
274 pub(crate) in_window: bool,
276 pub(crate) scalar_subqueries: Vec<PendingSubquery>,
278 pub(crate) joined_above: Vec<u32>,
284 pub(crate) outer_scopes: Vec<Scope>,
285 pub(crate) lateral_scopes: Vec<usize>,
292 pub(crate) correlations: Vec<Vec<ColumnBinding>>,
293 pub(crate) lambda_frames: Vec<crate::lambda::Frame>,
295 pub(crate) clause: &'static str,
297 pub(crate) outlined: bool,
306 expanding: Vec<String>,
308 materialized: Vec<Materialized>,
314 next_cte: u32,
316 started: Option<i64>,
318}
319
320impl<'a> Binder<'a> {
321 pub(crate) fn with(
322 catalog: &'a Catalog,
323 parameters: &'a Parameters,
324 session: &'a Session,
325 ) -> Self {
326 Self {
327 catalog,
328 parameters,
329 session,
330 semantics: session.semantics(),
331 plan: Plan::new(),
332 next_index: 0,
333 current_span: Span::new(0, 0),
334 aggregation: None,
335 want_ascending: false,
336 upsert: false,
337 insert_defaults: None,
338 default_as_null: false,
339 in_aggregate: false,
340 in_filter: false,
341 windows: Vec::new(),
342 in_window: false,
343 scalar_subqueries: Vec::new(),
344 joined_above: Vec::new(),
345 outer_scopes: Vec::new(),
346 lateral_scopes: Vec::new(),
347 correlations: Vec::new(),
348 lambda_frames: Vec::new(),
349 clause: "SELECT clause",
350 outlined: false,
351 expanding: Vec::new(),
352 materialized: Vec::new(),
353 next_cte: 0,
354 started: None,
355 }
356 }
357
358 pub(crate) fn catalog(&self) -> &Catalog {
359 self.catalog
360 }
361
362 pub(crate) fn instant(&mut self) -> i64 {
369 *self.started.get_or_insert_with(crate::context::micros_now)
370 }
371
372 pub(crate) fn plan(&self) -> &Plan {
373 &self.plan
374 }
375
376 pub(crate) fn plan_mut(&mut self) -> &mut Plan {
377 &mut self.plan
378 }
379
380 pub(crate) fn add_expr(&mut self, expr: Expr, ty: LogicalType) -> ExprRef {
381 self.plan.add_expr_at(expr, ty, self.current_span)
382 }
383
384 pub(crate) fn add_constant(&mut self, value: Value) -> ExprRef {
385 let ty = value.logical_type();
386 let reference = self.plan.add_value(value);
387 self.plan.add_expr_at(Expr::Constant(reference), ty, self.current_span)
388 }
389
390 pub(crate) fn add_node(&mut self, node: Node) -> NodeRef {
391 self.plan.add_node_at(node, self.current_span)
392 }
393
394 pub(crate) fn into_plan(self) -> Plan {
395 self.plan
396 }
397
398 pub(crate) fn fresh_index(&mut self) -> u32 {
400 let index = self.next_index;
401 self.next_index += 1;
402 index
403 }
404
405 fn column(&mut self, index: u32, position: usize, ty: LogicalType) -> ExprRef {
407 let binding = ColumnBinding::new(index, position as u32);
408 self.plan.add_expr(Expr::Column(binding), ty)
409 }
410
411 fn attach_scalar_subqueries(&mut self, mut input: NodeRef) -> NodeRef {
413 let subqueries = std::mem::take(&mut self.scalar_subqueries);
414 for pending in subqueries {
415 input = self.attach_subquery(input, pending);
416 }
417 input
418 }
419
420 fn attach_subquery(&mut self, input: NodeRef, pending: PendingSubquery) -> NodeRef {
427 let PendingSubquery {
428 node: mut right,
429 kind,
430 conditions,
431 dependent,
432 reads: _,
433 index: _,
434 inside_aggregate: _,
435 } = pending;
436 if kind == JoinKind::Single && !self.semantics.scalar_subquery_error_on_multiple_rows() {
437 right = self.add_node(Node::Limit {
438 input: right,
439 count: Bound::Rows(1),
440 offset: Bound::Rows(0),
441 });
442 }
443 let conditions = self.plan.add_expr_list(&conditions);
444 if dependent {
445 self.add_node(Node::DependentJoin { left: input, right, kind, conditions })
446 } else {
447 self.add_node(Node::Join {
448 left: input,
449 right,
450 kind,
451 conditions,
452 build: BuildSide::default(),
453 })
454 }
455 }
456
457 pub(crate) fn bind_query(
460 &mut self,
461 ast: &Ast,
462 query: ast::QueryRef,
463 ) -> Result<(NodeRef, Scope)> {
464 let span = ast.query_span(query);
465 let outer = std::mem::replace(&mut self.current_span, span);
466 let result =
467 self.bind_query_inner(ast, query).map_err(|error| error.with_fallback_span(span));
468 self.current_span = outer;
469 result
470 }
471
472 fn bind_query_inner(&mut self, ast: &Ast, query: ast::QueryRef) -> Result<(NodeRef, Scope)> {
473 let written = ast.query(query);
474 if written.ctes.is_empty() {
475 return self.bind_body(ast, &written);
476 }
477 let depth = self.materialized.len();
481 let result = self.bind_materialized(ast, &written);
482 self.materialized.truncate(depth);
483 result
484 }
485
486 fn bind_materialized(&mut self, ast: &Ast, written: &ast::Query) -> Result<(NodeRef, Scope)> {
492 let depth = self.materialized.len();
493 let held = ast.cte_list(written.ctes).to_vec();
494 let mut definitions = Vec::with_capacity(held.len());
495 for &index in &held {
496 definitions.push(self.bind_definition(ast, index)?);
497 }
498 let (mut node, scope) = self.bind_body(ast, written)?;
499 for (at, definition) in definitions.into_iter().enumerate().rev() {
500 let entry = &self.materialized[depth + at];
501 let cte = entry.cte;
502 let name = entry.name.clone();
503 let fields = entry.fields.clone();
504 let name = self.plan.intern(&name);
505 let columns = self.plan.add_fields(&fields);
506 node =
507 self.add_node(Node::MaterializedCte { definition, body: node, name, cte, columns });
508 }
509 Ok((node, scope))
510 }
511
512 fn bind_definition(&mut self, ast: &Ast, index: u32) -> Result<NodeRef> {
522 let held = ast.cte(index);
523 let name = ast.string(held.name).to_string();
524 let (node, mut scope) = self.bind_query(ast, held.query)?;
525 if !held.columns.is_empty() {
526 let names: Vec<&str> = ast.name(held.columns).collect();
527 scope.rename_prefix(&names);
528 }
529 let table = self.fresh_index();
530 let mut exprs = Vec::with_capacity(scope.len());
531 let mut names = Vec::with_capacity(scope.len());
532 for column in &scope.columns {
533 exprs.push(self.plan.add_expr(Expr::Column(column.binding), column.ty.clone()));
534 names.push(self.plan.intern(&column.name));
535 }
536 let exprs = self.plan.add_expr_list(&exprs);
537 let names = self.plan.add_name_list(&names);
538 let node = self.add_node(Node::Project { input: node, index: table, exprs, names });
539 let cte = self.next_cte;
540 self.next_cte += 1;
541 self.materialized.push(Materialized { written: index, cte, name, fields: scope.fields() });
542 Ok(node)
543 }
544
545 fn bind_body(&mut self, ast: &Ast, written: &ast::Query) -> Result<(NodeRef, Scope)> {
546 match written.body {
547 ast::QueryBody::Select(select) => self.bind_select(ast, select, written),
548 ast::QueryBody::SetOp { op, quantifier, by_name, left, right } => {
549 let operator = Operator { op, quantifier, by_name };
550 self.bind_set_op(ast, written, operator, left, right)
551 }
552 ast::QueryBody::Values(rows) => self.bind_values(ast, written, rows),
553 ast::QueryBody::Describe(inner) => self.bind_describe(ast, written, inner),
554 ast::QueryBody::Show { name, relation } => self.bind_show(ast, written, name, relation),
555 }
556 }
557
558 fn bind_show(
560 &mut self,
561 ast: &Ast,
562 query: &ast::Query,
563 name: ast::Slice,
564 relation: ast::QueryRef,
565 ) -> Result<(NodeRef, Scope)> {
566 let text = ast.name_text(name);
567 let parts: Vec<&str> = ast.name(name).collect();
568 let table_exists = self.catalog.resolve(&parts).is_ok();
569 let as_table = match self.semantics.show_behavior() {
570 ShowBehavior::Auto => table_exists,
571 ShowBehavior::Setting => false,
572 ShowBehavior::Table => true,
573 };
574 if as_table {
575 return self.bind_describe(ast, query, relation);
576 }
577 let shown = match self.session.iter().find(|(name, _)| name.eq_ignore_ascii_case(&text)) {
582 Some((_, value)) => value.to_string(),
583 None => match self.beyond(&text)? {
584 Some(Value::Varchar(declared)) => declared,
585 Some(other) => other.to_string(),
586 None => {
587 return Err(Error::catalog(format!(
588 "Setting with name \"{text}\" does not exist"
589 )));
590 }
591 },
592 };
593 let field = Field::new(text, LogicalType::Varchar);
594 let expr = self.plan.add_constant(Value::Varchar(shown));
595 let row = self.plan.add_expr_list(&[expr]);
596 let rows = self.plan.add_rows(&[row]);
597 let columns = self.plan.add_fields(std::slice::from_ref(&field));
598 let index = self.fresh_index();
599 let node = self.add_node(Node::Values { index, columns, rows });
600 let mut scope = Scope::empty();
601 scope.push(Visible {
602 table: String::new(),
603 name: field.name,
604 binding: ColumnBinding::new(index, 0),
605 ty: LogicalType::Varchar,
606 not_null: false,
607 key: None,
608 default: None,
609 qualified: false,
610 also: None,
611 });
612 Ok((node, scope))
613 }
614
615 fn bind_describe(
630 &mut self,
631 ast: &Ast,
632 query: &ast::Query,
633 inner: ast::QueryRef,
634 ) -> Result<(NodeRef, Scope)> {
635 let (_, described) = self.bind_query(ast, inner)?;
636 let fields: Vec<Field> = ["column_name", "column_type", "null", "key", "default", "extra"]
637 .iter()
638 .map(|name| Field::new(*name, LogicalType::Varchar))
639 .collect();
640 let mut slices = Vec::with_capacity(described.columns.len());
641 for column in described.columns.clone() {
642 let written = [
645 column.name.clone(),
646 column.ty.to_string(),
647 if column.not_null { "NO" } else { "YES" }.to_owned(),
648 ];
649 let mut items: Vec<ExprRef> = written
650 .into_iter()
651 .map(|text| self.plan.add_constant(Value::Varchar(text)))
652 .collect();
653 let mark = match column.key {
654 Some(mark) => self.plan.add_constant(Value::Varchar(mark.to_owned())),
655 None => {
656 let empty = self.plan.add_constant(Value::Null);
657 self.cast_to(empty, &LogicalType::Varchar)
658 }
659 };
660 items.push(mark);
661 if let Some(default) = &column.default {
662 items.push(self.plan.add_constant(Value::Varchar(default.clone())));
663 }
664 while items.len() < 6 {
665 let empty = self.plan.add_constant(Value::Null);
666 items.push(self.cast_to(empty, &LogicalType::Varchar));
667 }
668 slices.push(self.plan.add_expr_list(&items));
669 }
670 let rows = self.plan.add_rows(&slices);
671 let columns = self.plan.add_fields(&fields);
672 let index = self.fresh_index();
673 let mut node = self.add_node(Node::Values { index, columns, rows });
674 let mut scope = Scope::empty();
675 for (at, field) in fields.iter().enumerate() {
676 scope.push(Visible {
677 table: String::new(),
678 name: field.name.clone(),
679 binding: ColumnBinding::new(index, at as u32),
680 ty: field.ty.clone(),
681 not_null: false,
682 key: None,
683 default: None,
684 qualified: false,
685 also: None,
686 });
687 }
688 let keys = self.sort_keys(ast, query, &scope, &[])?;
689 if !keys.is_empty() {
690 let keys = self.plan.add_sort_keys(&keys);
691 node = self.add_node(Node::Sort { input: node, keys });
692 }
693 node = self.apply_limit(ast, query, node, &mut scope)?;
694 Ok((node, scope))
695 }
696
697 pub(crate) fn bind_default(&mut self, text: Option<&str>, ty: &LogicalType) -> Result<ExprRef> {
701 let Some(text) = text else {
702 let value = self.plan.add_value(Value::Null);
703 return Ok(self.plan.add_expr(Expr::Constant(value), ty.clone()));
704 };
705 let ast = rudb_parse::parse_ast(&format!("SELECT {text}"))?;
706 let found = match ast.statements.first() {
707 Some(&ast::Statement::Query(query)) => match ast.query(query).body {
708 ast::QueryBody::Select(select) => {
709 ast.target_list(ast.select(select).targets).first().map(|target| target.expr)
710 }
711 _ => None,
712 },
713 _ => None,
714 };
715 let Some(expr) = found else {
716 return Err(Error::internal(format!("a default that is not an expression: {text}")));
717 };
718 let expr = self.bind_expr(&ast, expr, &Scope::empty())?;
719 self.checked_cast_to(expr, ty, false)
720 }
721
722 fn passes_through(&self, expr: ExprRef, input: &Scope) -> bool {
728 let Expr::Column(binding) = *self.plan.expr(expr) else { return false };
729 input.columns.iter().any(|column| column.binding == binding && column.not_null)
730 }
731
732 fn key_through(&self, expr: ExprRef, input: &Scope) -> Option<&'static str> {
736 self.through(expr, input).and_then(|column| column.key)
737 }
738
739 fn through<'s>(&self, expr: ExprRef, input: &'s Scope) -> Option<&'s Visible> {
741 let Expr::Column(binding) = *self.plan.expr(expr) else { return None };
742 input.columns.iter().find(|column| column.binding == binding)
743 }
744
745 fn bind_values(
752 &mut self,
753 ast: &Ast,
754 query: &ast::Query,
755 rows: ast::Slice,
756 ) -> Result<(NodeRef, Scope)> {
757 let written = ast.rows(rows).to_vec();
758 let Some(first) = written.first() else {
759 return Err(Error::binder("VALUES needs at least one row"));
760 };
761 let width = first.len as usize;
762 for (at, row) in written.iter().enumerate() {
763 if row.len as usize != width {
764 return Err(Error::binder(format!(
765 "VALUES lists must all be the same length, expected {width} columns but row {} has {}",
766 at + 1,
767 row.len
768 )));
769 }
770 }
771 let empty = Scope::empty();
773 let defaults = self.insert_defaults.take();
774 let previous = std::mem::replace(&mut self.clause, "VALUES clause");
775 let mut bound: Vec<Vec<ExprRef>> = Vec::with_capacity(written.len());
776 for row in &written {
777 let mut items = Vec::with_capacity(width);
778 for (at, &expr) in ast.expr_list(*row).iter().enumerate() {
779 let column = defaults.as_ref().and_then(|defaults| defaults.get(at));
780 items.push(match (ast.expr(expr), column) {
781 (ast::Expr::Default, Some((ty, default))) => {
782 self.bind_default(default.as_deref(), ty)?
783 }
784 _ => self.bind_expr(ast, expr, &empty)?,
785 });
786 }
787 bound.push(items);
788 }
789 self.clause = previous;
790 let mut types = Vec::with_capacity(width);
791 for at in 0..width {
792 let mut ty = self.plan.expr_type(bound[0][at]).clone();
793 for row in &bound[1..] {
794 let other = self.plan.expr_type(row[at]).clone();
795 ty = ty.promote(&other).ok_or_else(|| {
796 Error::binder(format!(
797 "Cannot combine a value of type {ty} with a value of type {other} in column {} of a VALUES",
798 at + 1
799 ))
800 })?;
801 }
802 types.push(ty);
803 }
804 let mut slices = Vec::with_capacity(bound.len());
805 for row in &bound {
806 let items: Vec<ExprRef> = row
807 .iter()
808 .zip(&types)
809 .map(|(&expr, ty)| self.checked_cast_to(expr, ty, false))
810 .collect::<Result<_>>()?;
811 slices.push(self.plan.add_expr_list(&items));
812 }
813 let rows = self.plan.add_rows(&slices);
814 let fields: Vec<Field> = types
815 .iter()
816 .enumerate()
817 .map(|(at, ty)| Field::new(format!("col{at}"), ty.clone()))
818 .collect();
819 let columns = self.plan.add_fields(&fields);
820 let index = self.fresh_index();
821 let mut node = self.add_node(Node::Values { index, columns, rows });
822 let mut scope = Scope::empty();
823 for (at, field) in fields.iter().enumerate() {
824 scope.push(Visible {
825 table: String::new(),
826 name: field.name.clone(),
827 binding: ColumnBinding::new(index, at as u32),
828 ty: field.ty.clone(),
829 not_null: false,
830 key: None,
831 default: None,
832 qualified: false,
833 also: None,
834 });
835 }
836 let keys = self.sort_keys(ast, query, &scope, &[])?;
837 if !keys.is_empty() {
838 let keys = self.plan.add_sort_keys(&keys);
839 node = self.add_node(Node::Sort { input: node, keys });
840 }
841 node = self.apply_limit(ast, query, node, &mut scope)?;
842 Ok((node, scope))
843 }
844
845 fn bind_set_op(
846 &mut self,
847 ast: &Ast,
848 query: &ast::Query,
849 operator: Operator,
850 left: ast::QueryRef,
851 right: ast::QueryRef,
852 ) -> Result<(NodeRef, Scope)> {
853 let (left_node, left_scope) = self.bind_query(ast, left)?;
854 let (right_node, right_scope) = self.bind_query(ast, right)?;
855 let merged = if operator.by_name {
856 match_by_name(&left_scope, &right_scope)?
857 } else {
858 match_by_position(&left_scope, &right_scope)?
859 };
860 let left_node = self.conform(left_node, &left_scope, &merged, |column| column.left)?;
861 let right_node = self.conform(right_node, &right_scope, &merged, |column| column.right)?;
862 let index = self.fresh_index();
863 let kind = match operator.op {
864 SetOp::Union => SetOpKind::Union,
865 SetOp::Except => SetOpKind::Except,
866 SetOp::Intersect => SetOpKind::Intersect,
867 };
868 let all = operator.quantifier == Quantifier::All;
871 let mut node =
872 self.add_node(Node::SetOp { left: left_node, right: right_node, kind, all, index });
873 let mut scope = Scope::empty();
874 for (at, column) in merged.iter().enumerate() {
875 scope.push(Visible {
876 table: String::new(),
877 name: column.name.clone(),
878 binding: ColumnBinding::new(index, at as u32),
879 ty: column.ty.clone(),
880 not_null: false,
883 key: None,
884 default: None,
885 qualified: false,
886 also: None,
887 });
888 }
889 let keys = self.sort_keys(ast, query, &scope, &[])?;
893 if !keys.is_empty() {
894 let keys = self.plan.add_sort_keys(&keys);
895 node = self.add_node(Node::Sort { input: node, keys });
896 }
897 node = self.apply_limit(ast, query, node, &mut scope)?;
898 Ok((node, scope))
899 }
900
901 fn conform(
907 &mut self,
908 node: NodeRef,
909 scope: &Scope,
910 merged: &[Merged],
911 pick: impl Fn(&Merged) -> Option<usize>,
912 ) -> Result<NodeRef> {
913 let unchanged = merged.len() == scope.len()
914 && merged
915 .iter()
916 .enumerate()
917 .all(|(at, column)| pick(column) == Some(at) && column.ty == scope.columns[at].ty);
918 if unchanged {
919 return Ok(node);
920 }
921 let index = self.fresh_index();
922 let mut exprs = Vec::with_capacity(merged.len());
923 let mut names = Vec::with_capacity(merged.len());
924 for column in merged {
925 let expr = match pick(column) {
926 Some(at) => {
927 let held = &scope.columns[at];
928 self.plan.add_expr(Expr::Column(held.binding), held.ty.clone())
929 }
930 None => self.plan.add_constant(Value::Null),
931 };
932 exprs.push(self.checked_cast_to(expr, &column.ty, false)?);
933 names.push(self.plan.intern(&column.name));
934 }
935 let exprs = self.plan.add_expr_list(&exprs);
936 let names = self.plan.add_name_list(&names);
937 Ok(self.add_node(Node::Project { input: node, index, exprs, names }))
938 }
939
940 fn bind_select(
943 &mut self,
944 ast: &Ast,
945 select: ast::SelectRef,
946 query: &ast::Query,
947 ) -> Result<(NodeRef, Scope)> {
948 let written = ast.select(select);
949 self.want_ascending |= !written.group_by.is_empty() || written.group_by_all;
950 let outer_windows = std::mem::take(&mut self.windows);
954 let outer_joined_above = std::mem::take(&mut self.joined_above);
959 let (mut node, input) = self.bind_from(ast, written.from)?;
960 node = self.attach_scalar_subqueries(node);
961
962 if written.filter != NONE {
963 self.clause = "WHERE clause";
964 let predicate = self.bind_expr(ast, written.filter, &input)?;
965 let predicate = self.as_boolean(predicate, "WHERE")?;
966 node = self.attach_scalar_subqueries(node);
967 node = self.add_node(Node::Filter { input: node, predicate });
968 }
969
970 let targets = ast.target_list(written.targets).to_vec();
971 if targets.is_empty() {
972 return Err(Error::binder("a SELECT needs at least one expression to select"));
973 }
974
975 let group_items = self.group_items(ast, &written, &targets)?;
976 let aggregating = !group_items.is_empty()
977 || written.having != NONE
978 || targets.iter().any(|target| has_aggregate(ast, target.expr));
979 if aggregating {
980 self.clause = "GROUP BY clause";
981 let mut groups = Vec::with_capacity(group_items.len());
982 for item in &group_items {
983 groups.push(self.bind_expr(ast, *item, &input)?);
984 }
985 let index = self.fresh_index();
986 self.aggregation = Some(Aggregation { index, groups, aggregates: Vec::new() });
987 }
988
989 let mut above = Vec::new();
996
997 self.clause = "SELECT clause";
998 let (mut exprs, mut names) = self.bind_targets(ast, &targets, &input, &mut above)?;
999 let visible = exprs.len();
1000
1001 let mut having = None;
1002 if written.having != NONE {
1003 self.clause = "HAVING clause";
1004 let before = self.scalar_subqueries.len();
1005 let predicate = self.bind_expr(ast, written.having, &input)?;
1006 self.lift_over_aggregate(before, &mut above, &input)?;
1007 let predicate = self.over_aggregate(predicate, &input)?;
1008 having = Some(self.as_boolean(predicate, "HAVING")?);
1009 }
1010
1011 let project = self.fresh_index();
1014 let mut output = Scope::empty();
1015 for (at, (expr, name)) in exprs.iter().zip(&names).enumerate() {
1016 output.push(Visible {
1017 table: String::new(),
1018 name: name.clone(),
1019 binding: ColumnBinding::new(project, at as u32),
1020 ty: self.plan.expr_type(*expr).clone(),
1021 not_null: self.passes_through(*expr, &input),
1022 key: self.key_through(*expr, &input),
1023 default: self.through(*expr, &input).and_then(|column| column.default.clone()),
1024 qualified: false,
1025 also: None,
1026 });
1027 }
1028
1029 self.clause = "ORDER BY clause";
1030 let mut extra = Vec::new();
1031 let keys = self.select_sort_keys(
1032 ast, query, &input, &output, project, &mut exprs, &mut names, &mut extra, &mut above,
1033 )?;
1034 self.joined_above = outer_joined_above;
1035 if !extra.is_empty() && written.distinct != Distinct::No {
1036 return Err(Error::binder(
1037 "For SELECT DISTINCT, ORDER BY expressions must appear in the select list",
1038 ));
1039 }
1040 let on = self.distinct_on(ast, written.distinct, &output)?;
1041
1042 node = self.attach_scalar_subqueries(node);
1043
1044 if let Some(aggregation) = self.aggregation.take() {
1045 let index = aggregation.index;
1046 let groups = self.plan.add_expr_list(&aggregation.groups);
1047 let aggregates = self.plan.add_expr_list(&aggregation.aggregates);
1048 node = self.add_node(Node::Aggregate { input: node, index, groups, aggregates });
1049 }
1050 if !above.is_empty() {
1051 debug_assert!(self.scalar_subqueries.is_empty(), "a query is waiting to be joined");
1052 self.scalar_subqueries = above;
1053 node = self.attach_scalar_subqueries(node);
1054 }
1055 if let Some(predicate) = having {
1056 node = self.add_node(Node::Filter { input: node, predicate });
1057 }
1058
1059 for run in std::mem::replace(&mut self.windows, outer_windows) {
1063 let partition = self.plan.add_expr_list(&run.partition);
1064 let order = self.plan.add_sort_keys(&run.order);
1065 let expressions = self.plan.add_expr_list(&run.calls);
1066 node = self.add_node(Node::Window {
1067 input: node,
1068 index: run.index,
1069 partition,
1070 order,
1071 frame: run.frame,
1072 expressions,
1073 });
1074 }
1075
1076 let interned: Vec<u32> = names.iter().map(|name| self.plan.intern(name)).collect();
1077 let exprs_slice = self.plan.add_expr_list(&exprs);
1078 let names_slice = self.plan.add_name_list(&interned);
1079 node = self.add_node(Node::Project {
1080 input: node,
1081 index: project,
1082 exprs: exprs_slice,
1083 names: names_slice,
1084 });
1085
1086 if written.distinct != Distinct::No {
1087 let on = self.plan.add_expr_list(&on);
1088 node = self.add_node(Node::Distinct { input: node, on });
1089 }
1090 if !keys.is_empty() {
1091 let keys = self.plan.add_sort_keys(&keys);
1092 node = self.add_node(Node::Sort { input: node, keys });
1093 }
1094 node = self.apply_limit(ast, query, node, &mut output)?;
1095
1096 if extra.is_empty() {
1097 output.columns.truncate(visible);
1098 return Ok((node, output));
1099 }
1100 let index = self.fresh_index();
1103 let mut kept = Vec::with_capacity(visible);
1104 let mut kept_names = Vec::with_capacity(visible);
1105 let mut scope = Scope::empty();
1106 for (at, name) in names.iter().enumerate().take(visible) {
1107 let ty = output.columns[at].ty.clone();
1108 let binding = output.columns[at].binding;
1112 kept.push(self.plan.add_expr(Expr::Column(binding), ty.clone()));
1113 kept_names.push(self.plan.intern(name));
1114 scope.push(Visible {
1115 table: String::new(),
1116 name: name.clone(),
1117 binding: ColumnBinding::new(index, at as u32),
1118 ty,
1119 not_null: output.columns[at].not_null,
1120 key: output.columns[at].key,
1121 default: output.columns[at].default.clone(),
1122 qualified: false,
1123 also: None,
1124 });
1125 }
1126 let exprs = self.plan.add_expr_list(&kept);
1127 let names = self.plan.add_name_list(&kept_names);
1128 node = self.add_node(Node::Project { input: node, index, exprs, names });
1129 Ok((node, scope))
1130 }
1131
1132 fn lift_over_aggregate(
1151 &mut self,
1152 before: usize,
1153 above: &mut Vec<PendingSubquery>,
1154 scope: &Scope,
1155 ) -> Result<()> {
1156 if self.aggregation.is_none() {
1157 return Ok(());
1158 }
1159 let mut lifted = Vec::new();
1160 for mut pending in self.scalar_subqueries.split_off(before) {
1161 let stays = pending.inside_aggregate
1162 || (pending.dependent && !self.lift_correlated(&mut pending));
1163 if stays {
1164 self.scalar_subqueries.push(pending);
1165 } else {
1166 self.joined_above.push(pending.index);
1167 lifted.push(pending);
1168 }
1169 }
1170 for pending in &mut lifted {
1175 let conditions = std::mem::take(&mut pending.conditions);
1176 let mut over = Vec::with_capacity(conditions.len());
1177 for condition in conditions {
1178 over.push(self.over_aggregate(condition, scope)?);
1179 }
1180 pending.conditions = over;
1181 }
1182 above.append(&mut lifted);
1183 Ok(())
1184 }
1185
1186 fn lift_correlated(&mut self, pending: &mut PendingSubquery) -> bool {
1203 let Some(index) = self.aggregation.as_ref().map(|aggregation| aggregation.index) else {
1204 return false;
1205 };
1206 let mut moved = Vec::with_capacity(pending.reads.len());
1207 for read in &pending.reads {
1208 let Some(at) = self.group_of(*read) else {
1209 return false;
1210 };
1211 moved.push((*read, ColumnBinding::new(index, at as u32)));
1212 }
1213 let mut rewrites = Vec::new();
1214 self.plan.subtree_columns(pending.node, &mut |reference, binding| {
1215 if let Some(&(_, to)) = moved.iter().find(|(from, _)| *from == binding) {
1216 rewrites.push((reference, to));
1217 }
1218 });
1219 for (reference, to) in rewrites {
1220 self.plan.rebind(reference, to);
1221 }
1222 pending.reads = moved.into_iter().map(|(_, to)| to).collect();
1223 true
1224 }
1225
1226 fn bind_targets(
1227 &mut self,
1228 ast: &Ast,
1229 targets: &[ast::Target],
1230 input: &Scope,
1231 above: &mut Vec<PendingSubquery>,
1232 ) -> Result<(Vec<ExprRef>, Vec<String>)> {
1233 let mut exprs = Vec::with_capacity(targets.len());
1234 let mut names = Vec::with_capacity(targets.len());
1235 for target in targets {
1236 if let ast::Expr::Star { qualifier, replacements } = ast.expr(target.expr) {
1237 let table = ast.name(qualifier).last().map(str::to_string);
1238 let expanded: Vec<Visible> =
1239 input.star(table.as_deref())?.into_iter().cloned().collect();
1240 let replacements = ast.target_list(replacements).to_vec();
1241 let mut used = vec![false; replacements.len()];
1242 for column in expanded {
1243 let found = replacements.iter().zip(&mut used).find(|(replacement, _)| {
1244 same_name(ast.string(replacement.alias), &column.name)
1245 });
1246 let before = self.scalar_subqueries.len();
1251 let (expr, name) = match found {
1252 Some((replacement, used)) => {
1253 *used = true;
1254 let expr = self.bind_expr(ast, replacement.expr, input)?;
1255 (expr, ast.string(replacement.alias).to_string())
1256 }
1257 None => (
1258 self.plan.add_expr(Expr::Column(column.binding), column.ty),
1259 column.name,
1260 ),
1261 };
1262 self.lift_over_aggregate(before, above, input)?;
1263 exprs.push(self.over_aggregate(expr, input)?);
1264 names.push(name);
1265 }
1266 if let Some((replacement, _)) =
1270 replacements.iter().zip(&used).find(|(_, used)| !**used)
1271 {
1272 return Err(missing_replacement(ast.string(replacement.alias), input));
1273 }
1274 continue;
1275 }
1276 let before = self.scalar_subqueries.len();
1277 let expr = self.bind_expr(ast, target.expr, input)?;
1278 self.lift_over_aggregate(before, above, input)?;
1279 exprs.push(self.over_aggregate(expr, input)?);
1280 names.push(if target.alias == NONE {
1281 self.output_name(ast, target.expr, input)
1282 } else {
1283 ast.string(target.alias).to_string()
1284 });
1285 }
1286 Ok((exprs, names))
1287 }
1288
1289 fn output_name(&self, ast: &Ast, target: ast::ExprRef, input: &Scope) -> String {
1295 if let ast::Expr::Column { name } = ast.expr(target) {
1296 let parts: Vec<&str> = ast.name(name).collect();
1297 if let Ok(found) = input.resolve(&parts) {
1298 let written = parts.last().copied().unwrap_or_default();
1301 if let Some(also) = &found.also {
1302 if !same_name(&found.name, written) && same_name(also, written) {
1303 return also.clone();
1304 }
1305 }
1306 return found.name.clone();
1307 }
1308 }
1309 describe(ast, target, self.semantics)
1310 }
1311
1312 fn group_items(
1314 &self,
1315 ast: &Ast,
1316 select: &ast::Select,
1317 targets: &[ast::Target],
1318 ) -> Result<Vec<ast::ExprRef>> {
1319 if select.group_by_all {
1320 return Ok(targets
1323 .iter()
1324 .filter(|target| !has_aggregate(ast, target.expr))
1325 .map(|target| target.expr)
1326 .collect());
1327 }
1328 let mut items = Vec::new();
1329 for &item in ast.expr_list(select.group_by) {
1330 items.push(self.output_reference(ast, item, targets, "GROUP BY")?.unwrap_or(item));
1331 }
1332 Ok(items)
1333 }
1334
1335 fn output_reference(
1337 &self,
1338 ast: &Ast,
1339 item: ast::ExprRef,
1340 targets: &[ast::Target],
1341 clause: &str,
1342 ) -> Result<Option<ast::ExprRef>> {
1343 match ast.expr(item) {
1344 ast::Expr::Literal { kind: LiteralKind::Number, text } => {
1345 let written = ast.string(text);
1346 let position: usize = written.parse().map_err(|_| {
1347 Error::binder(format!("{clause} term {written} is not a column"))
1348 })?;
1349 if position == 0 || position > targets.len() {
1350 return Err(Error::binder(format!(
1351 "{clause} term out of range - should be between 1 and {}",
1352 targets.len()
1353 )));
1354 }
1355 Ok(Some(targets[position - 1].expr))
1356 }
1357 ast::Expr::Column { name } => {
1358 let parts: Vec<&str> = ast.name(name).collect();
1359 let [written] = parts.as_slice() else { return Ok(None) };
1360 let mut found = None;
1361 for target in targets {
1362 if target.alias != NONE && same_name(ast.string(target.alias), written) {
1363 if found.is_some() {
1364 return Ok(None);
1365 }
1366 found = Some(target.expr);
1367 }
1368 }
1369 Ok(found)
1370 }
1371 _ => Ok(None),
1372 }
1373 }
1374
1375 #[allow(clippy::too_many_arguments)]
1379 fn select_sort_keys(
1380 &mut self,
1381 ast: &Ast,
1382 query: &ast::Query,
1383 input: &Scope,
1384 output: &Scope,
1385 project: u32,
1386 exprs: &mut Vec<ExprRef>,
1387 names: &mut Vec<String>,
1388 extra: &mut Vec<usize>,
1389 above: &mut Vec<PendingSubquery>,
1390 ) -> Result<Vec<SortKey>> {
1391 if query.order_by_all {
1392 return Ok(self.every_column(output));
1393 }
1394 let items = ast.order_list(query.order_by).to_vec();
1395 let mut keys = Vec::with_capacity(items.len());
1396 for item in items {
1397 self.check_order_literal(ast, item.expr)?;
1398 let position = match self.output_position(ast, item.expr, output)? {
1399 Some(position) => position,
1400 None => {
1401 let before = self.scalar_subqueries.len();
1402 let bound = self.bind_expr(ast, item.expr, input)?;
1403 self.lift_over_aggregate(before, above, input)?;
1404 let bound = self.over_aggregate(bound, input)?;
1405 match exprs.iter().position(|&held| self.same_expr(held, bound)) {
1406 Some(position) => position,
1407 None => {
1408 exprs.push(bound);
1409 names.push(describe(ast, item.expr, self.semantics));
1410 extra.push(exprs.len() - 1);
1411 exprs.len() - 1
1412 }
1413 }
1414 }
1415 };
1416 let ty = self.plan.expr_type(exprs[position]).clone();
1417 let expr = self.column(project, position, ty);
1418 keys.push(self.sort_key(expr, item));
1419 }
1420 Ok(keys)
1421 }
1422
1423 fn sort_keys(
1425 &mut self,
1426 ast: &Ast,
1427 query: &ast::Query,
1428 output: &Scope,
1429 targets: &[ast::Target],
1430 ) -> Result<Vec<SortKey>> {
1431 if query.order_by_all {
1432 return Ok(self.every_column(output));
1433 }
1434 let items = ast.order_list(query.order_by).to_vec();
1435 let mut keys = Vec::with_capacity(items.len());
1436 for item in items {
1437 self.check_order_literal(ast, item.expr)?;
1438 let expr = match self.output_position(ast, item.expr, output)? {
1439 Some(position) => {
1440 let column = &output.columns[position];
1441 let (binding, ty) = (column.binding, column.ty.clone());
1442 self.plan.add_expr(Expr::Column(binding), ty)
1443 }
1444 None => {
1445 let _ = targets;
1446 self.bind_expr(ast, item.expr, output)?
1447 }
1448 };
1449 keys.push(self.sort_key(expr, item));
1450 }
1451 Ok(keys)
1452 }
1453
1454 fn every_column(&mut self, output: &Scope) -> Vec<SortKey> {
1455 let columns: Vec<(ColumnBinding, LogicalType)> =
1456 output.columns.iter().map(|column| (column.binding, column.ty.clone())).collect();
1457 columns
1458 .into_iter()
1459 .map(|(binding, ty)| {
1460 let expr = self.plan.add_expr(Expr::Column(binding), ty);
1461 let descending = self.semantics.default_descending();
1462 SortKey { expr, descending, nulls_first: self.semantics.nulls_first(descending) }
1463 })
1464 .collect()
1465 }
1466
1467 fn sort_key(&self, expr: ExprRef, item: ast::OrderItem) -> SortKey {
1469 let descending = match item.order {
1470 Order::Unstated => self.semantics.default_descending(),
1471 Order::Ascending => false,
1472 Order::Descending => true,
1473 };
1474 let nulls_first = match item.nulls {
1475 Nulls::First => true,
1476 Nulls::Last => false,
1477 Nulls::Unstated => self.semantics.nulls_first(descending),
1478 };
1479 SortKey { expr, descending, nulls_first }
1480 }
1481
1482 fn output_position(
1484 &self,
1485 ast: &Ast,
1486 item: ast::ExprRef,
1487 output: &Scope,
1488 ) -> Result<Option<usize>> {
1489 match ast.expr(item) {
1490 ast::Expr::Literal { kind: LiteralKind::Number, text } => {
1491 let written = ast.string(text);
1492 if written.contains(['.', 'e', 'E']) {
1493 return Ok(None);
1494 }
1495 let position: usize = written.parse().map_err(|_| {
1496 Error::binder(format!("ORDER BY term {written} is not a column"))
1497 })?;
1498 if position == 0 || position > output.len() {
1499 return Err(Error::binder(format!(
1500 "ORDER BY term out of range - should be between 1 and {}",
1501 output.len()
1502 )));
1503 }
1504 Ok(Some(position - 1))
1505 }
1506 ast::Expr::Column { name } => {
1507 let parts: Vec<&str> = ast.name(name).collect();
1508 let [written] = parts.as_slice() else { return Ok(None) };
1509 Ok(output.position_of(None, written))
1510 }
1511 _ => Ok(None),
1512 }
1513 }
1514
1515 fn check_order_literal(&self, ast: &Ast, item: ast::ExprRef) -> Result<()> {
1517 if !self.semantics.order_by_non_integer_literal()
1518 && matches!(
1519 ast.expr(item),
1520 ast::Expr::Literal { kind, text }
1521 if kind != LiteralKind::Number
1522 || ast.string(text).contains(['.', 'e', 'E'])
1523 )
1524 {
1525 return Err(Error::binder(
1526 "ORDER BY non-integer literal has no effect.\n* SET order_by_non_integer_literal=true to allow this behavior.",
1527 ));
1528 }
1529 Ok(())
1530 }
1531
1532 fn distinct_on(
1534 &mut self,
1535 ast: &Ast,
1536 distinct: Distinct,
1537 output: &Scope,
1538 ) -> Result<Vec<ExprRef>> {
1539 let Distinct::On(items) = distinct else {
1540 return Ok(Vec::new());
1541 };
1542 let items = ast.expr_list(items).to_vec();
1543 let mut on = Vec::with_capacity(items.len());
1544 for item in items {
1545 let Some(position) = self.output_position(ast, item, output)? else {
1546 return Err(Error::not_implemented(
1547 "DISTINCT ON an expression that is not in the select list",
1548 ));
1549 };
1550 let column = &output.columns[position];
1551 let (binding, ty) = (column.binding, column.ty.clone());
1552 on.push(self.plan.add_expr(Expr::Column(binding), ty));
1553 }
1554 Ok(on)
1555 }
1556
1557 fn apply_limit(
1564 &mut self,
1565 ast: &Ast,
1566 query: &ast::Query,
1567 input: NodeRef,
1568 scope: &mut Scope,
1569 ) -> Result<NodeRef> {
1570 let waiting = self.scalar_subqueries.len();
1571 if query.limit_percent {
1572 let percent = self.share(ast, query.limit)?;
1573 let offset = self.skipped(ast, query.offset)?;
1574 let node = |binder: &mut Self, input| match percent {
1575 Some(percent) => binder.add_node(Node::LimitPercent { input, percent, offset }),
1576 None => binder.limited(input, Bound::All, offset),
1579 };
1580 return self.over_subqueries(waiting, input, scope, node);
1581 }
1582 let count = self.count_bound(ast, query.limit, "LIMIT")?;
1583 let offset = self.skipped(ast, query.offset)?;
1584 let node = |binder: &mut Self, input| binder.limited(input, count, offset);
1585 self.over_subqueries(waiting, input, scope, node)
1586 }
1587
1588 fn skipped(&mut self, ast: &Ast, written: ast::ExprRef) -> Result<Bound> {
1593 Ok(match self.count_bound(ast, written, "OFFSET")? {
1594 Bound::All => Bound::Rows(0),
1595 named => named,
1596 })
1597 }
1598
1599 fn over_subqueries(
1606 &mut self,
1607 waiting: usize,
1608 input: NodeRef,
1609 scope: &mut Scope,
1610 node: impl FnOnce(&mut Self, NodeRef) -> NodeRef,
1611 ) -> Result<NodeRef> {
1612 let joined = self.scalar_subqueries.split_off(waiting);
1613 if joined.is_empty() {
1614 return Ok(node(self, input));
1615 }
1616 let mut input = input;
1617 for pending in joined {
1618 input = self.attach_subquery(input, pending);
1619 }
1620 let limit = node(self, input);
1621 Ok(self.reproject(limit, scope))
1622 }
1623
1624 fn limited(&mut self, input: NodeRef, count: Bound, offset: Bound) -> NodeRef {
1627 if count == Bound::All && offset == Bound::Rows(0) {
1628 return input;
1629 }
1630 self.add_node(Node::Limit { input, count, offset })
1631 }
1632
1633 fn reproject(&mut self, node: NodeRef, scope: &mut Scope) -> NodeRef {
1639 let index = self.fresh_index();
1640 let mut exprs = Vec::with_capacity(scope.columns.len());
1641 let mut names = Vec::with_capacity(scope.columns.len());
1642 for column in &scope.columns {
1643 exprs.push(self.plan.add_expr(Expr::Column(column.binding), column.ty.clone()));
1644 names.push(self.plan.intern(&column.name));
1645 }
1646 for (at, column) in scope.columns.iter_mut().enumerate() {
1647 column.binding = ColumnBinding::new(index, at as u32);
1648 }
1649 let exprs = self.plan.add_expr_list(&exprs);
1650 let names = self.plan.add_name_list(&names);
1651 self.add_node(Node::Project { input: node, index, exprs, names })
1652 }
1653
1654 fn share(&mut self, ast: &Ast, written: ast::ExprRef) -> Result<Option<Share>> {
1671 if written == NONE {
1672 return Ok(None);
1673 }
1674 self.clause = "LIMIT clause";
1675 let scope = Scope::empty();
1676 let bound = self.bind_expr(ast, written, &scope)?;
1677 let Some(value) = fold::value_of(&self.plan, bound)? else {
1678 return Ok(Some(Share::Read(bound)));
1679 };
1680 if value.is_null() {
1681 return Ok(None);
1682 }
1683 let percent = percentage(&value)?;
1684 if !(0.0..=100.0).contains(&percent) {
1685 return Err(Error::out_of_range(
1686 "Limit percent out of range, should be between 0% and 100%",
1687 ));
1688 }
1689 Ok(Some(Share::Percent(percent)))
1690 }
1691
1692 fn count_bound(&mut self, ast: &Ast, written: ast::ExprRef, clause: &str) -> Result<Bound> {
1711 if written == NONE {
1712 return Ok(Bound::All);
1713 }
1714 self.clause = "LIMIT clause";
1715 let scope = Scope::empty();
1716 let bound = self.bind_expr(ast, written, &scope)?;
1717 let Some(value) = fold::value_of(&self.plan, bound)? else {
1718 return Ok(Bound::Read(bound));
1719 };
1720 if value.is_null() {
1723 return Ok(Bound::All);
1724 }
1725 row_count(&value, clause).map(Bound::Rows)
1726 }
1727
1728 fn bind_from(&mut self, ast: &Ast, from: ast::Slice) -> Result<(NodeRef, Scope)> {
1731 let sources = ast.source_list(from).to_vec();
1732 let Some((first, rest)) = sources.split_first() else {
1733 return Ok((self.add_node(Node::Dummy), Scope::empty()));
1736 };
1737 let (mut node, mut scope) = self.bind_source(ast, *first)?;
1738 for source in rest {
1739 let (right, right_scope, correlations) = self.bind_lateral(ast, *source, &scope)?;
1740 node = if correlations.is_empty() {
1741 self.add_node(Node::CrossProduct { left: node, right })
1742 } else {
1743 let conditions = self.plan.add_expr_list(&[]);
1744 self.add_node(Node::DependentJoin {
1745 left: node,
1746 right,
1747 kind: JoinKind::Inner,
1748 conditions,
1749 })
1750 };
1751 scope = scope.concat(right_scope);
1752 }
1753 Ok((node, scope))
1754 }
1755
1756 fn bind_lateral(
1768 &mut self,
1769 ast: &Ast,
1770 source: ast::SourceRef,
1771 left: &Scope,
1772 ) -> Result<(NodeRef, Scope, Vec<ColumnBinding>)> {
1773 self.lateral_scopes.push(self.outer_scopes.len());
1774 self.outer_scopes.push(left.clone());
1775 self.correlations.push(Vec::new());
1776 let bound = self.bind_source(ast, source);
1777 let read = self.correlations.pop().expect("correlation frame");
1778 self.outer_scopes.pop();
1779 self.lateral_scopes.pop();
1780 let (node, scope) = bound?;
1781
1782 let mut here = Vec::new();
1783 for binding in read {
1784 if left.columns.iter().any(|column| column.binding == binding) {
1785 here.push(binding);
1786 } else if let Some(enclosing) = self.correlations.last_mut() {
1787 if !enclosing.contains(&binding) {
1788 enclosing.push(binding);
1789 }
1790 }
1791 }
1792 Ok((node, scope, here))
1802 }
1803
1804 fn bind_source(&mut self, ast: &Ast, source: ast::SourceRef) -> Result<(NodeRef, Scope)> {
1805 match ast.source(source) {
1806 ast::Source::Table { name, alias, columns } => {
1807 self.bind_table(ast, name, alias, columns)
1808 }
1809 ast::Source::Function { name, args, alias, columns, pragma } => {
1810 self.bind_table_function(ast, name, args, alias, columns, pragma)
1811 }
1812 ast::Source::Subquery { query, alias, columns } => {
1813 let (node, mut scope) = self.bind_query(ast, query)?;
1814 let label = if alias == NONE {
1815 "unnamed_subquery".to_string()
1816 } else {
1817 ast.string(alias).to_string()
1818 };
1819 scope.relabel(&label);
1820 if !columns.is_empty() {
1821 let names: Vec<&str> = ast.name(columns).collect();
1822 scope.rename(&names, &label)?;
1823 }
1824 Ok((node, scope))
1825 }
1826 ast::Source::Values { rows, alias, columns } => {
1827 let bare = ast::Query::bare(ast::QueryBody::Values(rows));
1828 let (node, mut scope) = self.bind_values(ast, &bare, rows)?;
1829 let label =
1830 if alias == NONE { String::new() } else { ast.string(alias).to_string() };
1831 scope.relabel(&label);
1832 if !columns.is_empty() {
1833 let names: Vec<&str> = ast.name(columns).collect();
1834 scope.rename(&names, &label)?;
1835 }
1836 Ok((node, scope))
1837 }
1838 ast::Source::Cte { cte, alias, columns } => {
1839 self.bind_cte_scan(ast, cte, alias, columns)
1840 }
1841 ast::Source::Join { left, right, kind, natural, on, using } => {
1842 self.bind_join(ast, left, right, kind, natural, on, using)
1843 }
1844 }
1845 }
1846
1847 fn bind_cte_scan(
1854 &mut self,
1855 ast: &Ast,
1856 written: u32,
1857 alias: ast::StrRef,
1858 columns: ast::Slice,
1859 ) -> Result<(NodeRef, Scope)> {
1860 let Some(held) = self.materialized.iter().rev().find(|held| held.written == written) else {
1861 let name = ast.string(ast.cte(written).name);
1862 return Err(Error::binder(format!("Table with name {name} does not exist!")));
1863 };
1864 let cte = held.cte;
1865 let fields = held.fields.clone();
1866 let text = held.name.clone();
1867 let label = if alias == NONE { text.clone() } else { ast.string(alias).to_string() };
1868 let name = self.plan.intern(&text);
1869 let index = self.fresh_index();
1870 let mut scope = Scope::empty();
1871 for (at, field) in fields.iter().enumerate() {
1872 scope.push(Visible {
1873 table: label.clone(),
1874 name: field.name.clone(),
1875 binding: ColumnBinding::new(index, at as u32),
1876 ty: field.ty.clone(),
1877 not_null: field.not_null,
1878 key: None,
1879 default: None,
1880 qualified: false,
1881 also: None,
1882 });
1883 }
1884 if !columns.is_empty() {
1885 let names: Vec<&str> = ast.name(columns).collect();
1886 scope.rename(&names, &label)?;
1887 }
1888 let columns = self.plan.add_fields(&fields);
1889 let node = self.add_node(Node::CteScan { index, cte, name, columns });
1890 Ok((node, scope))
1891 }
1892
1893 fn bind_table(
1894 &mut self,
1895 ast: &Ast,
1896 name: ast::Slice,
1897 alias: ast::StrRef,
1898 columns: ast::Slice,
1899 ) -> Result<(NodeRef, Scope)> {
1900 let parts: Vec<&str> = ast.name(name).collect();
1901 let catalog = self.catalog;
1902 let resolved = match catalog.resolve(&parts) {
1905 Ok(resolved) => resolved,
1906 Err(missing) => {
1907 return self.bind_replacement_scan(ast, &parts, alias, columns, missing);
1908 }
1909 };
1910 if catalog.entry(&resolved)? == Entry::View {
1911 return self.bind_view(ast, &resolved, alias, columns);
1912 }
1913 let label =
1914 if alias == NONE { resolved.table.clone() } else { ast.string(alias).to_string() };
1915 self.bind_catalog_table(ast, &resolved, label, columns)
1916 }
1917
1918 pub(crate) fn bind_catalog_table(
1921 &mut self,
1922 ast: &Ast,
1923 resolved: &QualifiedName,
1924 label: String,
1925 columns: ast::Slice,
1926 ) -> Result<(NodeRef, Scope)> {
1927 let table = self.catalog.table(resolved)?;
1928 let fields: Vec<Field> = table.columns().to_vec();
1929 let excluded = self.upsert && same_name(&label, "excluded");
1932 let mut marks = vec![None; fields.len()];
1935 for key in table.keys() {
1936 for &column in &key.columns {
1937 if key.primary || marks[column].is_none() {
1938 marks[column] = Some(if key.primary { "PRI" } else { "UNI" });
1939 }
1940 }
1941 }
1942 let index = self.fresh_index();
1943 let mut scope = Scope::empty();
1944 for (at, field) in fields.iter().enumerate() {
1945 scope.push(Visible {
1946 table: label.clone(),
1947 name: field.name.clone(),
1948 binding: ColumnBinding::new(index, at as u32),
1949 ty: field.ty.clone(),
1950 not_null: field.not_null,
1951 key: marks[at],
1952 default: table.default(at).map(str::to_owned),
1953 qualified: excluded,
1954 also: None,
1955 });
1956 }
1957 if !columns.is_empty() {
1958 let names: Vec<&str> = ast.name(columns).collect();
1959 scope.rename(&names, &label)?;
1960 }
1961 let resolved = if excluded { &QualifiedName::excluded() } else { resolved };
1962 let catalog_name = self.plan.intern(&resolved.catalog);
1963 let schema = self.plan.intern(&resolved.schema);
1964 let table_name = self.plan.intern(&resolved.table);
1965 let alias = self.plan.intern(&label);
1966 let columns = self.plan.add_fields(&fields);
1967 if let Some(zones) = table.rows().zones().filter(|_| !excluded) {
1972 self.plan.set_zones(index, zones);
1973 }
1974 if let Some(frequencies) = table.frequencies().filter(|_| !excluded) {
1975 self.plan.set_frequencies(index, frequencies);
1976 }
1977 for (column, distinct) in table.distincts() {
1978 if !excluded {
1979 self.plan.measure_distinct(index, &column, distinct);
1980 }
1981 }
1982 if self.want_ascending && !excluded {
1983 for column in table.ascending() {
1984 self.plan.mark_ascending(index, &column);
1985 }
1986 }
1987 let node = self.add_node(Node::Get {
1988 catalog: catalog_name,
1989 schema,
1990 table: table_name,
1991 alias,
1992 index,
1993 columns,
1994 });
1995 Ok((node, scope))
1996 }
1997
1998 fn bind_view(
2010 &mut self,
2011 ast: &Ast,
2012 name: &QualifiedName,
2013 alias: ast::StrRef,
2014 columns: ast::Slice,
2015 ) -> Result<(NodeRef, Scope)> {
2016 let view = self.catalog.view(name)?;
2017 let full = name.to_string();
2018 if self.expanding.contains(&full) {
2019 return Err(Error::binder(format!(
2023 "infinite recursion detected: attempting to recursively bind view \"\"{}\"\"",
2024 name.table
2025 )));
2026 }
2027 let body = parse_ast_with_case(view.sql(), self.semantics.identifier_case())?;
2028 let query = match body.statements.as_slice() {
2029 [ast::Statement::Query(query)] => *query,
2030 _ => return Err(Error::binder(format!("view \"{}\" is not a query", name.table))),
2033 };
2034 self.expanding.push(full);
2035 let bound = self.bind_query(&body, query);
2036 self.expanding.pop();
2037 let (node, mut scope) = bound?;
2038
2039 let aliases: Vec<&str> = view.aliases().iter().map(String::as_str).collect();
2040 if !aliases.is_empty() {
2041 scope.rename(&aliases, "unnamed_subquery")?;
2042 }
2043 view.remember(scope.fields());
2050 let label = if alias == NONE { name.table.clone() } else { ast.string(alias).to_string() };
2051 scope.relabel(&label);
2052 if !columns.is_empty() {
2053 let names: Vec<&str> = ast.name(columns).collect();
2054 scope.rename(&names, &label)?;
2055 }
2056 Ok((node, scope))
2057 }
2058
2059 fn bind_table_function(
2067 &mut self,
2068 ast: &Ast,
2069 name: ast::Slice,
2070 args: ast::Slice,
2071 alias: ast::StrRef,
2072 columns: ast::Slice,
2073 pragma: bool,
2074 ) -> Result<(NodeRef, Scope)> {
2075 let renamed = columns;
2078 let parts: Vec<&str> = ast.name(name).collect();
2079 let function_name = *parts.last().unwrap_or(&"");
2083 if let Some(schema) = parts.iter().rev().nth(1) {
2084 if !schema.eq_ignore_ascii_case("main") && !schema.eq_ignore_ascii_case("system") {
2085 return Err(Error::catalog(format!(
2086 "Table Function with name {} does not exist!",
2087 parts.join(".")
2088 )));
2089 }
2090 }
2091 let Some(called) = TableFunction::lookup(function_name) else {
2095 if pragma {
2096 if args.is_empty() && self.catalog.resolve(&parts).is_ok() {
2102 return self.bind_table(ast, name, alias, columns);
2103 }
2104 let spelled = function_name.strip_prefix("pragma_").unwrap_or(function_name);
2105 return Err(Error::catalog(format!(
2106 "Pragma Function with name {spelled} does not exist!"
2107 )));
2108 }
2109 return Err(Error::catalog(format!(
2110 "Table Function with name {function_name} does not exist!"
2111 )));
2112 };
2113 let written = ast.target_list(args).to_vec();
2114 let empty = Scope::empty();
2115 let previous = std::mem::replace(&mut self.clause, "table function arguments");
2116 let mut bound = Vec::new();
2117 let mut written_options = Vec::new();
2118 for argument in written {
2119 let expr = self.bind_expr(ast, argument.expr, &empty)?;
2120 if argument.alias == NONE {
2121 bound.push(expr);
2122 } else {
2123 let name = ast.string(argument.alias).to_string();
2124 let (parameter, value) = self.named_argument(called, &name, expr)?;
2125 written_options.push((parameter, value, expr));
2126 }
2127 }
2128 self.clause = previous;
2129 let options = Options::of(&written_options)?;
2130
2131 let given: Vec<LogicalType> =
2134 bound.iter().map(|&expr| self.plan.expr_type(expr).clone()).collect();
2135 let resolved = if pragma {
2136 resolve_pragma(function_name, &given)?
2137 } else {
2138 resolve_table(function_name, &given)?
2139 };
2140 let mut cast: Vec<ExprRef> = bound
2141 .iter()
2142 .zip(&resolved.arguments)
2143 .map(|(&expr, ty)| self.checked_cast_to(expr, ty, false))
2144 .collect::<Result<_>>()?;
2145
2146 if resolved.function.answered_when_bound() {
2147 let Columns::Fixed(fields) = resolved.columns else {
2148 return Err(Error::internal("a pragma that resolved to a file"));
2149 };
2150 let [argument] = cast[..] else {
2151 return Err(Error::internal("a pragma that resolved to more than one name"));
2152 };
2153 return self.bind_pragma(ast, resolved.function, &fields, argument, alias, columns);
2154 }
2155 let mut measured = Stat::Unknown;
2158 let mut counted: Vec<(String, Stat<u64>)> = Vec::new();
2159 let mut bounded: Option<Arc<dyn Zones>> = None;
2160 let fields = match resolved.columns {
2161 Columns::Fixed(fields) => fields,
2162 columns => {
2163 let paths = self.file_paths(cast[0], resolved.function.name())?;
2168 let mut mirrorable = None;
2169 if resolved.function == TableFunction::ReadParquet && !options.file_row_number {
2170 if let Some((path, stamp)) = mirror_target(&paths) {
2171 if let Some(name) =
2172 self.catalog.mirror(&path, options.binary_as_string, stamp)
2173 {
2174 let name = name.clone();
2175 let label = if alias == NONE {
2176 resolved.function.name().to_string()
2177 } else {
2178 ast.string(alias).to_string()
2179 };
2180 return self.bind_catalog_table(ast, &name, label, renamed);
2181 }
2182 mirrorable = Some(path);
2183 }
2184 }
2185 let mut fields = match columns {
2186 Columns::Csv => csv_fields(&paths, options.given)?,
2189 _ => {
2190 let footers = self.footers(&paths, mirrorable.as_deref())?;
2191 if let Some(path) = mirrorable.as_deref() {
2192 self.want_mirror(path, options.binary_as_string, &footers.rows);
2193 }
2194 measured = footers.rows;
2195 counted = footers.distincts;
2196 bounded = footers.zones;
2197 footers.fields
2198 }
2199 };
2200 if options.all_varchar {
2201 for field in &mut fields {
2206 field.ty = LogicalType::Varchar;
2207 }
2208 }
2209 if options.binary_as_string {
2210 for field in &mut fields {
2215 if field.ty == LogicalType::Blob {
2216 field.ty = LogicalType::Varchar;
2217 }
2218 }
2219 }
2220 if options.file_row_number {
2221 if fields.iter().any(|field| field.name == FILE_ROW_NUMBER) {
2227 return Err(Error::binder(format!(
2228 "Duplicate column name \"{FILE_ROW_NUMBER}\": the file already has a \
2229 column of that name, so file_row_number cannot add one"
2230 )));
2231 }
2232 fields.push(Field::required(FILE_ROW_NUMBER.to_string(), LogicalType::BigInt));
2233 }
2234 cast = paths.iter().map(|path| self.path_constant(path)).collect();
2235 fields
2236 }
2237 };
2238 let label = if alias == NONE {
2239 resolved.function.name().to_string()
2240 } else {
2241 ast.string(alias).to_string()
2242 };
2243 let names: Vec<&str> = ast.name(columns).collect();
2244 self.table_function_source(
2245 resolved.function,
2246 &cast,
2247 &written_options,
2248 Read { fields, rows: measured, distincts: counted, zones: bounded },
2249 &label,
2250 &names,
2251 )
2252 }
2253
2254 fn bind_pragma(
2267 &mut self,
2268 ast: &Ast,
2269 function: TableFunction,
2270 fields: &[Field],
2271 argument: ExprRef,
2272 alias: ast::StrRef,
2273 columns: ast::Slice,
2274 ) -> Result<(NodeRef, Scope)> {
2275 let written = self.pragma_name(argument, function)?;
2276 let parts = identifier_parts(&written);
2277 let spelled: Vec<&str> = parts.iter().map(String::as_str).collect();
2278 let name = self.catalog.resolve(&spelled)?;
2279 let described = self.described(ast, &name)?;
2280 let mut rows = Vec::with_capacity(described.len());
2281 for (at, field) in described.iter().enumerate() {
2282 let items = if matches!(function, TableFunction::PragmaShow) {
2283 self.describing(field)
2284 } else {
2285 self.table_info(at, field)
2286 };
2287 rows.push(self.plan.add_expr_list(&items));
2288 }
2289 let rows = self.plan.add_rows(&rows);
2290 let held = self.plan.add_fields(fields);
2291 let index = self.fresh_index();
2292 let node = self.add_node(Node::Values { index, columns: held, rows });
2293 let label =
2294 if alias == NONE { function.name().to_string() } else { ast.string(alias).to_string() };
2295 let mut scope = Scope::empty();
2296 for (at, field) in fields.iter().enumerate() {
2297 scope.push(Visible {
2298 table: label.clone(),
2299 name: field.name.clone(),
2300 binding: ColumnBinding::new(index, at as u32),
2301 ty: field.ty.clone(),
2302 not_null: false,
2303 key: None,
2304 default: None,
2305 qualified: false,
2306 also: None,
2307 });
2308 }
2309 if !columns.is_empty() {
2310 let names: Vec<&str> = ast.name(columns).collect();
2311 scope.rename(&names, &label)?;
2312 }
2313 Ok((node, scope))
2314 }
2315
2316 fn pragma_name(&self, argument: ExprRef, function: TableFunction) -> Result<String> {
2326 let Expr::Constant(reference) = *self.plan.expr(argument) else {
2327 return Err(Error::not_implemented(format!(
2328 "{}() given a name that is not a constant",
2329 function.name()
2330 )));
2331 };
2332 match self.plan.value(reference) {
2333 Value::Varchar(name) => Ok(name.clone()),
2334 Value::Null => Ok("NULL".to_string()),
2335 other => {
2336 Err(Error::internal(format!("a pragma name bound as VARCHAR arrived as {other}")))
2337 }
2338 }
2339 }
2340
2341 fn described(&mut self, ast: &Ast, name: &QualifiedName) -> Result<Vec<Field>> {
2352 if self.catalog.entry(name)? == Entry::Table {
2353 return Ok(self.catalog.table(name)?.columns().to_vec());
2354 }
2355 let (_, scope) = self.bind_view(ast, name, NONE, ast::Slice::default())?;
2356 Ok(scope.fields())
2357 }
2358
2359 fn describing(&mut self, field: &Field) -> Vec<ExprRef> {
2361 let written = [
2362 field.name.clone(),
2363 field.ty.to_string(),
2364 if field.not_null { "NO" } else { "YES" }.to_owned(),
2365 ];
2366 let mut items: Vec<ExprRef> =
2367 written.into_iter().map(|text| self.plan.add_constant(Value::Varchar(text))).collect();
2368 for _ in 0..3 {
2369 let empty = self.plan.add_constant(Value::Null);
2370 items.push(self.cast_to(empty, &LogicalType::Varchar));
2371 }
2372 items
2373 }
2374
2375 fn table_info(&mut self, at: usize, field: &Field) -> Vec<ExprRef> {
2381 let cid = self.plan.add_constant(Value::Integer(i32::try_from(at).unwrap_or(i32::MAX)));
2382 let name = self.plan.add_constant(Value::Varchar(field.name.clone()));
2383 let ty = self.plan.add_constant(Value::Varchar(field.ty.to_string()));
2384 let not_null = self.plan.add_constant(Value::Boolean(field.not_null));
2385 let default = self.plan.add_constant(Value::Null);
2386 let default = self.cast_to(default, &LogicalType::Varchar);
2387 let key = self.plan.add_constant(Value::Boolean(false));
2388 vec![cid, name, ty, not_null, default, key]
2389 }
2390
2391 fn named_argument(
2405 &mut self,
2406 function: TableFunction,
2407 name: &str,
2408 expr: ExprRef,
2409 ) -> Result<(&'static str, Value)> {
2410 let known = function
2411 .parameters()
2412 .iter()
2413 .find(|(parameter, _)| parameter.eq_ignore_ascii_case(name));
2414 let Some((parameter, wanted)) = known else {
2415 let candidates: Vec<String> = function
2416 .parameters()
2417 .iter()
2418 .map(|(parameter, ty)| format!(" {parameter} {ty}"))
2419 .collect();
2420 return Err(Error::binder(format!(
2421 "Invalid named parameter \"{name}\" for function {}\nCandidates:\n{}\n",
2422 function.name(),
2423 candidates.join("\n")
2424 )));
2425 };
2426 let Expr::Constant(reference) = *self.plan.expr(expr) else {
2427 return Err(Error::not_implemented(format!(
2428 "the named parameter {parameter} with a value that is not a constant"
2429 )));
2430 };
2431 let value = self.plan.value(reference).clone();
2432 if value == Value::Null {
2433 return Err(Error::binder(null_parameter(function, parameter)));
2434 }
2435 let given = self.plan.expr_type(expr).clone();
2436 if given != *wanted {
2437 return Err(Error::not_implemented(format!(
2438 "the named parameter {parameter} given a {given} where a {wanted} was wanted"
2439 )));
2440 }
2441 Ok((parameter, value))
2442 }
2443
2444 fn bind_replacement_scan(
2455 &mut self,
2456 ast: &Ast,
2457 parts: &[&str],
2458 alias: ast::StrRef,
2459 columns: ast::Slice,
2460 missing: Error,
2461 ) -> Result<(NodeRef, Scope)> {
2462 let [path] = parts else { return Err(missing) };
2463 let path = *path;
2464 let extension = path.rsplit_once('.').map(|(_, after)| after).unwrap_or_default();
2465 let Some(function) = Self::reader_for(extension) else {
2466 if is_file(path) {
2467 return Err(Error::binder(format!(
2472 "No extension found that is capable of reading the file \"{path}\"\n* If this \
2473 file is a supported file format you can explicitly use the reader functions, \
2474 such as read_csv, read_json or read_parquet"
2475 )));
2476 }
2477 return Err(missing);
2478 };
2479 let paths = files(path)?;
2484 let label = if alias == NONE {
2490 if is_pattern(path) {
2491 path.to_string()
2492 } else {
2493 let file = path.rsplit_once('/').map_or(path, |(_, file)| file);
2494 file.rsplit_once('.').map_or(file, |(stem, _)| stem).to_string()
2495 }
2496 } else {
2497 ast.string(alias).to_string()
2498 };
2499 let mut mirrorable = None;
2500 if function == TableFunction::ReadParquet {
2501 if let Some((canonical, stamp)) = mirror_target(&paths) {
2502 if let Some(name) = self.catalog.mirror(&canonical, false, stamp) {
2503 let name = name.clone();
2504 return self.bind_catalog_table(ast, &name, label, columns);
2505 }
2506 mirrorable = Some(canonical);
2507 }
2508 }
2509 let read = match function {
2510 TableFunction::ReadParquet => {
2511 let footers = self.footers(&paths, mirrorable.as_deref())?;
2512 if let Some(canonical) = mirrorable.as_deref() {
2513 self.want_mirror(canonical, false, &footers.rows);
2514 }
2515 Read {
2516 fields: footers.fields,
2517 rows: footers.rows,
2518 distincts: footers.distincts,
2519 zones: footers.zones,
2520 }
2521 }
2522 _ => Read::uncounted(csv_fields(&paths, Given::default())?),
2523 };
2524 let arguments: Vec<ExprRef> = paths.iter().map(|path| self.path_constant(path)).collect();
2525 let names: Vec<&str> = ast.name(columns).collect();
2526 self.table_function_source(function, &arguments, &[], read, &label, &names)
2527 }
2528
2529 fn footers(&self, paths: &[String], mirrorable: Option<&str>) -> Result<Footers> {
2535 if let Some(path) = mirrorable.filter(|_| self.outlined) {
2536 let outline = parquet_outline(path)?;
2537 if outline.rows.value().is_some() {
2538 return Ok(outline);
2539 }
2540 }
2541 parquet_footers(paths)
2542 }
2543
2544 fn want_mirror(&mut self, path: &str, binary_as_string: bool, rows: &Stat<u64>) {
2548 if let Some(&rows) = rows.value() {
2549 self.plan.want_mirror(path, binary_as_string, rows);
2550 }
2551 }
2552
2553 fn path_constant(&mut self, path: &str) -> ExprRef {
2555 let value = self.plan.add_value(Value::Varchar(path.to_string()));
2556 self.plan.add_expr(Expr::Constant(value), LogicalType::Varchar)
2557 }
2558
2559 fn reader_for(extension: &str) -> Option<TableFunction> {
2566 if extension.eq_ignore_ascii_case("parquet") {
2567 return Some(TableFunction::ReadParquet);
2568 }
2569 if extension.eq_ignore_ascii_case("csv") || extension.eq_ignore_ascii_case("tsv") {
2570 return Some(TableFunction::ReadCsv);
2571 }
2572 None
2573 }
2574
2575 fn table_function_source(
2585 &mut self,
2586 function: TableFunction,
2587 args: &[ExprRef],
2588 written: &[(&'static str, Value, ExprRef)],
2589 read: Read,
2590 label: &str,
2591 names: &[&str],
2592 ) -> Result<(NodeRef, Scope)> {
2593 let Read { fields, rows, distincts, zones } = read;
2594 let index = self.fresh_index();
2595 if rows.is_known() {
2600 self.plan.measure(index, rows);
2601 }
2602 for (column, distinct) in distincts {
2603 self.plan.measure_distinct(index, &column, distinct);
2604 }
2605 if let Some(zones) = zones {
2606 self.plan.set_zones(index, zones);
2607 }
2608 let mut scope = Scope::empty();
2609 for (at, field) in fields.iter().enumerate() {
2610 scope.push(Visible {
2611 table: label.to_string(),
2612 name: field.name.clone(),
2613 binding: ColumnBinding::new(index, at as u32),
2614 ty: field.ty.clone(),
2615 not_null: false,
2618 key: None,
2619 default: None,
2620 qualified: false,
2621 also: None,
2622 });
2623 }
2624 if !names.is_empty() {
2625 scope.rename(names, label)?;
2626 } else if matches!(function, TableFunction::Range | TableFunction::GenerateSeries) {
2627 for column in &mut scope.columns {
2630 column.also = Some(std::mem::replace(&mut column.name, label.to_string()));
2631 }
2632 }
2633 let function = self.plan.intern(function.name());
2634 let args = self.plan.add_expr_list(args);
2635 let named: Vec<u32> =
2636 written.iter().map(|(parameter, _, _)| self.plan.intern(parameter)).collect();
2637 let settings: Vec<ExprRef> = written.iter().map(|(_, _, expr)| *expr).collect();
2638 let options = self.plan.add_name_list(&named);
2639 let settings = self.plan.add_expr_list(&settings);
2640 let columns = self.plan.add_fields(&fields);
2641 let node = self.add_node(Node::TableFunction {
2642 index,
2643 function,
2644 args,
2645 options,
2646 settings,
2647 columns,
2648 });
2649 Ok((node, scope))
2650 }
2651
2652 fn file_paths(&self, expr: ExprRef, name: &str) -> Result<Vec<String>> {
2659 let mut paths = Vec::new();
2660 for pattern in self.file_patterns(expr, name)? {
2661 paths.extend(files(&pattern)?);
2662 }
2663 Ok(paths)
2664 }
2665
2666 fn file_patterns(&self, expr: ExprRef, name: &str) -> Result<Vec<String>> {
2685 let Some(value) = fold::value_of(&self.plan, expr)? else {
2686 return Err(Error::not_implemented(
2687 "a table function file name that is not a constant",
2688 ));
2689 };
2690 match value {
2691 Value::Varchar(path) => Ok(vec![path]),
2692 Value::Null => Err(Error::parser(format!("{name} cannot take NULL list as parameter"))),
2694 Value::List { values, .. } if values.is_empty() => {
2699 Err(Error::io(format!("\"{name}\" needs at least one file to read")))
2700 }
2701 Value::List { values, .. } => values
2702 .iter()
2703 .map(|value| match value {
2704 Value::Varchar(path) => Ok(path.clone()),
2705 _ => Err(Error::parser(format!(
2706 "{name} reader cannot take NULL input as parameter"
2707 ))),
2708 })
2709 .collect(),
2710 other => {
2711 Err(Error::internal(format!("a file name bound as VARCHAR arrived as {other}")))
2712 }
2713 }
2714 }
2715
2716 fn side_of(
2734 &self,
2735 pending: &PendingSubquery,
2736 left_tables: &[u32],
2737 right_tables: &[u32],
2738 ) -> Option<Side> {
2739 let mut needs_left = false;
2740 let mut needs_right = false;
2741 let mut note = |binding: ColumnBinding| {
2742 needs_left |= left_tables.contains(&binding.table);
2743 needs_right |= right_tables.contains(&binding.table);
2744 };
2745 for &binding in &pending.reads {
2746 note(binding);
2747 }
2748 for &condition in &pending.conditions {
2753 self.plan.read_columns(condition, &mut |_, binding| note(binding));
2754 }
2755 match (needs_left, needs_right) {
2756 (true, true) => None,
2757 (_, true) => Some(Side::Right),
2758 _ => Some(Side::Left),
2759 }
2760 }
2761
2762 #[allow(clippy::too_many_arguments)]
2783 fn bind_pair_dependent_join(
2784 &mut self,
2785 kind: ast::JoinKind,
2786 independent: bool,
2787 left: NodeRef,
2788 right: NodeRef,
2789 pair: Vec<PendingSubquery>,
2790 conditions: Vec<ExprRef>,
2791 scope: Scope,
2792 ) -> Result<(NodeRef, Scope)> {
2793 if kind != ast::JoinKind::Inner {
2794 return Err(Error::not_implemented(
2795 "a subquery that reads both sides of that join, written in the condition of a join \
2796 that is not an inner join"
2797 .to_string(),
2798 ));
2799 }
2800 if !independent {
2803 return Err(Error::not_implemented(
2804 "a subquery that reads both sides of that join, written in the condition of a join \
2805 whose right side is lateral"
2806 .to_string(),
2807 ));
2808 }
2809 let mut node = self.add_node(Node::CrossProduct { left, right });
2810 for pending in pair {
2811 node = self.attach_subquery(node, pending);
2812 }
2813 let mut conditions = conditions.into_iter();
2817 let mut predicate = conditions.next().expect("a join condition was bound");
2818 for next in conditions {
2819 let children = self.plan.add_expr_list(&[predicate, next]);
2820 let conjunction = Expr::Conjunction { op: ConjunctionOp::And, children };
2821 predicate = self.plan.add_expr(conjunction, LogicalType::Boolean);
2822 }
2823 let node = self.add_node(Node::Filter { input: node, predicate });
2824 Ok((node, scope))
2825 }
2826
2827 #[allow(clippy::too_many_arguments)]
2828 fn bind_join(
2829 &mut self,
2830 ast: &Ast,
2831 left: ast::SourceRef,
2832 right: ast::SourceRef,
2833 kind: ast::JoinKind,
2834 natural: bool,
2835 on: ast::ExprRef,
2836 using: ast::Slice,
2837 ) -> Result<(NodeRef, Scope)> {
2838 let (left_node, left_scope) = self.bind_source(ast, left)?;
2839 let (right_node, right_scope, correlated) = self.bind_lateral(ast, right, &left_scope)?;
2840 if !correlated.is_empty()
2844 && !matches!(kind, ast::JoinKind::Inner | ast::JoinKind::Cross | ast::JoinKind::Left)
2845 {
2846 return Err(Error::binder(
2847 "The combining JOIN type must be INNER or LEFT for a LATERAL reference",
2848 ));
2849 }
2850 let split = left_scope.len();
2851 let left_tables: Vec<u32> =
2857 left_scope.columns.iter().map(|column| column.binding.table).collect();
2858 let right_tables: Vec<u32> =
2859 right_scope.columns.iter().map(|column| column.binding.table).collect();
2860 let mut scope = left_scope.concat(right_scope);
2861
2862 let merged: Vec<String> = if natural {
2865 let mut names = Vec::new();
2866 for (at, column) in scope.columns.iter().enumerate().take(split) {
2867 if scope.columns[split..].iter().any(|right| same_name(&right.name, &column.name))
2868 && !names.iter().any(|held: &String| same_name(held, &column.name))
2869 {
2870 let _ = at;
2871 names.push(column.name.clone());
2872 }
2873 }
2874 names
2875 } else {
2876 let mut names: Vec<String> = Vec::new();
2882 for name in ast.name(using) {
2883 if !names.iter().any(|held| same_name(held, name)) {
2884 names.push(name.to_string());
2885 }
2886 }
2887 names
2888 };
2889
2890 let mut conditions = Vec::new();
2891 let mut dropped = Vec::new();
2892 for name in &merged {
2893 let left_at = scope.columns[..split]
2894 .iter()
2895 .position(|column| same_name(&column.name, name))
2896 .ok_or_else(|| {
2897 Error::binder(format!(
2898 "column \"{name}\" specified in USING clause does not exist in left table"
2899 ))
2900 })?;
2901 let right_at = scope.columns[split..]
2902 .iter()
2903 .position(|column| same_name(&column.name, name))
2904 .map(|at| at + split)
2905 .ok_or_else(|| {
2906 Error::binder(format!(
2907 "column \"{name}\" specified in USING clause does not exist in right table"
2908 ))
2909 })?;
2910 let left_column = &scope.columns[left_at];
2911 let (left_binding, left_type) = (left_column.binding, left_column.ty.clone());
2912 let right_column = &scope.columns[right_at];
2913 let (right_binding, right_type) = (right_column.binding, right_column.ty.clone());
2914 let left_expr = self.plan.add_expr(Expr::Column(left_binding), left_type);
2915 let right_expr = self.plan.add_expr(Expr::Column(right_binding), right_type);
2916 conditions.push(self.compare(rudb_plan::CompareOp::Equal, left_expr, right_expr)?);
2917 dropped.push(right_at);
2918 }
2919 dropped.sort_unstable();
2922 for at in dropped.into_iter().rev() {
2923 scope.remove(at);
2924 }
2925
2926 let mut left_node = left_node;
2927 let mut right_node = right_node;
2928 let mut pair = Vec::new();
2929 if on != NONE {
2930 if !merged.is_empty() {
2931 return Err(Error::binder("a join cannot have both ON and USING"));
2932 }
2933 self.clause = "JOIN condition";
2934 let waiting = self.scalar_subqueries.len();
2935 let predicate = self.bind_expr(ast, on, &scope)?;
2936 conditions.push(self.as_boolean(predicate, "JOIN")?);
2937 for pending in self.scalar_subqueries.split_off(waiting) {
2938 match self.side_of(&pending, &left_tables, &right_tables) {
2939 Some(Side::Right) => right_node = self.attach_subquery(right_node, pending),
2940 Some(Side::Left) => left_node = self.attach_subquery(left_node, pending),
2941 None => pair.push(pending),
2942 }
2943 }
2944 }
2945
2946 if kind == ast::JoinKind::Cross && !conditions.is_empty() {
2947 return Err(Error::binder("a CROSS JOIN cannot have a condition"));
2948 }
2949 if !pair.is_empty() {
2950 return self.bind_pair_dependent_join(
2951 kind,
2952 correlated.is_empty(),
2953 left_node,
2954 right_node,
2955 pair,
2956 conditions,
2957 scope,
2958 );
2959 }
2960 if correlated.is_empty()
2964 && conditions.is_empty()
2965 && matches!(kind, ast::JoinKind::Cross | ast::JoinKind::Inner)
2966 {
2967 let node = self.add_node(Node::CrossProduct { left: left_node, right: right_node });
2968 return Ok((node, scope));
2969 }
2970 if matches!(kind, ast::JoinKind::Semi | ast::JoinKind::Anti) {
2979 scope.truncate(split);
2980 }
2981 let kind = match kind {
2982 ast::JoinKind::Inner | ast::JoinKind::Cross => JoinKind::Inner,
2983 ast::JoinKind::Left => JoinKind::Left,
2984 ast::JoinKind::Right => JoinKind::Right,
2985 ast::JoinKind::Full => JoinKind::Full,
2986 ast::JoinKind::Semi => JoinKind::Semi,
2987 ast::JoinKind::Anti => JoinKind::Anti,
2988 ast::JoinKind::Positional => JoinKind::Positional,
2989 };
2990 let conditions = self.plan.add_expr_list(&conditions);
2991 let node = if correlated.is_empty() {
2992 self.add_node(Node::Join {
2993 left: left_node,
2994 right: right_node,
2995 kind,
2996 conditions,
2997 build: BuildSide::default(),
2998 })
2999 } else {
3000 self.add_node(Node::DependentJoin {
3001 left: left_node,
3002 right: right_node,
3003 kind,
3004 conditions,
3005 })
3006 };
3007 Ok((node, scope))
3008 }
3009
3010 fn bind_filter(
3018 &mut self,
3019 ast: &Ast,
3020 filter: ast::ExprRef,
3021 scope: &Scope,
3022 ) -> Result<Option<ExprRef>> {
3023 if filter == NONE {
3024 return Ok(None);
3025 }
3026 let bound = self.bind_expr(ast, filter, scope)?;
3027 Ok(Some(self.checked_cast_to(bound, &LogicalType::Boolean, false)?))
3028 }
3029
3030 pub(crate) fn bind_aggregate(
3035 &mut self,
3036 ast: &Ast,
3037 name: &str,
3038 args: &[ast::ExprRef],
3039 distinct: bool,
3040 filter: ast::ExprRef,
3041 scope: &Scope,
3042 ) -> Result<ExprRef> {
3043 let frames = std::mem::take(&mut self.lambda_frames);
3044 let bound = self.bind_aggregate_over_rows(ast, name, args, distinct, filter, scope);
3045 self.lambda_frames = frames;
3046 bound
3047 }
3048
3049 fn bind_aggregate_over_rows(
3050 &mut self,
3051 ast: &Ast,
3052 name: &str,
3053 args: &[ast::ExprRef],
3054 distinct: bool,
3055 filter: ast::ExprRef,
3056 scope: &Scope,
3057 ) -> Result<ExprRef> {
3058 if self.in_filter {
3059 return Err(Error::binder("aggregate functions are not allowed in FILTER"));
3060 }
3061 if self.in_aggregate {
3062 return Err(Error::binder(format!(
3063 "aggregate function calls cannot be nested, and {name}() is inside one"
3064 )));
3065 }
3066 if self.aggregation.is_none() {
3067 return Err(Error::binder(format!(
3068 "aggregate function calls cannot be used in the {}",
3069 self.clause
3070 )));
3071 }
3072 self.in_aggregate = true;
3077 self.in_filter = true;
3078 let filter = self.bind_filter(ast, filter, scope);
3079 self.in_filter = false;
3080 self.in_aggregate = false;
3081 let filter = filter?;
3082
3083 self.in_aggregate = true;
3084 let mut bound = Vec::with_capacity(args.len());
3085 let mut failure = None;
3086 for &arg in args {
3087 match self.bind_expr(ast, arg, scope) {
3088 Ok(expr) => bound.push(expr),
3089 Err(error) => {
3090 failure = Some(error);
3091 break;
3092 }
3093 }
3094 }
3095 self.in_aggregate = false;
3096 if let Some(error) = failure {
3097 return Err(error);
3098 }
3099
3100 let types: Vec<LogicalType> =
3101 bound.iter().map(|&arg| self.plan.expr_type(arg).clone()).collect();
3102 let resolved = resolve(name, &types)?;
3103 if resolved.name == "string_agg"
3106 && bound.len() == 2
3107 && !matches!(fold::value_of(&self.plan, bound[1]), Ok(Some(_)))
3108 {
3109 return Err(Error::binder(
3110 "The \"separator\" argument in function \"string_agg\" must be a constant expression",
3111 ));
3112 }
3113 let mut cast = Vec::with_capacity(bound.len());
3114 for (arg, wanted) in bound.iter().zip(&resolved.arguments) {
3115 cast.push(self.checked_cast_to(*arg, wanted, false)?);
3116 }
3117 let args = self.plan.add_expr_list(&cast);
3118 let name = self.plan.intern(resolved.name);
3119 let ty = resolved.returns;
3120 let call = self.plan.add_expr(Expr::Aggregate { name, args, distinct, filter }, ty.clone());
3121
3122 let existing = self.aggregation.as_ref().map(|held| held.aggregates.clone());
3125 let existing = existing.unwrap_or_default();
3126 let at = match existing.iter().position(|&held| self.same_expr(held, call)) {
3127 Some(at) => at,
3128 None => {
3129 let aggregation = self.aggregation.as_mut().expect("checked above");
3130 aggregation.aggregates.push(call);
3131 aggregation.aggregates.len() - 1
3132 }
3133 };
3134 let aggregation = self.aggregation.as_ref().expect("checked above");
3135 let (index, groups) = (aggregation.index, aggregation.groups.len());
3136 Ok(self.column(index, groups + at, ty))
3137 }
3138
3139 pub(crate) fn bind_window(
3150 &mut self,
3151 ast: &Ast,
3152 written: &WindowCall<'_>,
3153 scope: &Scope,
3154 ) -> Result<ExprRef> {
3155 let frames = std::mem::take(&mut self.lambda_frames);
3156 let bound = self.bind_window_over_rows(ast, written, scope);
3157 self.lambda_frames = frames;
3158 bound
3159 }
3160
3161 fn bind_window_over_rows(
3162 &mut self,
3163 ast: &Ast,
3164 written: &WindowCall<'_>,
3165 scope: &Scope,
3166 ) -> Result<ExprRef> {
3167 let WindowCall { name, args, distinct, filter, ignore_nulls, spec, .. } = *written;
3168 if self.in_aggregate {
3169 return Err(Error::binder(
3170 "aggregate function calls cannot contain window function calls",
3171 ));
3172 }
3173 if self.in_window {
3174 return Err(Error::binder("window function calls cannot be nested"));
3175 }
3176 let clause = if self.clause == "JOIN condition" { "WHERE clause" } else { self.clause };
3180 if clause != "SELECT clause" && clause != "ORDER BY clause" {
3181 return Err(Error::binder(format!("{clause} cannot contain window functions!")));
3182 }
3183
3184 let starred = args.iter().any(|&arg| {
3188 matches!(ast.expr(arg), ast::Expr::Star { qualifier, replacements }
3189 if qualifier.is_empty() && replacements.is_empty())
3190 });
3191 let (name, args): (&str, &[ast::ExprRef]) = if starred {
3192 if !same_name(name, "count") || args.len() != 1 {
3193 return Err(Error::binder(format!("* is not allowed in {name}()")));
3194 }
3195 ("count_star", &[])
3196 } else if same_name(name, "count") && args.is_empty() {
3197 ("count_star", &[])
3200 } else {
3201 (name, args)
3202 };
3203
3204 let held = ast.window(spec);
3205 self.in_window = true;
3206 let parts = self.window_parts(ast, written, args, held, scope);
3207 let filter = if parts.is_ok() { self.bind_filter(ast, filter, scope) } else { Ok(None) };
3212 self.in_window = false;
3213 let parts = parts?;
3214 let filter = filter?;
3215 let offsets = [parts.frame.start, parts.frame.end]
3218 .iter()
3219 .any(|end| matches!(end, WindowBound::Preceding(_) | WindowBound::Following(_)));
3220 if parts.frame.unit == WindowUnit::Range && offsets && parts.order.len() != 1 {
3221 return Err(Error::binder("RANGE frames must have only one ORDER BY expression"));
3222 }
3223
3224 let types: Vec<LogicalType> =
3225 parts.args.iter().map(|&arg| self.plan.expr_type(arg).clone()).collect();
3226 let resolved = window_signature(name, &types)?;
3227 if resolved.name == "fill" {
3230 let keys: Vec<LogicalType> =
3231 parts.order.iter().map(|key| self.plan.expr_type(key.expr).clone()).collect();
3232 refuse_fill(&types[0], &keys, distinct, ignore_nulls)?;
3233 }
3234 if distinct && kind_of(resolved.name) == Some(FunctionKind::Window) {
3238 return Err(Error::binder(format!(
3239 "DISTINCT is not implemented for the window function \"\"{name}\"\""
3240 )));
3241 }
3242 if filter.is_some() && kind_of(resolved.name) == Some(FunctionKind::Window) {
3245 return Err(Error::binder(format!(
3246 "FILTER is not implemented for the window function \"\"{name}\"\""
3247 )));
3248 }
3249 if !parts.inner.is_empty() && kind_of(resolved.name) == Some(FunctionKind::Window) {
3257 let counts = matches!(resolved.name, "first_value" | "last_value" | "nth_value");
3258 if !counts {
3259 if parts.frame.exclude != WindowExclude::NoOthers {
3260 return Err(Error::binder(format!(
3261 "EXCLUDE is not supported for the window function \"\"{}\"\"",
3262 resolved.name
3263 )));
3264 }
3265 return Err(Error::not_implemented(format!(
3266 "ORDER BY inside the arguments of the window function \"{}\"",
3267 resolved.name
3268 )));
3269 }
3270 }
3271 let mut cast = Vec::with_capacity(parts.args.len());
3272 for (arg, wanted) in parts.args.iter().zip(&resolved.arguments) {
3273 cast.push(self.checked_cast_to(*arg, wanted, false)?);
3274 }
3275 let args = self.plan.add_expr_list(&cast);
3276 let order = self.plan.add_sort_keys(&parts.inner);
3277 let name = self.plan.intern(resolved.name);
3278 let ty = resolved.returns;
3279 let call = self.plan.add_expr(
3280 Expr::Window { name, args, distinct, filter, ignore_nulls, order },
3281 ty.clone(),
3282 );
3283
3284 let at = self.window_run(parts.partition, parts.order, parts.frame, call);
3285 let index = self.windows.last().expect("the run was just filed").index;
3286 Ok(self.column(index, at, ty))
3287 }
3288
3289 fn window_run(
3296 &mut self,
3297 partition: Vec<ExprRef>,
3298 order: Vec<SortKey>,
3299 frame: WindowFrame,
3300 call: ExprRef,
3301 ) -> usize {
3302 let matches = self.windows.last().is_some_and(|run| {
3303 run.frame == frame
3304 && run.partition.len() == partition.len()
3305 && run.order.len() == order.len()
3306 && run.partition.iter().zip(&partition).all(|(&l, &r)| self.same_expr(l, r))
3307 && run.order.iter().zip(&order).all(|(l, r)| {
3308 l.descending == r.descending
3309 && l.nulls_first == r.nulls_first
3310 && self.same_expr(l.expr, r.expr)
3311 })
3312 });
3313 if !matches {
3314 let index = self.fresh_index();
3315 self.windows.push(WindowRun { index, partition, order, frame, calls: Vec::new() });
3316 }
3317 let calls = self.windows.last().expect("a run is open").calls.clone();
3320 if let Some(at) = calls.iter().position(|&held| self.same_expr(held, call)) {
3321 return at;
3322 }
3323 let run = self.windows.last_mut().expect("a run is open");
3324 run.calls.push(call);
3325 run.calls.len() - 1
3326 }
3327
3328 fn window_parts(
3334 &mut self,
3335 ast: &Ast,
3336 written: &WindowCall<'_>,
3337 args: &[ast::ExprRef],
3338 held: ast::WindowSpec,
3339 scope: &Scope,
3340 ) -> Result<WindowParts> {
3341 let mut bound = Vec::with_capacity(args.len());
3342 for &arg in args {
3343 let expr = self.bind_expr(ast, arg, scope)?;
3344 bound.push(self.over_aggregate(expr, scope)?);
3345 }
3346 let mut inner = Vec::new();
3350 for item in ast.order_list(written.order).to_vec() {
3351 let expr = self.bind_expr(ast, item.expr, scope)?;
3352 let expr = self.over_aggregate(expr, scope)?;
3353 inner.push(self.sort_key(expr, item));
3354 }
3355 let mut partition = Vec::new();
3356 for &key in ast.expr_list(held.partition) {
3357 let expr = self.bind_expr(ast, key, scope)?;
3358 partition.push(self.over_aggregate(expr, scope)?);
3359 }
3360 let mut order = Vec::new();
3361 for item in ast.order_list(held.order).to_vec() {
3362 let expr = self.bind_expr(ast, item.expr, scope)?;
3363 let expr = self.over_aggregate(expr, scope)?;
3364 order.push(self.sort_key(expr, item));
3365 }
3366 let frame = WindowFrame {
3367 unit: match held.unit {
3368 ast::WindowUnit::Rows => WindowUnit::Rows,
3369 ast::WindowUnit::Range => WindowUnit::Range,
3370 ast::WindowUnit::Groups => WindowUnit::Groups,
3371 },
3372 start: self.window_bound(ast, held.start, scope)?,
3373 end: self.window_bound(ast, held.end, scope)?,
3374 exclude: match held.exclude {
3375 ast::WindowExclude::NoOthers => WindowExclude::NoOthers,
3376 ast::WindowExclude::CurrentRow => WindowExclude::CurrentRow,
3377 ast::WindowExclude::Group => WindowExclude::Group,
3378 ast::WindowExclude::Ties => WindowExclude::Ties,
3379 },
3380 };
3381 Ok(WindowParts { args: bound, partition, order, inner, frame })
3382 }
3383
3384 fn window_bound(
3386 &mut self,
3387 ast: &Ast,
3388 bound: ast::WindowBound,
3389 scope: &Scope,
3390 ) -> Result<WindowBound> {
3391 let offset = |binder: &mut Self, written| {
3392 let expr = binder.bind_expr(ast, written, scope)?;
3393 binder.over_aggregate(expr, scope)
3394 };
3395 Ok(match bound {
3396 ast::WindowBound::UnboundedPreceding => WindowBound::UnboundedPreceding,
3397 ast::WindowBound::CurrentRow => WindowBound::CurrentRow,
3398 ast::WindowBound::UnboundedFollowing => WindowBound::UnboundedFollowing,
3399 ast::WindowBound::Preceding(written) => WindowBound::Preceding(offset(self, written)?),
3400 ast::WindowBound::Following(written) => WindowBound::Following(offset(self, written)?),
3401 })
3402 }
3403
3404 fn group_of(&self, read: ColumnBinding) -> Option<usize> {
3410 self.aggregation.as_ref()?.groups.iter().position(
3411 |group| matches!(*self.plan.expr(*group), Expr::Column(binding) if binding == read),
3412 )
3413 }
3414
3415 fn ungrouped_correlation(&self, binding: ColumnBinding) -> Option<ColumnBinding> {
3421 let pending =
3422 self.scalar_subqueries.iter().find(|pending| pending.index == binding.table)?;
3423 pending.reads.iter().copied().find(|read| self.group_of(*read).is_none())
3424 }
3425
3426 fn is_window_output(&self, binding: ColumnBinding) -> bool {
3428 self.windows.iter().any(|run| run.index == binding.table)
3429 }
3430
3431 fn is_correlation(&self, binding: ColumnBinding) -> bool {
3437 self.correlations.last().is_some_and(|frame| frame.contains(&binding))
3438 }
3439
3440 fn name_of(&self, binding: ColumnBinding, scope: &Scope) -> String {
3446 std::iter::once(scope)
3447 .chain(self.outer_scopes.iter().rev())
3448 .flat_map(|visible| visible.columns.iter())
3449 .find(|column| column.binding == binding)
3450 .map_or_else(|| "a column".to_string(), |column| format!("\"{}\"", column.name))
3451 }
3452
3453 pub(crate) fn over_aggregate(&mut self, expr: ExprRef, scope: &Scope) -> Result<ExprRef> {
3459 let Some(aggregation) = self.aggregation.as_ref() else {
3460 return Ok(expr);
3461 };
3462 let index = aggregation.index;
3463 let groups = aggregation.groups.clone();
3464 for (at, group) in groups.iter().enumerate() {
3465 if self.same_expr(expr, *group) {
3466 let ty = self.plan.expr_type(*group).clone();
3467 return Ok(self.column(index, at, ty));
3468 }
3469 }
3470 let ty = self.plan.expr_type(expr).clone();
3471 match self.plan.expr(expr).clone() {
3472 Expr::Column(binding) if binding.table == index => Ok(expr),
3473 Expr::Column(binding) if self.is_window_output(binding) => Ok(expr),
3478 Expr::Column(binding) if self.joined_above.contains(&binding.table) => Ok(expr),
3483 Expr::Column(binding) if self.is_correlation(binding) => Ok(expr),
3489 Expr::Column(binding) => {
3498 let read = self.ungrouped_correlation(binding).unwrap_or(binding);
3499 let name = self.name_of(read, scope);
3500 Err(Error::binder(format!(
3501 "column {name} must appear in the GROUP BY clause or must be part of an aggregate function"
3502 )))
3503 }
3504 Expr::Constant(_)
3505 | Expr::Aggregate { .. }
3506 | Expr::Window { .. }
3507 | Expr::LambdaParam(_) => Ok(expr),
3508 Expr::Lambda { table, params, body } => {
3512 let body = self.over_aggregate(body, scope)?;
3513 Ok(self.plan.add_expr(Expr::Lambda { table, params, body }, ty))
3514 }
3515 Expr::Cast { input, try_cast } => {
3516 let input = self.over_aggregate(input, scope)?;
3517 Ok(self.plan.add_expr(Expr::Cast { input, try_cast }, ty))
3518 }
3519 Expr::Compare { op, left, right } => {
3520 let left = self.over_aggregate(left, scope)?;
3521 let right = self.over_aggregate(right, scope)?;
3522 Ok(self.plan.add_expr(Expr::Compare { op, left, right }, ty))
3523 }
3524 Expr::Conjunction { op, children } => {
3525 let written = self.plan.expr_list(children).to_vec();
3526 let mut rewritten = Vec::with_capacity(written.len());
3527 for child in written {
3528 rewritten.push(self.over_aggregate(child, scope)?);
3529 }
3530 let children = self.plan.add_expr_list(&rewritten);
3531 Ok(self.plan.add_expr(Expr::Conjunction { op, children }, ty))
3532 }
3533 Expr::Function { name, args } => {
3534 let written = self.plan.expr_list(args).to_vec();
3535 let mut rewritten = Vec::with_capacity(written.len());
3536 for arg in written {
3537 rewritten.push(self.over_aggregate(arg, scope)?);
3538 }
3539 let args = self.plan.add_expr_list(&rewritten);
3540 Ok(self.plan.add_expr(Expr::Function { name, args }, ty))
3541 }
3542 Expr::Case { arms, otherwise } => {
3543 let written = self.plan.arm_list(arms).to_vec();
3544 let mut rewritten = Vec::with_capacity(written.len());
3545 for arm in written {
3546 let when = self.over_aggregate(arm.when, scope)?;
3547 let then = self.over_aggregate(arm.then, scope)?;
3548 rewritten.push(rudb_plan::Arm { when, then });
3549 }
3550 let otherwise = match otherwise {
3551 Some(expr) => Some(self.over_aggregate(expr, scope)?),
3552 None => None,
3553 };
3554 let arms = self.plan.add_arms(&rewritten);
3555 Ok(self.plan.add_expr(Expr::Case { arms, otherwise }, ty))
3556 }
3557 }
3558 }
3559
3560 pub(crate) fn same_expr(&self, left: ExprRef, right: ExprRef) -> bool {
3562 same_expr(&self.plan, left, right)
3563 }
3564}
3565
3566#[derive(Debug, Default)]
3576struct Options {
3577 binary_as_string: bool,
3580 all_varchar: bool,
3582 file_row_number: bool,
3587 given: Given,
3589}
3590
3591impl Options {
3592 fn of(written: &[(&'static str, Value, ExprRef)]) -> Result<Self> {
3599 let mut options = Self::default();
3600 for (parameter, value, _) in written {
3601 match (*parameter, value) {
3602 ("binary_as_string", Value::Boolean(on)) => options.binary_as_string = *on,
3603 ("all_varchar", Value::Boolean(on)) => options.all_varchar = *on,
3604 ("file_row_number", Value::Boolean(on)) => options.file_row_number = *on,
3605 _ => {}
3606 }
3607 }
3608 let named: Vec<(&str, Value)> =
3609 written.iter().map(|(parameter, value, _)| (*parameter, value.clone())).collect();
3610 options.given = csv_given(&named)?;
3611 Ok(options)
3612 }
3613}
3614
3615fn mirror_target(paths: &[String]) -> Option<(String, FileStamp)> {
3618 let [path] = paths else { return None };
3619 let canonical = std::fs::canonicalize(path).ok()?;
3620 let stamp = FileStamp::of(&canonical)?;
3621 Some((canonical.to_str()?.to_string(), stamp))
3622}
3623
3624#[derive(Clone, Copy)]
3626struct Operator {
3627 op: SetOp,
3629 quantifier: Quantifier,
3631 by_name: bool,
3633}
3634
3635struct Merged {
3637 name: String,
3639 ty: LogicalType,
3641 left: Option<usize>,
3643 right: Option<usize>,
3645}
3646
3647fn match_by_position(left: &Scope, right: &Scope) -> Result<Vec<Merged>> {
3651 if left.len() != right.len() {
3652 return Err(Error::binder(format!(
3653 "Set operations can only apply to expressions with the same number of result columns, but left side has {} and right side has {}",
3654 left.len(),
3655 right.len()
3656 )));
3657 }
3658 let mut merged = Vec::with_capacity(left.len());
3659 for (at, (held, other)) in left.columns.iter().zip(&right.columns).enumerate() {
3660 merged.push(Merged {
3661 name: held.name.clone(),
3662 ty: meet(&held.ty, &other.ty)?,
3663 left: Some(at),
3664 right: Some(at),
3665 });
3666 }
3667 Ok(merged)
3668}
3669
3670fn match_by_name(left: &Scope, right: &Scope) -> Result<Vec<Merged>> {
3678 named_once(left)?;
3679 named_once(right)?;
3680 let mut merged = Vec::with_capacity(left.len() + right.len());
3681 for (at, held) in left.columns.iter().enumerate() {
3682 let other = right.columns.iter().position(|column| same_name(&column.name, &held.name));
3683 let ty = match other {
3684 Some(other) => meet(&held.ty, &right.columns[other].ty)?,
3685 None => held.ty.clone(),
3686 };
3687 merged.push(Merged { name: held.name.clone(), ty, left: Some(at), right: other });
3688 }
3689 for (at, held) in right.columns.iter().enumerate() {
3690 if left.columns.iter().any(|column| same_name(&column.name, &held.name)) {
3691 continue;
3692 }
3693 merged.push(Merged {
3694 name: held.name.clone(),
3695 ty: held.ty.clone(),
3696 left: None,
3697 right: Some(at),
3698 });
3699 }
3700 Ok(merged)
3701}
3702
3703fn named_once(scope: &Scope) -> Result<()> {
3709 for (at, held) in scope.columns.iter().enumerate() {
3710 if scope.columns[..at].iter().any(|column| same_name(&column.name, &held.name)) {
3711 return Err(Error::binder(format!(
3712 "UNION (ALL) BY NAME operation doesn't support duplicate names in the SELECT list - the name \"\"{}\"\" occurs multiple times",
3713 held.name
3714 )));
3715 }
3716 }
3717 Ok(())
3718}
3719
3720fn meet(left: &LogicalType, right: &LogicalType) -> Result<LogicalType> {
3722 left.promote(right).ok_or_else(|| {
3723 Error::binder(format!(
3724 "Cannot combine a column of type {left} with a column of type {right} in a set operation"
3725 ))
3726 })
3727}
3728
3729fn null_parameter(function: TableFunction, parameter: &str) -> String {
3738 match parameter {
3739 "header" => format!("\"{parameter}\" expects a non-null boolean value (e.g. TRUE or 1)"),
3740 "all_varchar" => format!("{} \"{parameter}\" cannot be NULL", function.name()),
3741 _ => format!("Cannot use NULL as argument to \"{parameter}\""),
3742 }
3743}
3744
3745fn missing_replacement(name: &str, input: &Scope) -> Error {
3750 Error::binder(format!(
3751 "Column \"{name}\" in REPLACE list not found in FROM clause{}",
3752 input.candidates()
3753 ))
3754}
3755
3756fn subtractable(ty: &LogicalType, ordering: bool) -> bool {
3765 if ty.is_numeric() {
3766 return true;
3767 }
3768 match ty {
3769 LogicalType::Date
3770 | LogicalType::Time
3771 | LogicalType::Timestamp
3772 | LogicalType::TimestampS
3773 | LogicalType::TimestampMs
3774 | LogicalType::TimestampNs
3775 | LogicalType::TimestampTz => true,
3776 LogicalType::TimeTz => ordering,
3777 _ => false,
3778 }
3779}
3780
3781fn refuse_fill(
3790 argument: &LogicalType,
3791 order: &[LogicalType],
3792 distinct: bool,
3793 ignore_nulls: bool,
3794) -> Result<()> {
3795 if !subtractable(argument, false) {
3796 return Err(Error::binder("FILL argument must support subtraction"));
3797 }
3798 let [key] = order else {
3799 return Err(Error::binder("FILL functions must have only one ORDER BY expression"));
3800 };
3801 if !subtractable(key, true) {
3802 return Err(Error::binder("FILL ordering must support subtraction"));
3803 }
3804 if distinct {
3805 return Err(Error::binder(
3806 "DISTINCT is not implemented for the window function \"\"fill\"\"",
3807 ));
3808 }
3809 if ignore_nulls {
3810 return Err(Error::binder(
3811 "RESPECT/IGNORE NULLS is not supported for the window function \"fill\"",
3812 ));
3813 }
3814 Ok(())
3815}
3816
3817fn window_signature(name: &str, types: &[LogicalType]) -> Result<Resolved> {
3824 match kind_of(name) {
3825 Some(FunctionKind::Aggregate | FunctionKind::Window) => resolve(name, types),
3826 Some(FunctionKind::Scalar) => {
3827 Err(Error::catalog(format!("{name} is not an aggregate function")))
3828 }
3829 None => Err(Error::catalog(format!("Aggregate Function with name {name} does not exist!"))),
3830 }
3831}
3832
3833fn same_expr(plan: &Plan, left: ExprRef, right: ExprRef) -> bool {
3835 if left == right {
3836 return true;
3837 }
3838 if plan.expr_type(left) != plan.expr_type(right) {
3839 return false;
3840 }
3841 let lists = |left, right| {
3842 let left: &[ExprRef] = plan.expr_list(left);
3843 let right: &[ExprRef] = plan.expr_list(right);
3844 left.len() == right.len()
3845 && left.iter().zip(right).all(|(&left, &right)| same_expr(plan, left, right))
3846 };
3847 match (plan.expr(left), plan.expr(right)) {
3848 (Expr::Column(left), Expr::Column(right)) => left == right,
3849 (Expr::Constant(left), Expr::Constant(right)) => plan.value(*left) == plan.value(*right),
3850 (
3851 Expr::Cast { input: left, try_cast: left_try },
3852 Expr::Cast { input: right, try_cast: right_try },
3853 ) => left_try == right_try && same_expr(plan, *left, *right),
3854 (
3855 Expr::Compare { op: left_op, left: left_a, right: left_b },
3856 Expr::Compare { op: right_op, left: right_a, right: right_b },
3857 ) => {
3858 left_op == right_op
3859 && same_expr(plan, *left_a, *right_a)
3860 && same_expr(plan, *left_b, *right_b)
3861 }
3862 (
3863 Expr::Conjunction { op: left_op, children: left_children },
3864 Expr::Conjunction { op: right_op, children: right_children },
3865 ) => left_op == right_op && lists(*left_children, *right_children),
3866 (
3867 Expr::Function { name: left_name, args: left_args },
3868 Expr::Function { name: right_name, args: right_args },
3869 ) => plan.string(*left_name) == plan.string(*right_name) && lists(*left_args, *right_args),
3870 (
3871 Expr::Aggregate {
3872 name: left_name,
3873 args: left_args,
3874 distinct: left_distinct,
3875 filter: left_filter,
3876 },
3877 Expr::Aggregate {
3878 name: right_name,
3879 args: right_args,
3880 distinct: right_distinct,
3881 filter: right_filter,
3882 },
3883 ) => {
3884 plan.string(*left_name) == plan.string(*right_name)
3885 && left_distinct == right_distinct
3886 && match (left_filter, right_filter) {
3887 (None, None) => true,
3888 (Some(left), Some(right)) => same_expr(plan, *left, *right),
3889 _ => false,
3890 }
3891 && lists(*left_args, *right_args)
3892 }
3893 (
3897 Expr::Window {
3898 name: left_name,
3899 args: left_args,
3900 distinct: left_distinct,
3901 filter: left_filter,
3902 ignore_nulls: left_nulls,
3903 order: left_order,
3904 },
3905 Expr::Window {
3906 name: right_name,
3907 args: right_args,
3908 distinct: right_distinct,
3909 filter: right_filter,
3910 ignore_nulls: right_nulls,
3911 order: right_order,
3912 },
3913 ) => {
3914 let left_keys = plan.sort_key_list(*left_order);
3917 let right_keys = plan.sort_key_list(*right_order);
3918 plan.string(*left_name) == plan.string(*right_name)
3919 && left_distinct == right_distinct
3920 && left_nulls == right_nulls
3921 && left_keys.len() == right_keys.len()
3922 && left_keys.iter().zip(right_keys).all(|(left, right)| {
3923 left.descending == right.descending
3924 && left.nulls_first == right.nulls_first
3925 && same_expr(plan, left.expr, right.expr)
3926 })
3927 && match (left_filter, right_filter) {
3928 (None, None) => true,
3929 (Some(left), Some(right)) => same_expr(plan, *left, *right),
3930 _ => false,
3931 }
3932 && lists(*left_args, *right_args)
3933 }
3934 (
3935 Expr::Case { arms: left_arms, otherwise: left_otherwise },
3936 Expr::Case { arms: right_arms, otherwise: right_otherwise },
3937 ) => {
3938 let left_arms = plan.arm_list(*left_arms);
3939 let right_arms = plan.arm_list(*right_arms);
3940 left_arms.len() == right_arms.len()
3941 && left_arms.iter().zip(right_arms).all(|(left, right)| {
3942 same_expr(plan, left.when, right.when) && same_expr(plan, left.then, right.then)
3943 })
3944 && match (left_otherwise, right_otherwise) {
3945 (None, None) => true,
3946 (Some(left), Some(right)) => same_expr(plan, *left, *right),
3947 _ => false,
3948 }
3949 }
3950 _ => false,
3951 }
3952}