Skip to main content

risingwave_frontend/planner/
relation.rs

1// Copyright 2022 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashMap;
16use std::ops::Deref;
17use std::rc::Rc;
18
19use iceberg::spec::{Operation, TableMetadata};
20use itertools::Itertools;
21use risingwave_common::bail_not_implemented;
22use risingwave_common::catalog::{
23    CdcTableDesc, ColumnCatalog, Engine, Field, RISINGWAVE_ICEBERG_COMMIT_EPOCH,
24    RISINGWAVE_ICEBERG_ROW_ID, ROW_ID_COLUMN_NAME, Schema,
25};
26use risingwave_common::constants::log_store::{
27    EPOCH_COLUMN_NAME, INSERT_OP_CODE, ROW_OP_COLUMN_NAME, encode_epoch,
28};
29use risingwave_common::session_config::IcebergQueryStorageMode;
30use risingwave_common::types::{DataType, Interval, ScalarImpl};
31use risingwave_common::util::iter_util::ZipEqFast;
32use risingwave_connector::WithOptionsSecResolved;
33use risingwave_connector::source::ConnectorProperties;
34use risingwave_connector::source::cdc::CdcScanOptions;
35use risingwave_connector::source::iceberg::IcebergTimeTravelInfo;
36use risingwave_sqlparser::ast::AsOf;
37use thiserror_ext::AsReport;
38
39use crate::TableCatalog;
40use crate::binder::{
41    BoundBaseTable, BoundGapFill, BoundIcebergMetadataTable, BoundJoin, BoundMatchRecognize,
42    BoundShare, BoundShareInput, BoundSource, BoundSystemTable, BoundWatermark,
43    BoundWindowTableFunction, Relation, WindowTableFunctionKind,
44};
45use crate::catalog::source_catalog::SourceCatalog;
46use crate::error::{ErrorCode, Result};
47use crate::expr::{CastContext, Expr, ExprImpl, ExprType, FunctionCall, InputRef, Literal};
48use crate::handler::cdc::derive_with_options_for_cdc_table;
49use crate::optimizer::IcebergSnapshotInfo;
50use crate::optimizer::plan_node::generic::{self, GenericPlanRef, SourceNodeKind};
51use crate::optimizer::plan_node::utils::to_iceberg_time_travel_as_of;
52use crate::optimizer::plan_node::{
53    LogicalApply, LogicalCdcScan, LogicalGapFill, LogicalHopWindow, LogicalIcebergIntermediateScan,
54    LogicalIcebergMetadataScan, LogicalJoin, LogicalMatchRecognize, LogicalPlanRef as PlanRef,
55    LogicalProject, LogicalScan, LogicalShare, LogicalSource, LogicalSysScan, LogicalTableFunction,
56    LogicalUnion, LogicalValues,
57};
58use crate::optimizer::property::Cardinality;
59use crate::planner::{PlanFor, Planner};
60use crate::utils::{ColIndexMapping, Condition};
61
62const ERROR_WINDOW_SIZE_ARG: &str =
63    "The size arg of window table function should be an interval literal.";
64
65impl Planner {
66    pub fn plan_relation(&mut self, relation: Relation) -> Result<PlanRef> {
67        match relation {
68            Relation::BaseTable(t) => self.plan_base_table(&t),
69            Relation::SystemTable(st) => self.plan_sys_table(*st),
70            Relation::IcebergMetadataTable(table) => self.plan_iceberg_metadata_table(*table),
71            // TODO: order is ignored in the subquery
72            Relation::Subquery(q) => Ok(self.plan_query(q.query)?.into_unordered_subplan()),
73            Relation::Join(join) => self.plan_join(*join),
74            Relation::Apply(join) => self.plan_apply(*join),
75            Relation::WindowTableFunction(tf) => self.plan_window_table_function(*tf),
76            Relation::Source(s) => self.plan_source(*s),
77            Relation::TableFunction {
78                expr: tf,
79                with_ordinality,
80            } => self.plan_table_function(tf, with_ordinality),
81            Relation::Watermark(tf) => self.plan_watermark(*tf),
82            Relation::Share(share) => self.plan_share(*share),
83            Relation::GapFill(bound_gap_fill) => self.plan_gap_fill(*bound_gap_fill),
84            Relation::MatchRecognize(mr) => self.plan_match_recognize(*mr),
85        }
86    }
87
88    pub(crate) fn plan_sys_table(&mut self, sys_table: BoundSystemTable) -> Result<PlanRef> {
89        Ok(LogicalSysScan::create(
90            sys_table.sys_table_catalog,
91            self.ctx(),
92            Cardinality::unknown(), // TODO(card): cardinality of system table
93        )
94        .into())
95    }
96
97    fn plan_iceberg_metadata_table(&mut self, table: BoundIcebergMetadataTable) -> Result<PlanRef> {
98        let timezone = self.ctx().get_session_timezone();
99        let time_travel_info = to_iceberg_time_travel_as_of(&table.as_of, &timezone)?;
100        let core = crate::optimizer::plan_node::generic::IcebergMetadataScan {
101            metadata_type: table.metadata_type,
102            properties: table.properties,
103            secret_refs: table.secret_refs,
104            time_travel_info,
105            ctx: self.ctx(),
106        };
107        Ok(LogicalIcebergMetadataScan::new(core).into())
108    }
109
110    pub(super) fn plan_base_table(&mut self, base_table: &BoundBaseTable) -> Result<PlanRef> {
111        let as_of = base_table.as_of.clone();
112        let scan = LogicalScan::from_base_table(base_table, self.ctx(), as_of.clone());
113
114        match base_table.table_catalog.engine {
115            Engine::Hummock => {
116                match as_of {
117                    None
118                    | Some(AsOf::ProcessTime)
119                    | Some(AsOf::ProcessTimeBroadcast)
120                    | Some(AsOf::TimestampNum(_))
121                    | Some(AsOf::TimestampString(_))
122                    | Some(AsOf::ProcessTimeWithInterval(_)) => {}
123                    Some(AsOf::VersionNum(_)) | Some(AsOf::VersionString(_)) => {
124                        bail_not_implemented!("As Of Version is not supported yet.")
125                    }
126                };
127                Ok(scan.into())
128            }
129            Engine::Iceberg => self.plan_iceberg_table(base_table, scan, as_of),
130        }
131    }
132
133    fn plan_iceberg_table(
134        &mut self,
135        base_table: &BoundBaseTable,
136        scan: LogicalScan,
137        as_of: Option<AsOf>,
138    ) -> Result<PlanRef> {
139        let is_append_only = base_table.table_catalog.append_only;
140        let iceberg_query_storage_mode = self
141            .ctx()
142            .session_ctx()
143            .config()
144            .iceberg_query_storage_mode();
145
146        enum PlanTarget {
147            TableScan,
148            Source,
149            IntermediateScan,
150        }
151        let plan_target = match self.plan_for() {
152            PlanFor::StreamIcebergEngineInternal => PlanTarget::TableScan,
153            PlanFor::BatchDql => match iceberg_query_storage_mode {
154                // Append-only Iceberg engine tables use a dummy Hummock materialization and keep
155                // their rows in Iceberg plus the sink log store.
156                IcebergQueryStorageMode::Hummock if !is_append_only => PlanTarget::TableScan,
157                _ => PlanTarget::IntermediateScan,
158            },
159            PlanFor::Stream => {
160                if is_append_only {
161                    PlanTarget::Source
162                } else {
163                    PlanTarget::TableScan
164                }
165            }
166            PlanFor::Batch => {
167                if is_append_only {
168                    PlanTarget::IntermediateScan
169                } else {
170                    PlanTarget::TableScan
171                }
172            }
173        };
174        match as_of {
175            None
176            | Some(AsOf::VersionNum(_))
177            | Some(AsOf::TimestampString(_))
178            | Some(AsOf::TimestampNum(_)) => {}
179            Some(AsOf::ProcessTime)
180            | Some(AsOf::ProcessTimeBroadcast)
181            | Some(AsOf::ProcessTimeWithInterval(_)) => {
182                bail_not_implemented!("As Of ProcessTime() is not supported yet.")
183            }
184            Some(AsOf::VersionString(_)) => {
185                bail_not_implemented!("As Of Version is not supported yet.")
186            }
187        }
188
189        if matches!(plan_target, PlanTarget::TableScan) {
190            return Ok(scan.into());
191        }
192
193        let source_catalog = self.get_iceberg_source_by_table_catalog(&base_table.table_catalog)
194            .ok_or_else(|| {
195                ErrorCode::BindError(format!(
196                    "failed to plan an iceberg engine table: {}. Can't find the corresponding iceberg source. Maybe you need to recreate the table",
197                    base_table.table_catalog.name()
198                    ))
199            })?;
200
201        // Build type mapping: source column name → Hummock table type.
202        // This lets the intermediate scan output Hummock types directly,
203        // while the Iceberg scan path still adds explicit casts after
204        // materialization to match the expected output types.
205        let mut table_column_type_mapping = HashMap::new();
206        let table_column_map: HashMap<&str, &DataType> = base_table
207            .table_catalog
208            .columns
209            .iter()
210            .map(|c| (c.name.as_str(), &c.column_desc.data_type))
211            .collect();
212        for source_col in &source_catalog.columns {
213            let source_name = source_col.name();
214            let table_name = if source_name == RISINGWAVE_ICEBERG_ROW_ID {
215                ROW_ID_COLUMN_NAME
216            } else {
217                source_name
218            };
219            if let Some(&table_type) = table_column_map.get(table_name)
220                && source_col.column_desc.data_type != *table_type
221            {
222                table_column_type_mapping.insert(source_name.to_owned(), table_type.clone());
223            }
224        }
225
226        let column_map: HashMap<String, (usize, ColumnCatalog)> = source_catalog
227            .columns
228            .clone()
229            .into_iter()
230            .enumerate()
231            .map(|(i, column)| (column.name().to_owned(), (i, column)))
232            .collect();
233        // For intermediate scan, it will output Hummock types directly. But for source, it still
234        // outputs original types. So only when the plan target is source and the column type is
235        // different, we need to add cast.
236        let exprs = scan
237            .table()
238            .column_schema()
239            .fields()
240            .iter()
241            .map(|field| {
242                let source_filed_name = if field.name == ROW_ID_COLUMN_NAME {
243                    RISINGWAVE_ICEBERG_ROW_ID
244                } else {
245                    &field.name
246                };
247                if let Some((i, source_column)) = column_map.get(source_filed_name) {
248                    let input_type = &source_column.column_desc.data_type;
249                    if matches!(plan_target, PlanTarget::Source) && input_type != &field.data_type {
250                        let mut input_ref =
251                            ExprImpl::InputRef(InputRef::new(*i, input_type.clone()).into());
252                        FunctionCall::cast_mut(
253                            &mut input_ref,
254                            &field.data_type,
255                            CastContext::Explicit,
256                        )
257                        .unwrap();
258                        input_ref
259                    } else {
260                        ExprImpl::InputRef(InputRef::new(*i, field.data_type.clone()).into())
261                    }
262                } else {
263                    // fields like `_rw_timestamp`, would not be found in source.
264                    ExprImpl::Literal(Literal::new(None, field.data_type.clone()).into())
265                }
266            })
267            .collect_vec();
268
269        // Build source→table column index mapping for Hummock rewrite.
270        // Must be built before source_catalog is moved into Rc.
271        let table_col_index: HashMap<&str, usize> = base_table
272            .table_catalog
273            .columns
274            .iter()
275            .enumerate()
276            .map(|(i, c)| (c.name.as_str(), i))
277            .collect();
278        let source_to_table_mapping = ColIndexMapping::new(
279            source_catalog
280                .columns
281                .iter()
282                .map(|c| {
283                    let table_name = if c.name() == RISINGWAVE_ICEBERG_ROW_ID {
284                        ROW_ID_COLUMN_NAME
285                    } else {
286                        c.name()
287                    };
288                    table_col_index.get(table_name).copied()
289                })
290                .collect(),
291            base_table.table_catalog.columns.len(),
292        );
293
294        let logical_source = LogicalSource::with_catalog(
295            Rc::new(source_catalog),
296            SourceNodeKind::CreateMViewOrBatch,
297            self.ctx(),
298            as_of.clone(),
299        )?;
300        if matches!(plan_target, PlanTarget::Source) {
301            return Ok(LogicalProject::new(logical_source.into(), exprs).into());
302        }
303
304        // Pin Hummock before loading the latest Iceberg snapshot. If a sink commit races with
305        // planning, this ordering guarantees that the Iceberg snapshot plus the filtered log
306        // store still represents a coherent prefix of the table.
307        let read_pending_log_store =
308            matches!(self.plan_for(), PlanFor::BatchDql) && is_append_only && as_of.is_none();
309        if read_pending_log_store {
310            drop(self.ctx().session_ctx().pinned_snapshot());
311        }
312
313        let mut logical_iceberg_intermediate_scan = self.plan_iceberg_intermediate_scan(
314            &logical_source,
315            table_column_type_mapping,
316            source_to_table_mapping,
317        )?;
318
319        if read_pending_log_store {
320            let snapshot_info = self.fetch_current_snapshot_info(&logical_source)?;
321            // A table without a snapshot has no committed rows, so all log-store rows are
322            // pending. For an existing snapshot without a RisingWave epoch marker, fall back to
323            // Iceberg-only reads to avoid returning duplicates from legacy/external snapshots.
324            let committed_epoch = match snapshot_info {
325                None => Some(None),
326                Some(IcebergSnapshotInfo {
327                    commit_epoch: Some(epoch),
328                    ..
329                }) => Some(Some(epoch)),
330                Some(_) => None,
331            };
332            if let Some(committed_epoch) = committed_epoch
333                && let Some(log_store_scan) = self.plan_iceberg_log_store_scan(
334                    &base_table.table_catalog,
335                    logical_iceberg_intermediate_scan.schema(),
336                    &logical_source.core.column_catalog,
337                    committed_epoch,
338                )?
339            {
340                logical_iceberg_intermediate_scan = LogicalUnion::create(
341                    true,
342                    vec![logical_iceberg_intermediate_scan, log_store_scan],
343                );
344            }
345        }
346
347        Ok(LogicalProject::new(logical_iceberg_intermediate_scan, exprs).into())
348    }
349
350    pub(super) fn plan_source(&mut self, source: BoundSource) -> Result<PlanRef> {
351        if source.catalog.is_cdc_table_source() {
352            if !matches!(self.plan_for(), PlanFor::Stream) {
353                return Err(ErrorCode::NotSupported(
354                    "batch queries on CDC table sources are not supported".to_owned(),
355                    "Create a materialized view to consume the CDC table source".to_owned(),
356                )
357                .into());
358            }
359            if source.as_of.is_some() {
360                return Err(ErrorCode::NotSupported(
361                    "AS OF on CDC table sources is not supported".to_owned(),
362                    "Remove the AS OF clause".to_owned(),
363                )
364                .into());
365            }
366
367            let catalog_desc = source
368                .catalog
369                .info
370                .external_table
371                .as_ref()
372                .expect("checked by is_cdc_table_source");
373            let desc = CdcTableDesc::from_protobuf(catalog_desc)?;
374            let upstream_source = {
375                let session = self.ctx.session_ctx();
376                let catalog_reader = session.env().catalog_reader().read_guard();
377                catalog_reader
378                    .get_source_by_id_with_db(&session.database(), desc.source_id)?
379                    .clone()
380            };
381            let desc = resolve_current_cdc_table_desc(desc, &upstream_source.with_properties)?;
382            let scan = LogicalCdcScan::create(
383                source.catalog.name.clone(),
384                Rc::new(desc),
385                self.ctx(),
386                CdcScanOptions {
387                    disable_backfill: true,
388                    ..Default::default()
389                },
390            );
391            Ok(scan.into())
392        } else if source.is_shareable_cdc_connector() {
393            Err(ErrorCode::InternalError(
394                "Should not create MATERIALIZED VIEW or SELECT directly on shared CDC source. HINT: create TABLE from the source instead.".to_owned(),
395            )
396            .into())
397        } else {
398            let as_of = source.as_of.clone();
399            match as_of {
400                None
401                | Some(AsOf::VersionNum(_))
402                | Some(AsOf::TimestampString(_))
403                | Some(AsOf::TimestampNum(_)) => {}
404                Some(AsOf::ProcessTime)
405                | Some(AsOf::ProcessTimeBroadcast)
406                | Some(AsOf::ProcessTimeWithInterval(_)) => {
407                    bail_not_implemented!("As Of ProcessTime() is not supported yet.")
408                }
409                Some(AsOf::VersionString(_)) => {
410                    bail_not_implemented!("As Of Version is not supported yet.")
411                }
412            }
413            let is_iceberg = source.catalog.is_iceberg_connector();
414
415            // validate the source has pk. We raise an error here to avoid panic in expect_stream_key later
416            // for a nicer error message.
417            if matches!(self.plan_for(), PlanFor::Stream) {
418                let has_pk =
419                    source.catalog.row_id_index.is_some() || !source.catalog.pk_col_ids.is_empty();
420                if !has_pk {
421                    // in older version, iceberg source doesn't have row_id, thus may hit this
422                    // only iceberg should hit this.
423                    debug_assert!(is_iceberg);
424                    if is_iceberg {
425                        return Err(ErrorCode::BindError(format!(
426                        "Cannot create a stream job from an iceberg source without a primary key.\nThe iceberg source might be created in an older version of RisingWave. Please try recreating the source.\nSource: {:?}",
427                        source.catalog
428                    ))
429                    .into());
430                    } else {
431                        return Err(ErrorCode::BindError(format!(
432                            "Cannot create a stream job from a source without a primary key.
433This is a bug. We would appreciate a bug report at:
434https://github.com/risingwavelabs/risingwave/issues/new?labels=type%2Fbug&template=bug_report.yml
435
436source: {:?}",
437                            source.catalog
438                        ))
439                        .into());
440                    }
441                }
442            }
443
444            let source = LogicalSource::with_catalog(
445                Rc::new(source.catalog),
446                SourceNodeKind::CreateMViewOrBatch,
447                self.ctx(),
448                as_of,
449            )?;
450            if is_iceberg && !matches!(self.plan_for(), PlanFor::Stream) {
451                let num_cols = source.core.column_catalog.len();
452                let intermediate_scan = self.plan_iceberg_intermediate_scan(
453                    &source,
454                    HashMap::new(),
455                    ColIndexMapping::identity(num_cols),
456                )?;
457                Ok(intermediate_scan)
458            } else {
459                Ok(source.into())
460            }
461        }
462    }
463
464    pub(super) fn plan_join(&mut self, join: BoundJoin) -> Result<PlanRef> {
465        let left = self.plan_relation(join.left)?;
466        let right = self.plan_relation(join.right)?;
467        let join_type = join.join_type;
468        let on_clause = join.cond;
469        if on_clause.has_subquery() {
470            bail_not_implemented!("Subquery in join on condition");
471        } else {
472            Ok(LogicalJoin::create(left, right, join_type, on_clause))
473        }
474    }
475
476    pub(super) fn plan_apply(&mut self, mut join: BoundJoin) -> Result<PlanRef> {
477        let join_type = join.join_type;
478        let on_clause = join.cond;
479        if on_clause.has_subquery() {
480            bail_not_implemented!("Subquery in join on condition");
481        }
482
483        let correlated_id = self.ctx.next_correlated_id();
484        let correlated_indices = join
485            .right
486            .collect_correlated_indices_by_depth_and_assign_id(0, correlated_id);
487        let left = self.plan_relation(join.left)?;
488        let right = self.plan_relation(join.right)?;
489
490        Ok(LogicalApply::create(
491            left,
492            right,
493            join_type,
494            Condition::with_expr(on_clause),
495            correlated_id,
496            correlated_indices,
497            false,
498        ))
499    }
500
501    pub(super) fn plan_window_table_function(
502        &mut self,
503        table_function: BoundWindowTableFunction,
504    ) -> Result<PlanRef> {
505        use WindowTableFunctionKind::*;
506        match table_function.kind {
507            Tumble => self.plan_tumble_window(
508                table_function.input,
509                table_function.time_col,
510                table_function.args,
511            ),
512            Hop => self.plan_hop_window(
513                table_function.input,
514                table_function.time_col,
515                table_function.args,
516            ),
517        }
518    }
519
520    pub(super) fn plan_table_function(
521        &mut self,
522        table_function: ExprImpl,
523        with_ordinality: bool,
524    ) -> Result<PlanRef> {
525        // TODO: maybe we can unify LogicalTableFunction with LogicalValues
526        match table_function {
527            ExprImpl::TableFunction(tf) => {
528                Ok(LogicalTableFunction::new(*tf, with_ordinality, self.ctx()).into())
529            }
530            expr => {
531                let schema = Schema {
532                    // TODO: should be named
533                    fields: vec![Field::unnamed(expr.return_type())],
534                };
535                let expr_return_type = expr.return_type();
536                let root = LogicalValues::create(vec![vec![expr]], schema, self.ctx());
537                let input_ref = ExprImpl::from(InputRef::new(0, expr_return_type.clone()));
538                let mut exprs = if let DataType::Struct(st) = expr_return_type {
539                    st.iter()
540                        .enumerate()
541                        .map(|(i, (_, ty))| {
542                            let idx = ExprImpl::literal_int(i.try_into().unwrap());
543                            let args = vec![input_ref.clone(), idx];
544                            FunctionCall::new_unchecked(ExprType::Field, args, ty.clone()).into()
545                        })
546                        .collect()
547                } else {
548                    vec![input_ref]
549                };
550                if with_ordinality {
551                    exprs.push(ExprImpl::literal_bigint(1));
552                }
553                Ok(LogicalProject::create(root, exprs))
554            }
555        }
556    }
557
558    pub(super) fn plan_share(&mut self, share: BoundShare) -> Result<PlanRef> {
559        match share.input {
560            BoundShareInput::Query(query) => {
561                let id = share.share_id;
562                match self.share_cache.get(&id) {
563                    None => {
564                        let result = self.plan_query(query)?.into_unordered_subplan();
565                        let logical_share = LogicalShare::create(result);
566                        self.share_cache.insert(id, logical_share.clone());
567                        Ok(logical_share)
568                    }
569                    Some(result) => Ok(result.clone()),
570                }
571            }
572            BoundShareInput::ChangeLog {
573                relation,
574                key_indices,
575            } => {
576                let id = share.share_id;
577                let result = self.plan_changelog(relation, key_indices)?;
578                let logical_share = LogicalShare::create(result);
579                self.share_cache.insert(id, logical_share.clone());
580                Ok(logical_share)
581            }
582        }
583    }
584
585    pub(super) fn plan_watermark(&mut self, _watermark: BoundWatermark) -> Result<PlanRef> {
586        todo!("plan watermark");
587    }
588
589    pub(super) fn plan_gap_fill(&mut self, gap_fill: BoundGapFill) -> Result<PlanRef> {
590        let input = self.plan_relation(gap_fill.input)?;
591        Ok(LogicalGapFill::new(
592            input,
593            gap_fill.time_col,
594            gap_fill.interval,
595            gap_fill.fill_strategies,
596            gap_fill.partition_by_cols,
597        )
598        .into())
599    }
600
601    pub(super) fn plan_match_recognize(&mut self, mr: BoundMatchRecognize) -> Result<PlanRef> {
602        let input = self.plan_relation(mr.input)?;
603        Ok(LogicalMatchRecognize::new(
604            input,
605            mr.partition_by,
606            mr.order_by,
607            mr.measures,
608            mr.rows_per_match,
609            mr.after_match_skip,
610            mr.pattern,
611            mr.defines,
612            mr.within,
613            mr.within_deadline,
614        )
615        .into())
616    }
617
618    fn collect_col_data_types_for_tumble_window(relation: &Relation) -> Result<Vec<DataType>> {
619        let col_data_types = match relation {
620            Relation::Source(s) => s
621                .catalog
622                .columns
623                .iter()
624                .map(|col| col.data_type().clone())
625                .collect(),
626            Relation::BaseTable(t) => t
627                .table_catalog
628                .columns
629                .iter()
630                .map(|col| col.data_type().clone())
631                .collect(),
632            Relation::Subquery(q) => q.query.schema().data_types(),
633            Relation::Share(share) => share
634                .input
635                .fields()?
636                .into_iter()
637                .map(|(_, f)| f.data_type)
638                .collect(),
639            r => {
640                return Err(ErrorCode::BindError(format!(
641                    "Invalid input relation to tumble: {r:?}"
642                ))
643                .into());
644            }
645        };
646        Ok(col_data_types)
647    }
648
649    fn plan_tumble_window(
650        &mut self,
651        input: Relation,
652        time_col: InputRef,
653        args: Vec<ExprImpl>,
654    ) -> Result<PlanRef> {
655        let mut args = args.into_iter();
656        let col_data_types: Vec<_> = Self::collect_col_data_types_for_tumble_window(&input)?;
657
658        match (args.next(), args.next(), args.next()) {
659            (Some(window_size @ ExprImpl::Literal(_)), None, None) => {
660                let mut exprs = Vec::with_capacity(col_data_types.len() + 2);
661                for (idx, col_dt) in col_data_types.iter().enumerate() {
662                    exprs.push(InputRef::new(idx, col_dt.clone()).into());
663                }
664                let window_start: ExprImpl = FunctionCall::new(
665                    ExprType::TumbleStart,
666                    vec![ExprImpl::InputRef(Box::new(time_col)), window_size.clone()],
667                )?
668                .into();
669                // TODO: `window_end` may be optimized to avoid double calculation of
670                // `tumble_start`, or we can depends on common expression
671                // optimization.
672                let window_end =
673                    FunctionCall::new(ExprType::Add, vec![window_start.clone(), window_size])?
674                        .into();
675                exprs.push(window_start);
676                exprs.push(window_end);
677                let base = self.plan_relation(input)?;
678                let project = LogicalProject::create(base, exprs);
679                Ok(project)
680            }
681            (
682                Some(window_size @ ExprImpl::Literal(_)),
683                Some(window_offset @ ExprImpl::Literal(_)),
684                None,
685            ) => {
686                let mut exprs = Vec::with_capacity(col_data_types.len() + 2);
687                for (idx, col_dt) in col_data_types.iter().enumerate() {
688                    exprs.push(InputRef::new(idx, col_dt.clone()).into());
689                }
690                let window_start: ExprImpl = FunctionCall::new(
691                    ExprType::TumbleStart,
692                    vec![
693                        ExprImpl::InputRef(Box::new(time_col)),
694                        window_size.clone(),
695                        window_offset,
696                    ],
697                )?
698                .into();
699                // TODO: `window_end` may be optimized to avoid double calculation of
700                // `tumble_start`, or we can depends on common expression
701                // optimization.
702                let window_end =
703                    FunctionCall::new(ExprType::Add, vec![window_start.clone(), window_size])?
704                        .into();
705                exprs.push(window_start);
706                exprs.push(window_end);
707                let base = self.plan_relation(input)?;
708                let project = LogicalProject::create(base, exprs);
709                Ok(project)
710            }
711            _ => Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into()),
712        }
713    }
714
715    fn plan_hop_window(
716        &mut self,
717        input: Relation,
718        time_col: InputRef,
719        args: Vec<ExprImpl>,
720    ) -> Result<PlanRef> {
721        let input = self.plan_relation(input)?;
722        let mut args = args.into_iter();
723        let Some((ExprImpl::Literal(window_slide), ExprImpl::Literal(window_size))) =
724            args.next_tuple()
725        else {
726            return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into());
727        };
728
729        let Some(ScalarImpl::Interval(window_slide)) = *window_slide.get_data() else {
730            return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into());
731        };
732        let Some(ScalarImpl::Interval(window_size)) = *window_size.get_data() else {
733            return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into());
734        };
735
736        let window_offset = match (args.next(), args.next()) {
737            (Some(ExprImpl::Literal(window_offset)), None) => match *window_offset.get_data() {
738                Some(ScalarImpl::Interval(window_offset)) => window_offset,
739                _ => return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into()),
740            },
741            (None, None) => Interval::from_month_day_usec(0, 0, 0),
742            _ => return Err(ErrorCode::BindError(ERROR_WINDOW_SIZE_ARG.to_owned()).into()),
743        };
744
745        if !window_size.is_positive() || !window_slide.is_positive() {
746            return Err(ErrorCode::BindError(format!(
747                "window_size {} and window_slide {} must be positive",
748                window_size, window_slide
749            ))
750            .into());
751        }
752
753        if window_size.exact_div(&window_slide).is_none() {
754            return Err(ErrorCode::BindError(format!("Invalid arguments for HOP window function: window_size {} cannot be divided by window_slide {}",window_size, window_slide)).into());
755        }
756
757        Ok(LogicalHopWindow::create(
758            input,
759            time_col,
760            window_slide,
761            window_size,
762            window_offset,
763        ))
764    }
765
766    fn plan_iceberg_intermediate_scan(
767        &self,
768        source: &LogicalSource,
769        table_column_type_mapping: HashMap<String, DataType>,
770        source_to_table_mapping: ColIndexMapping,
771    ) -> Result<PlanRef> {
772        // If time travel is not specified, we use current timestamp to get the latest snapshot
773        let timezone = self.ctx().get_session_timezone();
774        let mut time_travel_info = to_iceberg_time_travel_as_of(&source.core.as_of, &timezone)?;
775        if time_travel_info.is_none() {
776            time_travel_info = self
777                .fetch_current_snapshot_info(source)?
778                .map(|info| IcebergTimeTravelInfo::Version(info.snapshot_id));
779        }
780        let Some(time_travel_info) = time_travel_info else {
781            let mut schema = source.schema().clone();
782            for field in &mut schema.fields {
783                if let Some(target_type) = table_column_type_mapping.get(&field.name) {
784                    field.data_type = target_type.clone();
785                }
786            }
787            return Ok(LogicalValues::new(vec![], schema, self.ctx()).into());
788        };
789        let intermediate_scan = LogicalIcebergIntermediateScan::new(
790            source,
791            time_travel_info,
792            table_column_type_mapping,
793            source_to_table_mapping,
794        );
795        Ok(intermediate_scan.into())
796    }
797
798    fn fetch_current_snapshot_info(
799        &self,
800        source: &LogicalSource,
801    ) -> Result<Option<IcebergSnapshotInfo>> {
802        let mut map = self.ctx.iceberg_snapshot_info_map();
803        let catalog = source.source_catalog().ok_or_else(|| {
804            crate::error::ErrorCode::InternalError(
805                "Iceberg source must have a valid source catalog".to_owned(),
806            )
807        })?;
808        if let Some(&snapshot_info) = map.get(&catalog.id) {
809            return Ok(snapshot_info);
810        }
811
812        #[cfg(madsim)]
813        return Err(crate::error::ErrorCode::BindError(
814            "iceberg source time travel can't be used in the madsim mode".to_string(),
815        )
816        .into());
817
818        #[cfg(not(madsim))]
819        {
820            let ConnectorProperties::Iceberg(prop) =
821                ConnectorProperties::extract(catalog.with_properties.clone(), false)?
822            else {
823                return Err(crate::error::ErrorCode::InternalError(
824                    "Iceberg source must have Iceberg connector properties".to_owned(),
825                )
826                .into());
827            };
828
829            let snapshot_info = tokio::task::block_in_place(|| {
830                crate::utils::FRONTEND_RUNTIME.block_on(async {
831                    prop.load_table().await.map(|table| {
832                        let metadata = table.metadata();
833                        metadata
834                            .current_snapshot()
835                            .map(|snapshot| IcebergSnapshotInfo {
836                                snapshot_id: snapshot.snapshot_id(),
837                                commit_epoch: risingwave_iceberg_commit_epoch(
838                                    metadata,
839                                    snapshot.snapshot_id(),
840                                ),
841                            })
842                    })
843                })
844            })?;
845            map.insert(catalog.id, snapshot_info);
846            Ok(snapshot_info)
847        }
848    }
849
850    /// Build a scan of insert rows that have not yet reached the current Iceberg snapshot.
851    /// Returns `None` for legacy/in-memory sinks whose persisted KV log store cannot be found.
852    fn plan_iceberg_log_store_scan(
853        &self,
854        table_catalog: &TableCatalog,
855        target_schema: &Schema,
856        source_columns: &[ColumnCatalog],
857        committed_epoch: Option<u64>,
858    ) -> Result<Option<PlanRef>> {
859        let Some(sink_name) = table_catalog.iceberg_sink_name() else {
860            return Ok(None);
861        };
862        let log_store_table = {
863            let catalog_reader = self.ctx.session_ctx().env().catalog_reader().read_guard();
864            let Ok(schema) =
865                catalog_reader.get_schema_by_id(table_catalog.database_id, table_catalog.schema_id)
866            else {
867                return Ok(None);
868            };
869            let Some(sink) = schema.get_created_sink_by_name(&sink_name) else {
870                return Ok(None);
871            };
872            schema
873                .iter_internal_table()
874                .find(|table| {
875                    table.job_id == Some(sink.id.as_job_id())
876                        && table
877                            .columns()
878                            .iter()
879                            .any(|column| column.name() == EPOCH_COLUMN_NAME)
880                        && table
881                            .columns()
882                            .iter()
883                            .any(|column| column.name() == ROW_OP_COLUMN_NAME)
884                })
885                .cloned()
886        };
887        let Some(log_store_table) = log_store_table else {
888            return Ok(None);
889        };
890
891        let Some(epoch_idx) = log_store_table
892            .columns()
893            .iter()
894            .position(|column| column.name() == EPOCH_COLUMN_NAME)
895        else {
896            return Ok(None);
897        };
898        let Some(row_op_idx) = log_store_table
899            .columns()
900            .iter()
901            .position(|column| column.name() == ROW_OP_COLUMN_NAME)
902        else {
903            return Ok(None);
904        };
905        // The log-store payload can contain hidden upstream columns such as `_rw_timestamp`
906        // that are not part of the Iceberg table schema. Do not let those implementation-only
907        // columns prevent the pending rows from being planned.
908        let payload_col_idx = log_store_table
909            .columns()
910            .iter()
911            .enumerate()
912            .skip(row_op_idx + 1)
913            .filter_map(|(idx, column)| (!column.is_hidden).then_some(idx))
914            .collect_vec();
915        if source_columns.len() != target_schema.len()
916            || payload_col_idx.len()
917                != source_columns
918                    .iter()
919                    .filter(|column| !column.is_hidden)
920                    .count()
921        {
922            return Ok(None);
923        }
924
925        let row_op_predicate: ExprImpl = FunctionCall::new(
926            ExprType::Equal,
927            vec![
928                InputRef::new(row_op_idx, DataType::Int16).into(),
929                Literal::new(Some(ScalarImpl::Int16(INSERT_OP_CODE)), DataType::Int16).into(),
930            ],
931        )?
932        .into();
933        let mut predicate = Condition::with_expr(row_op_predicate);
934        if let Some(epoch) = committed_epoch {
935            let epoch_predicate: ExprImpl = FunctionCall::new(
936                ExprType::GreaterThan,
937                vec![
938                    InputRef::new(epoch_idx, DataType::Int64).into(),
939                    Literal::new(
940                        Some(ScalarImpl::Int64(encode_epoch(epoch))),
941                        DataType::Int64,
942                    )
943                    .into(),
944                ],
945            )?
946            .into();
947            predicate = predicate.and(Condition::with_expr(epoch_predicate));
948        }
949
950        let log_store_scan: PlanRef = LogicalScan::from(generic::TableScan::new(
951            payload_col_idx,
952            log_store_table,
953            vec![],
954            vec![],
955            self.ctx(),
956            predicate,
957            None,
958        ))
959        .into();
960
961        let source_types = log_store_scan.schema().data_types();
962        let target_types = target_schema.data_types();
963        let mut payload_idx = 0;
964        let project_exprs = source_columns
965            .iter()
966            .zip_eq_fast(target_types)
967            .map(|(source_column, target_type)| {
968                if source_column.is_hidden {
969                    return Ok(Literal::new(None, target_type).into());
970                }
971
972                let source_type = source_types[payload_idx].clone();
973                let mut expr: ExprImpl = InputRef::new(payload_idx, source_type).into();
974                payload_idx += 1;
975                FunctionCall::cast_mut(&mut expr, &target_type, CastContext::Explicit).map_err(
976                    |error| {
977                        ErrorCode::InternalError(format!(
978                            "failed to align Iceberg log-store payload type: {}",
979                            error.as_report()
980                        ))
981                    },
982                )?;
983                Ok(expr)
984            })
985            .collect::<Result<Vec<_>>>()?;
986        Ok(Some(LogicalProject::create(log_store_scan, project_exprs)))
987    }
988
989    fn get_iceberg_source_by_table_catalog(
990        &self,
991        table_catalog: &TableCatalog,
992    ) -> Option<SourceCatalog> {
993        let catalog_reader = self.ctx.session_ctx().env().catalog_reader().read_guard();
994
995        let iceberg_source_name = table_catalog.iceberg_source_name()?;
996        let schema = catalog_reader
997            .get_schema_by_id(table_catalog.database_id, table_catalog.schema_id)
998            .ok()?;
999        let source_catalog = schema.get_source_by_name(&iceberg_source_name)?;
1000        Some(source_catalog.deref().clone())
1001    }
1002}
1003
1004/// Find the latest known RisingWave commit boundary. Iceberg compaction produces `replace`
1005/// snapshots without changing table contents, so the marker is inherited across a chain of those
1006/// snapshots. We deliberately stop at any other unmarked operation: an external or legacy append
1007/// cannot be assigned a safe log-store boundary.
1008fn risingwave_iceberg_commit_epoch(metadata: &TableMetadata, snapshot_id: i64) -> Option<u64> {
1009    let mut snapshot = metadata.snapshot_by_id(snapshot_id)?;
1010    loop {
1011        if let Some(epoch) = snapshot
1012            .summary()
1013            .additional_properties
1014            .get(RISINGWAVE_ICEBERG_COMMIT_EPOCH)
1015        {
1016            return epoch.parse().ok();
1017        }
1018        if snapshot.summary().operation != Operation::Replace {
1019            return None;
1020        }
1021        snapshot = metadata.snapshot_by_id(snapshot.parent_snapshot_id()?)?;
1022    }
1023}
1024
1025fn resolve_current_cdc_table_desc(
1026    mut desc: CdcTableDesc,
1027    upstream_properties: &WithOptionsSecResolved,
1028) -> Result<CdcTableDesc> {
1029    let (current_properties, normalized_external_table_name) =
1030        derive_with_options_for_cdc_table(upstream_properties, desc.external_table_name)?;
1031    let (connect_properties, secret_refs) = current_properties.into_parts();
1032    desc.external_table_name = normalized_external_table_name;
1033    desc.connect_properties = connect_properties;
1034    desc.secret_refs = secret_refs;
1035    Ok(desc)
1036}
1037
1038#[cfg(test)]
1039mod tests {
1040    use std::collections::BTreeMap;
1041
1042    use iceberg::spec::{
1043        FormatVersion, MAIN_BRANCH, NestedField, PrimitiveType, Schema as IcebergSchema, Snapshot,
1044        SortOrder, Summary, TableMetadataBuilder, Type, UnboundPartitionSpec,
1045    };
1046
1047    use super::*;
1048
1049    #[test]
1050    fn test_risingwave_commit_epoch_inherits_only_across_replace_snapshots() {
1051        let append = snapshot(
1052            1,
1053            None,
1054            Operation::Append,
1055            HashMap::from([(RISINGWAVE_ICEBERG_COMMIT_EPOCH.to_owned(), "42".to_owned())]),
1056        );
1057        let replace = snapshot(2, Some(1), Operation::Replace, HashMap::new());
1058        let metadata = metadata_with_current(vec![append, replace]);
1059        assert_eq!(risingwave_iceberg_commit_epoch(&metadata, 2), Some(42));
1060
1061        let external_append = snapshot(3, Some(2), Operation::Append, HashMap::new());
1062        let metadata = metadata_with_current(vec![
1063            snapshot(
1064                1,
1065                None,
1066                Operation::Append,
1067                HashMap::from([(RISINGWAVE_ICEBERG_COMMIT_EPOCH.to_owned(), "42".to_owned())]),
1068            ),
1069            snapshot(2, Some(1), Operation::Replace, HashMap::new()),
1070            external_append,
1071        ]);
1072        assert_eq!(risingwave_iceberg_commit_epoch(&metadata, 3), None);
1073    }
1074
1075    fn metadata_with_current(mut snapshots: Vec<Snapshot>) -> TableMetadata {
1076        let current = snapshots.pop().unwrap();
1077        let mut builder = TableMetadataBuilder::new(
1078            IcebergSchema::builder()
1079                .with_fields(vec![
1080                    NestedField::new(1, "id", Type::Primitive(PrimitiveType::Long), false).into(),
1081                ])
1082                .build()
1083                .unwrap(),
1084            UnboundPartitionSpec::builder().build(),
1085            SortOrder::unsorted_order(),
1086            "s3://warehouse/db/table".to_owned(),
1087            FormatVersion::V2,
1088            HashMap::new(),
1089        )
1090        .unwrap();
1091        for snapshot in snapshots {
1092            builder = builder.add_snapshot(snapshot).unwrap();
1093        }
1094        builder
1095            .set_branch_snapshot(current, MAIN_BRANCH)
1096            .unwrap()
1097            .build()
1098            .unwrap()
1099            .metadata
1100    }
1101
1102    fn snapshot(
1103        snapshot_id: i64,
1104        parent_snapshot_id: Option<i64>,
1105        operation: Operation,
1106        additional_properties: HashMap<String, String>,
1107    ) -> Snapshot {
1108        Snapshot::builder()
1109            .with_snapshot_id(snapshot_id)
1110            .with_parent_snapshot_id(parent_snapshot_id)
1111            .with_sequence_number(snapshot_id)
1112            .with_timestamp_ms(snapshot_id)
1113            .with_manifest_list(format!("/snap-{snapshot_id}.avro"))
1114            .with_summary(Summary {
1115                operation,
1116                additional_properties,
1117            })
1118            .with_schema_id(0)
1119            .build()
1120    }
1121
1122    #[test]
1123    fn test_resolve_current_cdc_table_desc_replaces_stale_properties() {
1124        let desc = CdcTableDesc {
1125            external_table_name: "risedev.orders".to_owned(),
1126            connect_properties: BTreeMap::from([("hostname".to_owned(), "old-host".to_owned())]),
1127            ..Default::default()
1128        };
1129        let upstream_properties = WithOptionsSecResolved::new(
1130            BTreeMap::from([
1131                ("connector".to_owned(), "mysql-cdc".to_owned()),
1132                ("database.name".to_owned(), "risedev".to_owned()),
1133                ("hostname".to_owned(), "new-host".to_owned()),
1134            ]),
1135            BTreeMap::new(),
1136        );
1137
1138        let desc = resolve_current_cdc_table_desc(desc, &upstream_properties).unwrap();
1139
1140        assert_eq!(desc.connect_properties.get("hostname").unwrap(), "new-host");
1141        assert_eq!(desc.connect_properties.get("table.name").unwrap(), "orders");
1142        assert_eq!(
1143            desc.connect_properties.get("database.name").unwrap(),
1144            "risedev"
1145        );
1146    }
1147}