1use std::collections::hash_map::Entry;
16use std::ops::Deref;
17
18use itertools::{EitherOrBoth, Itertools};
19use risingwave_common::bail;
20use risingwave_common::catalog::{Field, TableId};
21use risingwave_sqlparser::ast::{
22 AsOf, Expr as ParserExpr, FunctionArg, FunctionArgExpr, Ident, ObjectName, TableAlias,
23 TableFactor,
24};
25use thiserror::Error;
26use thiserror_ext::AsReport;
27
28use super::bind_context::ColumnBinding;
29use super::statement::RewriteExprsRecursive;
30use crate::binder::Binder;
31use crate::binder::bind_context::{BindingCte, BindingCteState};
32use crate::error::{ErrorCode, Result, RwError};
33use crate::expr::{ExprImpl, InputRef};
34
35mod gap_fill;
36mod join;
37mod match_recognize;
38mod share;
39mod subquery;
40mod table_function;
41mod table_or_source;
42mod watermark;
43mod window_table_function;
44
45pub use gap_fill::BoundGapFill;
46pub use join::BoundJoin;
47pub use match_recognize::{
48 BoundMatchRecognize, BoundMeasure, BoundSymbolDefinition, MeasureSlotKind,
49};
50pub use share::{BoundShare, BoundShareInput};
51pub use subquery::BoundSubquery;
52pub use table_or_source::{
53 BoundBaseTable, BoundIcebergMetadataTable, BoundSource, BoundSystemTable,
54};
55pub use watermark::BoundWatermark;
56pub use window_table_function::{BoundWindowTableFunction, WindowTableFunctionKind};
57
58use crate::expr::{CorrelatedId, Depth};
59
60#[derive(Debug, Clone)]
63pub enum Relation {
64 Source(Box<BoundSource>),
65 BaseTable(Box<BoundBaseTable>),
66 SystemTable(Box<BoundSystemTable>),
67 IcebergMetadataTable(Box<BoundIcebergMetadataTable>),
68 Subquery(Box<BoundSubquery>),
69 Join(Box<BoundJoin>),
70 Apply(Box<BoundJoin>),
71 WindowTableFunction(Box<BoundWindowTableFunction>),
72 TableFunction {
74 expr: ExprImpl,
75 with_ordinality: bool,
76 },
77 Watermark(Box<BoundWatermark>),
78 Share(Box<BoundShare>),
79 GapFill(Box<BoundGapFill>),
80 MatchRecognize(Box<BoundMatchRecognize>),
81}
82
83impl RewriteExprsRecursive for Relation {
84 fn rewrite_exprs_recursive(&mut self, rewriter: &mut impl crate::expr::ExprRewriter) {
85 match self {
86 Relation::Subquery(inner) => inner.rewrite_exprs_recursive(rewriter),
87 Relation::Join(inner) => inner.rewrite_exprs_recursive(rewriter),
88 Relation::Apply(inner) => inner.rewrite_exprs_recursive(rewriter),
89 Relation::WindowTableFunction(inner) => inner.rewrite_exprs_recursive(rewriter),
90 Relation::Watermark(inner) => inner.rewrite_exprs_recursive(rewriter),
91 Relation::Share(inner) => inner.rewrite_exprs_recursive(rewriter),
92 Relation::TableFunction { expr: inner, .. } => {
93 *inner = rewriter.rewrite_expr(inner.take())
94 }
95 Relation::MatchRecognize(inner) => {
96 inner.input.rewrite_exprs_recursive(rewriter);
97 for e in inner.exprs_mut() {
98 *e = rewriter.rewrite_expr(e.take());
99 }
100 }
101 _ => {}
102 }
103 }
104}
105
106impl Relation {
107 pub fn is_correlated_by_depth(&self, depth: Depth) -> bool {
108 match self {
109 Relation::Subquery(subquery) => subquery.query.is_correlated_by_depth(depth),
110 Relation::Join(join) => {
111 join.cond.has_correlated_input_ref_by_depth(depth)
112 || join.left.is_correlated_by_depth(depth)
113 || join.right.is_correlated_by_depth(depth)
114 }
115 Relation::Apply(join) => {
118 join.cond.has_correlated_input_ref_by_depth(depth)
119 || join.left.is_correlated_by_depth(depth)
120 || join.right.is_correlated_by_depth(depth + 1)
121 }
122 Relation::TableFunction {
123 expr: table_function,
124 with_ordinality: _,
125 } => table_function.has_correlated_input_ref_by_depth(depth + 1),
126 Relation::Share(share) => match &share.input {
127 BoundShareInput::Query(query) => query.is_correlated_by_depth(depth),
128 BoundShareInput::ChangeLog { relation, .. } => {
129 relation.is_correlated_by_depth(depth)
130 }
131 },
132 Relation::MatchRecognize(inner) => {
133 inner.input.is_correlated_by_depth(depth)
134 || inner
135 .exprs()
136 .any(|e| e.has_correlated_input_ref_by_depth(depth))
137 }
138 _ => false,
139 }
140 }
141
142 pub fn is_correlated_by_correlated_id(&self, correlated_id: CorrelatedId) -> bool {
143 match self {
144 Relation::Subquery(subquery) => {
145 subquery.query.is_correlated_by_correlated_id(correlated_id)
146 }
147 Relation::Join(join) | Relation::Apply(join) => {
148 join.cond
149 .has_correlated_input_ref_by_correlated_id(correlated_id)
150 || join.left.is_correlated_by_correlated_id(correlated_id)
151 || join.right.is_correlated_by_correlated_id(correlated_id)
152 }
153 Relation::TableFunction {
154 expr: table_function,
155 with_ordinality: _,
156 } => table_function.has_correlated_input_ref_by_correlated_id(correlated_id),
157 Relation::Share(share) => match &share.input {
158 BoundShareInput::Query(query) => {
159 query.is_correlated_by_correlated_id(correlated_id)
160 }
161 BoundShareInput::ChangeLog { relation, .. } => {
162 relation.is_correlated_by_correlated_id(correlated_id)
163 }
164 },
165 Relation::MatchRecognize(inner) => {
166 inner.input.is_correlated_by_correlated_id(correlated_id)
167 || inner
168 .exprs()
169 .any(|e| e.has_correlated_input_ref_by_correlated_id(correlated_id))
170 }
171 _ => false,
172 }
173 }
174
175 pub fn collect_correlated_indices_by_depth_and_assign_id(
176 &mut self,
177 depth: Depth,
178 correlated_id: CorrelatedId,
179 ) -> Vec<usize> {
180 match self {
181 Relation::Subquery(subquery) => subquery
182 .query
183 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
184 Relation::Join(join) => {
185 let mut correlated_indices = vec![];
186 correlated_indices.extend(
187 join.cond
188 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
189 );
190 correlated_indices.extend(
191 join.left
192 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
193 );
194 correlated_indices.extend(
195 join.right
196 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
197 );
198 correlated_indices
199 }
200 Relation::Apply(join) => {
201 let mut correlated_indices = vec![];
202 correlated_indices.extend(
203 join.cond
204 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
205 );
206 correlated_indices.extend(
207 join.left
208 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
209 );
210 correlated_indices.extend(
211 join.right
212 .collect_correlated_indices_by_depth_and_assign_id(
213 depth + 1,
214 correlated_id,
215 ),
216 );
217 correlated_indices
218 }
219 Relation::TableFunction {
220 expr: table_function,
221 with_ordinality: _,
222 } => table_function
223 .collect_correlated_indices_by_depth_and_assign_id(depth + 1, correlated_id),
224 Relation::Share(share) => {
225 match &mut share.input {
226 BoundShareInput::Query(query) => query
227 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
228 BoundShareInput::ChangeLog { relation, .. } => relation
229 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
230 }
231 }
232 Relation::MatchRecognize(inner) => {
233 let mut indices = inner
234 .input
235 .collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id);
236 for e in inner.exprs_mut() {
237 indices.extend(
238 e.collect_correlated_indices_by_depth_and_assign_id(depth, correlated_id),
239 );
240 }
241 indices
242 }
243 _ => vec![],
244 }
245 }
246}
247
248#[derive(Debug)]
249#[non_exhaustive]
250pub enum ResolveQualifiedNameErrorKind {
251 QualifiedNameTooLong,
252 NotCurrentDatabase,
253}
254
255#[derive(Debug, Error)]
256pub struct ResolveQualifiedNameError {
257 qualified: String,
258 kind: ResolveQualifiedNameErrorKind,
259}
260
261impl std::fmt::Display for ResolveQualifiedNameError {
262 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
263 match self.kind {
264 ResolveQualifiedNameErrorKind::QualifiedNameTooLong => write!(
265 f,
266 "improper qualified name (too many dotted names): {}",
267 self.qualified
268 ),
269 ResolveQualifiedNameErrorKind::NotCurrentDatabase => write!(
270 f,
271 "cross-database references are not implemented: \"{}\"",
272 self.qualified
273 ),
274 }
275 }
276}
277
278impl ResolveQualifiedNameError {
279 pub fn new(qualified: String, kind: ResolveQualifiedNameErrorKind) -> Self {
280 Self { qualified, kind }
281 }
282}
283
284impl From<ResolveQualifiedNameError> for RwError {
285 fn from(e: ResolveQualifiedNameError) -> Self {
286 ErrorCode::InvalidInputSyntax(format!("{}", e.as_report())).into()
287 }
288}
289
290impl Binder {
291 pub fn resolve_schema_qualified_name(
293 db_name: &str,
294 name: &ObjectName,
295 ) -> std::result::Result<(Option<String>, String), ResolveQualifiedNameError> {
296 let formatted_name = name.to_string();
297 let mut identifiers = name.0.clone();
298
299 if identifiers.len() > 3 {
300 return Err(ResolveQualifiedNameError::new(
301 formatted_name,
302 ResolveQualifiedNameErrorKind::QualifiedNameTooLong,
303 ));
304 }
305
306 let name = identifiers.pop().unwrap().real_value();
307
308 let schema_name = identifiers.pop().map(|ident| ident.real_value());
309 let database_name = identifiers.pop().map(|ident| ident.real_value());
310
311 if let Some(database_name) = database_name
312 && database_name != db_name
313 {
314 return Err(ResolveQualifiedNameError::new(
315 formatted_name,
316 ResolveQualifiedNameErrorKind::NotCurrentDatabase,
317 ));
318 }
319
320 Ok((schema_name, name))
321 }
322
323 pub fn validate_cross_db_reference(
325 db_name: &str,
326 name: &ObjectName,
327 ) -> std::result::Result<(), ResolveQualifiedNameError> {
328 let formatted_name = name.to_string();
329 let identifiers = &name.0;
330 if identifiers.len() > 3 {
331 return Err(ResolveQualifiedNameError::new(
332 formatted_name,
333 ResolveQualifiedNameErrorKind::QualifiedNameTooLong,
334 ));
335 }
336
337 if identifiers.len() == 3 && identifiers[0].real_value() != db_name {
338 return Err(ResolveQualifiedNameError::new(
339 formatted_name,
340 ResolveQualifiedNameErrorKind::NotCurrentDatabase,
341 ));
342 }
343
344 Ok(())
345 }
346
347 pub fn resolve_db_schema_qualified_name(
349 name: &ObjectName,
350 ) -> std::result::Result<(Option<String>, Option<String>, String), ResolveQualifiedNameError>
351 {
352 let formatted_name = name.to_string();
353 let mut identifiers = name.0.clone();
354
355 if identifiers.len() > 3 {
356 return Err(ResolveQualifiedNameError::new(
357 formatted_name,
358 ResolveQualifiedNameErrorKind::QualifiedNameTooLong,
359 ));
360 }
361
362 let name = identifiers.pop().unwrap().real_value();
363 let schema_name = identifiers.pop().map(|ident| ident.real_value());
364 let database_name = identifiers.pop().map(|ident| ident.real_value());
365
366 Ok((database_name, schema_name, name))
367 }
368
369 fn resolve_single_name(mut identifiers: Vec<Ident>, ident_desc: &str) -> Result<String> {
371 if identifiers.len() > 1 {
372 bail!("{} must contain 1 argument", ident_desc);
373 }
374 let name = identifiers.pop().unwrap().real_value();
375
376 Ok(name)
377 }
378
379 pub fn resolve_database_name(name: ObjectName) -> Result<String> {
381 Self::resolve_single_name(name.0, "database name")
382 }
383
384 pub fn resolve_schema_name(name: ObjectName) -> Result<String> {
386 Self::resolve_single_name(name.0, "schema name")
387 }
388
389 pub fn resolve_index_name(name: ObjectName) -> Result<String> {
391 Self::resolve_single_name(name.0, "index name")
392 }
393
394 pub fn resolve_view_name(name: ObjectName) -> Result<String> {
396 Self::resolve_single_name(name.0, "view name")
397 }
398
399 pub fn resolve_sink_name(name: ObjectName) -> Result<String> {
401 Self::resolve_single_name(name.0, "sink name")
402 }
403
404 pub fn resolve_subscription_name(name: ObjectName) -> Result<String> {
406 Self::resolve_single_name(name.0, "subscription name")
407 }
408
409 pub fn resolve_table_name(name: ObjectName) -> Result<String> {
411 Self::resolve_single_name(name.0, "table name")
412 }
413
414 pub fn resolve_source_name(name: ObjectName) -> Result<String> {
416 Self::resolve_single_name(name.0, "source name")
417 }
418
419 pub fn resolve_user_name(name: ObjectName) -> Result<String> {
421 Self::resolve_single_name(name.0, "user name")
422 }
423
424 pub(super) fn bind_table_to_context(
426 &mut self,
427 columns: impl IntoIterator<Item = (bool, Field)>, table_name: String,
429 schema_name: Option<String>,
430 alias: Option<&TableAlias>,
431 ) -> Result<()> {
432 const EMPTY: [Ident; 0] = [];
433 let (resolved_schema_name, table_name, column_aliases, table_alias) = match alias {
434 None => (schema_name.clone(), table_name, &EMPTY[..], None),
435 Some(TableAlias { name, columns }) => (
436 None,
437 name.real_value(),
438 columns.as_slice(),
439 Some(table_name),
440 ),
441 };
442
443 let num_col_aliases = column_aliases.len();
444
445 let begin = self.context.columns.len();
446 let mut alias_iter = column_aliases.iter().fuse();
449 let mut index = 0;
450 columns.into_iter().for_each(|(is_hidden, mut field)| {
451 let name = match is_hidden {
452 true => field.name.clone(),
453 false => alias_iter
454 .next()
455 .map(|t| t.real_value())
456 .unwrap_or_else(|| field.name.clone()),
457 };
458 field.name.clone_from(&name);
459 self.context.columns.push(ColumnBinding::new(
460 table_name.clone(),
461 schema_name.clone(),
462 table_alias.clone(),
463 begin + index,
464 is_hidden,
465 field,
466 ));
467 self.context
468 .indices_of
469 .entry(name)
470 .or_default()
471 .push(self.context.columns.len() - 1);
472 index += 1;
473 });
474
475 let num_cols = index;
476 if num_cols < num_col_aliases {
477 return Err(ErrorCode::BindError(format!(
478 "table \"{table_name}\" has {num_cols} columns available but {num_col_aliases} column aliases specified",
479 ))
480 .into());
481 }
482
483 match self
484 .context
485 .range_of
486 .entry((resolved_schema_name, table_name.clone()))
487 {
488 Entry::Occupied(_) => Err(ErrorCode::InternalError(format!(
489 "Duplicated table name while binding table to context: {}",
490 table_name
491 ))
492 .into()),
493 Entry::Vacant(entry) => {
494 entry.insert((begin, self.context.columns.len()));
495 Ok(())
496 }
497 }
498 }
499
500 pub fn bind_relation_by_name(
505 &mut self,
506 name: &ObjectName,
507 alias: Option<&TableAlias>,
508 as_of: Option<&AsOf>,
509 allow_cross_db: bool,
510 ) -> Result<Relation> {
511 let (db_name, schema_name, table_name) = if allow_cross_db {
512 Self::resolve_db_schema_qualified_name(name)?
513 } else {
514 let (schema_name, table_name) =
515 Self::resolve_schema_qualified_name(&self.db_name, name)?;
516 (None, schema_name, table_name)
517 };
518
519 if schema_name.is_none()
520 && let Some(item) = self.context.cte_to_relation.get(&table_name)
522 {
523 if as_of.is_some() {
526 return Err(ErrorCode::BindError(
527 "Right table of a temporal join should not be a CTE. \
528 It should be a table, index, or materialized view"
529 .to_owned(),
530 )
531 .into());
532 }
533
534 let BindingCte {
535 share_id,
536 state: cte_state,
537 alias: mut original_alias,
538 } = item.deref().borrow().clone();
539
540 debug_assert_eq!(original_alias.name.real_value(), table_name);
542
543 if let Some(from_alias) = alias {
544 original_alias.name = from_alias.name.clone();
545 original_alias.columns = original_alias
546 .columns
547 .into_iter()
548 .zip_longest(from_alias.columns.iter().cloned())
549 .map(EitherOrBoth::into_right)
550 .collect();
551 }
552
553 let exposed_table_name = original_alias.name.real_value();
554 self.context
555 .check_relation_name_conflict(&exposed_table_name)?;
556
557 match cte_state {
558 BindingCteState::Bound { query } => {
559 let input = BoundShareInput::Query(query);
560 self.bind_table_to_context(
561 input.fields()?,
562 table_name,
563 None,
564 Some(&original_alias),
565 )?;
566 self.context.add_cte_name(exposed_table_name);
567 Ok(Relation::Share(Box::new(BoundShare { share_id, input })))
570 }
571 BindingCteState::ChangeLog { table, key_indices } => {
572 let input = BoundShareInput::ChangeLog {
573 relation: table,
574 key_indices,
575 };
576 self.bind_table_to_context(
577 input.fields()?,
578 table_name,
579 None,
580 Some(&original_alias),
581 )?;
582 self.context.add_cte_name(exposed_table_name);
583 Ok(Relation::Share(Box::new(BoundShare { share_id, input })))
584 }
585 }
586 } else {
587 let exposed_table_name = alias
588 .map(|alias| alias.name.real_value())
589 .unwrap_or_else(|| table_name.clone());
590 self.context.check_catalog_name(&exposed_table_name)?;
591 self.bind_catalog_relation_by_name(
592 db_name.as_deref(),
593 schema_name.as_deref(),
594 &table_name,
595 alias,
596 as_of,
597 false,
598 )
599 }
600 }
601
602 fn bind_relation_by_function_arg(
604 &mut self,
605 arg: Option<&FunctionArg>,
606 err_msg: &str,
607 ) -> Result<(Relation, ObjectName)> {
608 let Some(FunctionArg::Unnamed(FunctionArgExpr::Expr(expr))) = arg else {
609 return Err(ErrorCode::BindError(err_msg.to_owned()).into());
610 };
611 let table_name = match expr {
612 ParserExpr::Identifier(ident) => Ok::<_, RwError>(ObjectName(vec![ident.clone()])),
613 ParserExpr::CompoundIdentifier(idents) => Ok(ObjectName(idents.clone())),
614 _ => Err(ErrorCode::BindError(err_msg.to_owned()).into()),
615 }?;
616
617 Ok((
618 self.bind_relation_by_name(&table_name, None, None, true)?,
619 table_name,
620 ))
621 }
622
623 fn bind_column_by_function_args(
625 &mut self,
626 arg: Option<&FunctionArg>,
627 err_msg: &str,
628 ) -> Result<Box<InputRef>> {
629 if let Some(time_col_arg) = arg
630 && let Some(ExprImpl::InputRef(time_col)) =
631 self.bind_function_arg(time_col_arg)?.into_iter().next()
632 {
633 Ok(time_col)
634 } else {
635 Err(ErrorCode::BindError(err_msg.to_owned()).into())
636 }
637 }
638
639 fn bind_internal_table(
641 &mut self,
642 args: &[FunctionArg],
643 alias: Option<&TableAlias>,
644 ) -> Result<Relation> {
645 if args.is_empty() || args.len() > 2 {
646 return Err(
647 ErrorCode::BindError("usage: rw_table(table_id[,schema_name])".to_owned()).into(),
648 );
649 }
650
651 let table_id: TableId = args[0]
652 .to_string()
653 .parse::<u32>()
654 .map_err(|err| {
655 RwError::from(ErrorCode::BindError(format!(
656 "invalid table id: {}",
657 err.as_report()
658 )))
659 })?
660 .into();
661
662 let schema = args.get(1).map(|arg| arg.to_string());
663
664 let table_name = self.catalog.get_table_name_by_id(table_id)?;
665 self.bind_catalog_relation_by_name(None, schema.as_deref(), &table_name, alias, None, false)
666 }
667
668 pub(super) fn bind_table_factor(&mut self, table_factor: &TableFactor) -> Result<Relation> {
669 match table_factor {
670 TableFactor::Table { name, alias, as_of } => {
671 self.bind_relation_by_name(name, alias.as_ref(), as_of.as_ref(), true)
672 }
673 TableFactor::TableFunction {
674 name,
675 alias,
676 args,
677 with_ordinality,
678 } => {
679 let visibility = self.mark_lateral_contexts_visible();
680 let result = self.bind_table_function(name, alias.as_ref(), args, *with_ordinality);
681 self.restore_lateral_contexts_visibility(visibility);
682 result
683 }
684 TableFactor::Derived {
685 lateral,
686 subquery,
687 alias,
688 } => {
689 if *lateral {
690 let visibility = self.mark_lateral_contexts_visible();
691
692 let result = self.bind_subquery_relation(subquery, alias.as_ref(), true);
694
695 self.restore_lateral_contexts_visibility(visibility);
696 result.map(|subquery| Relation::Subquery(Box::new(subquery)))
697 } else {
698 self.push_lateral_context();
700 let bound_subquery =
701 self.bind_subquery_relation(subquery, alias.as_ref(), false)?;
702 self.pop_and_merge_lateral_context()?;
703 Ok(Relation::Subquery(Box::new(bound_subquery)))
704 }
705 }
706 TableFactor::NestedJoin(table_with_joins) => {
707 self.push_lateral_context();
708 let bound_join = self.bind_table_with_joins(table_with_joins)?;
709 self.pop_and_merge_lateral_context()?;
710 Ok(bound_join)
711 }
712 TableFactor::MatchRecognize {
713 table,
714 partition_by,
715 order_by,
716 measures,
717 rows_per_match,
718 after_match_skip,
719 pattern,
720 within,
721 subsets,
722 symbols,
723 alias,
724 } => Ok(Relation::MatchRecognize(Box::new(
725 self.bind_match_recognize(
726 table,
727 partition_by,
728 order_by,
729 measures,
730 rows_per_match,
731 after_match_skip,
732 pattern,
733 within,
734 subsets,
735 symbols,
736 alias.as_ref(),
737 )?,
738 ))),
739 }
740 }
741}