Skip to main content

risingwave_meta_service/
ddl_service.rs

1// Copyright 2023 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, HashSet};
16use std::pin::pin;
17use std::sync::Arc;
18
19use anyhow::{Context, anyhow};
20use futures::future::select;
21use rand::rng as thread_rng;
22use rand::seq::IndexedRandom;
23use replace_job_plan::{ReplaceSource, ReplaceTable};
24use risingwave_common::catalog::cdc_type_compatibility::cdc_auto_schema_change_existing_type_compatible;
25use risingwave_common::catalog::{AlterDatabaseParam, ColumnCatalog};
26use risingwave_common::id::{ObjectId, TableId};
27use risingwave_common::system_param::adaptive_parallelism_strategy::parse_strategy;
28use risingwave_common::types::DataType;
29use risingwave_common::util::stream_graph_visitor;
30use risingwave_connector::sink::catalog::SinkId;
31use risingwave_connector::sink::iceberg::ENABLE_PK_INDEX;
32use risingwave_meta::barrier::{BarrierScheduler, Command, ResumeBackfillTarget};
33use risingwave_meta::manager::{EventLogManagerRef, MetadataManager, iceberg_compaction};
34use risingwave_meta::model::TableParallelism as ModelTableParallelism;
35use risingwave_meta::rpc::metrics::MetaMetrics;
36use risingwave_meta::stream::{ParallelismPolicy, ReschedulePolicy, ResourceGroupPolicy};
37use risingwave_meta::{MetaResult, bail_invalid_parameter, bail_unavailable};
38use risingwave_meta_model::StreamingParallelism;
39use risingwave_pb::catalog::connection::Info as ConnectionInfo;
40use risingwave_pb::catalog::table::{CdcTableType as PbCdcTableType, OptionalAssociatedSourceId};
41use risingwave_pb::catalog::{Comment, Connection, PbCreateType, Secret, Table};
42use risingwave_pb::common::WorkerType;
43use risingwave_pb::common::worker_node::State;
44use risingwave_pb::ddl_service::create_iceberg_table_request::{PbSinkJobInfo, PbTableJobInfo};
45use risingwave_pb::ddl_service::ddl_service_server::DdlService;
46use risingwave_pb::ddl_service::drop_table_request::PbSourceId;
47use risingwave_pb::ddl_service::replace_job_plan::ReplaceMaterializedView;
48use risingwave_pb::ddl_service::{streaming_job_resource_type, *};
49use risingwave_pb::frontend_service::GetTableReplacePlanRequest;
50use risingwave_pb::meta::event_log;
51use risingwave_pb::meta::table_parallelism::{FixedParallelism, Parallelism};
52use risingwave_pb::stream_plan::stream_node::NodeBody;
53use risingwave_pb::stream_plan::throttle_mutation::ThrottleConfig;
54use thiserror_ext::AsReport;
55use tokio::sync::oneshot::Sender;
56use tokio::task::JoinHandle;
57use tonic::{Request, Response, Status};
58
59use crate::MetaError;
60use crate::barrier::BarrierManagerRef;
61use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
62use crate::manager::sink_coordination::SinkCoordinatorManager;
63use crate::manager::{MetaSrvEnv, StreamingJob};
64use crate::rpc::ddl_controller::{
65    DdlCommand, DdlController, DropMode, ReplaceStreamJobInfo, StreamingJobId,
66};
67use crate::stream::{GlobalStreamManagerRef, SourceManagerRef};
68
69#[derive(Clone)]
70pub struct DdlServiceImpl {
71    env: MetaSrvEnv,
72
73    metadata_manager: MetadataManager,
74    sink_manager: SinkCoordinatorManager,
75    ddl_controller: DdlController,
76    meta_metrics: Arc<MetaMetrics>,
77    iceberg_compaction_manager: iceberg_compaction::IcebergCompactionManagerRef,
78    barrier_scheduler: BarrierScheduler,
79    iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
80}
81
82impl DdlServiceImpl {
83    pub async fn new(
84        env: MetaSrvEnv,
85        metadata_manager: MetadataManager,
86        stream_manager: GlobalStreamManagerRef,
87        source_manager: SourceManagerRef,
88        barrier_manager: BarrierManagerRef,
89        sink_manager: SinkCoordinatorManager,
90        meta_metrics: Arc<MetaMetrics>,
91        iceberg_compaction_manager: iceberg_compaction::IcebergCompactionManagerRef,
92        barrier_scheduler: BarrierScheduler,
93        iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
94    ) -> Self {
95        let ddl_controller = DdlController::new(
96            env.clone(),
97            metadata_manager.clone(),
98            stream_manager,
99            source_manager,
100            barrier_manager,
101            sink_manager.clone(),
102            iceberg_compaction_manager.clone(),
103            iceberg_pk_index_sink_manager.clone(),
104        )
105        .await;
106        Self {
107            env,
108            metadata_manager,
109            sink_manager,
110            ddl_controller,
111            meta_metrics,
112            iceberg_compaction_manager,
113            barrier_scheduler,
114            iceberg_pk_index_sink_manager,
115        }
116    }
117
118    fn extract_replace_table_info(
119        ReplaceJobPlan {
120            fragment_graph,
121            replace_job,
122        }: ReplaceJobPlan,
123    ) -> ReplaceStreamJobInfo {
124        let replace_streaming_job: StreamingJob = match replace_job.unwrap() {
125            replace_job_plan::ReplaceJob::ReplaceTable(ReplaceTable {
126                table,
127                source,
128                job_type,
129            }) => StreamingJob::Table(
130                source,
131                table.unwrap(),
132                TableJobType::try_from(job_type).unwrap(),
133            ),
134            replace_job_plan::ReplaceJob::ReplaceSource(ReplaceSource { source }) => {
135                StreamingJob::Source(source.unwrap())
136            }
137            replace_job_plan::ReplaceJob::ReplaceMaterializedView(ReplaceMaterializedView {
138                table,
139            }) => StreamingJob::MaterializedView(table.unwrap()),
140            replace_job_plan::ReplaceJob::ReplaceSink(_) => unreachable!("use replace sink path"),
141        };
142
143        ReplaceStreamJobInfo {
144            streaming_job: replace_streaming_job,
145            fragment_graph: fragment_graph.unwrap(),
146        }
147    }
148
149    fn default_streaming_job_resource_type() -> streaming_job_resource_type::ResourceType {
150        streaming_job_resource_type::ResourceType::Regular(true)
151    }
152
153    pub fn start_migrate_table_fragments(&self) -> (JoinHandle<()>, Sender<()>) {
154        tracing::info!("start migrate legacy table fragments task");
155        let env = self.env.clone();
156        let metadata_manager = self.metadata_manager.clone();
157        let ddl_controller = self.ddl_controller.clone();
158
159        let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel();
160        let join_handle = tokio::spawn(async move {
161            async fn migrate_inner(
162                env: &MetaSrvEnv,
163                metadata_manager: &MetadataManager,
164                ddl_controller: &DdlController,
165            ) -> MetaResult<()> {
166                let tables = metadata_manager
167                    .catalog_controller
168                    .list_unmigrated_tables()
169                    .await?;
170
171                if tables.is_empty() {
172                    tracing::info!("no legacy table fragments need migration");
173                    return Ok(());
174                }
175
176                let client = {
177                    let workers = metadata_manager
178                        .list_worker_node(Some(WorkerType::Frontend), Some(State::Running))
179                        .await?;
180                    if workers.is_empty() {
181                        return Err(anyhow::anyhow!("no active frontend nodes found").into());
182                    }
183                    let worker = workers.choose(&mut thread_rng()).unwrap();
184                    env.frontend_client_pool().get(worker).await?
185                };
186
187                for table in tables {
188                    let start = tokio::time::Instant::now();
189                    let req = GetTableReplacePlanRequest {
190                        database_id: table.database_id,
191                        table_id: table.id,
192                        cdc_table_change: None,
193                    };
194                    let resp = client
195                        .get_table_replace_plan(req)
196                        .await
197                        .context("failed to get table replace plan from frontend")?;
198
199                    let plan = resp.into_inner().replace_plan.unwrap();
200                    let replace_info = DdlServiceImpl::extract_replace_table_info(plan);
201                    ddl_controller
202                        .run_command(DdlCommand::ReplaceStreamJob(replace_info))
203                        .await?;
204                    tracing::info!(elapsed=?start.elapsed(), table_id=%table.id, "migrated table fragments");
205                }
206                tracing::info!("successfully migrated all legacy table fragments");
207
208                Ok(())
209            }
210
211            let migrate_future = async move {
212                let mut attempt = 0;
213                loop {
214                    match migrate_inner(&env, &metadata_manager, &ddl_controller).await {
215                        Ok(_) => break,
216                        Err(e) => {
217                            attempt += 1;
218                            tracing::error!(
219                                "failed to migrate legacy table fragments: {}, attempt {}, retrying in 5 secs",
220                                e.as_report(),
221                                attempt
222                            );
223                            tokio::time::sleep(std::time::Duration::from_secs(5)).await;
224                        }
225                    }
226                }
227            };
228
229            select(pin!(migrate_future), shutdown_rx).await;
230        });
231
232        (join_handle, shutdown_tx)
233    }
234}
235
236#[async_trait::async_trait]
237impl DdlService for DdlServiceImpl {
238    async fn create_database(
239        &self,
240        request: Request<CreateDatabaseRequest>,
241    ) -> Result<Response<CreateDatabaseResponse>, Status> {
242        let req = request.into_inner();
243        let database = req.get_db()?.clone();
244        let version = self
245            .ddl_controller
246            .run_command(DdlCommand::CreateDatabase(database))
247            .await?;
248
249        Ok(Response::new(CreateDatabaseResponse {
250            status: None,
251            version,
252        }))
253    }
254
255    async fn drop_database(
256        &self,
257        request: Request<DropDatabaseRequest>,
258    ) -> Result<Response<DropDatabaseResponse>, Status> {
259        let req = request.into_inner();
260        let database_id = req.get_database_id();
261
262        let version = self
263            .ddl_controller
264            .run_command(DdlCommand::DropDatabase(database_id))
265            .await?;
266
267        Ok(Response::new(DropDatabaseResponse {
268            status: None,
269            version,
270        }))
271    }
272
273    async fn create_secret(
274        &self,
275        request: Request<CreateSecretRequest>,
276    ) -> Result<Response<CreateSecretResponse>, Status> {
277        let req = request.into_inner();
278        let pb_secret = Secret {
279            id: 0.into(),
280            name: req.get_name().clone(),
281            database_id: req.get_database_id(),
282            value: req.get_value().clone(),
283            owner: req.get_owner_id(),
284            schema_id: req.get_schema_id(),
285        };
286        let version = self
287            .ddl_controller
288            .run_command(DdlCommand::CreateSecret(pb_secret))
289            .await?;
290
291        Ok(Response::new(CreateSecretResponse { version }))
292    }
293
294    async fn drop_secret(
295        &self,
296        request: Request<DropSecretRequest>,
297    ) -> Result<Response<DropSecretResponse>, Status> {
298        let req = request.into_inner();
299        let secret_id = req.get_secret_id();
300        let drop_mode = DropMode::from_request_setting(req.cascade);
301        let version = self
302            .ddl_controller
303            .run_command(DdlCommand::DropSecret(secret_id, drop_mode))
304            .await?;
305
306        Ok(Response::new(DropSecretResponse {
307            status: None,
308            version,
309        }))
310    }
311
312    async fn alter_secret(
313        &self,
314        request: Request<AlterSecretRequest>,
315    ) -> Result<Response<AlterSecretResponse>, Status> {
316        let req = request.into_inner();
317        let pb_secret = Secret {
318            id: req.get_secret_id(),
319            name: req.get_name().clone(),
320            database_id: req.get_database_id(),
321            value: req.get_value().clone(),
322            owner: req.get_owner_id(),
323            schema_id: req.get_schema_id(),
324        };
325        let version = self
326            .ddl_controller
327            .run_command(DdlCommand::AlterSecret(pb_secret))
328            .await?;
329
330        Ok(Response::new(AlterSecretResponse { version }))
331    }
332
333    async fn create_schema(
334        &self,
335        request: Request<CreateSchemaRequest>,
336    ) -> Result<Response<CreateSchemaResponse>, Status> {
337        let req = request.into_inner();
338        let schema = req.get_schema()?.clone();
339        let version = self
340            .ddl_controller
341            .run_command(DdlCommand::CreateSchema(schema))
342            .await?;
343
344        Ok(Response::new(CreateSchemaResponse {
345            status: None,
346            version,
347        }))
348    }
349
350    async fn drop_schema(
351        &self,
352        request: Request<DropSchemaRequest>,
353    ) -> Result<Response<DropSchemaResponse>, Status> {
354        let req = request.into_inner();
355        let schema_id = req.get_schema_id();
356        let drop_mode = DropMode::from_request_setting(req.cascade);
357        let version = self
358            .ddl_controller
359            .run_command(DdlCommand::DropSchema(schema_id, drop_mode))
360            .await?;
361        Ok(Response::new(DropSchemaResponse {
362            status: None,
363            version,
364        }))
365    }
366
367    async fn create_source(
368        &self,
369        request: Request<CreateSourceRequest>,
370    ) -> Result<Response<CreateSourceResponse>, Status> {
371        let req = request.into_inner();
372        let source = req.get_source()?.clone();
373
374        match req.fragment_graph {
375            None => {
376                let version = self
377                    .ddl_controller
378                    .run_command(DdlCommand::CreateNonSharedSource(source, None))
379                    .await?;
380                Ok(Response::new(CreateSourceResponse {
381                    status: None,
382                    version,
383                }))
384            }
385            Some(fragment_graph) => {
386                // The id of stream job has been set above
387                let stream_job = StreamingJob::Source(source);
388                let version = self
389                    .ddl_controller
390                    .run_command(DdlCommand::CreateStreamingJob {
391                        stream_job,
392                        fragment_graph,
393                        dependencies: HashSet::new(),
394                        resource_type: Self::default_streaming_job_resource_type(),
395                        if_not_exists: req.if_not_exists,
396                        refresh_interval_sec: None,
397                        replace_sink: None,
398                        since_timestamp_epoch: None,
399                    })
400                    .await?;
401                Ok(Response::new(CreateSourceResponse {
402                    status: None,
403                    version,
404                }))
405            }
406        }
407    }
408
409    async fn drop_source(
410        &self,
411        request: Request<DropSourceRequest>,
412    ) -> Result<Response<DropSourceResponse>, Status> {
413        let request = request.into_inner();
414        let source_id = request.source_id;
415        let drop_mode = DropMode::from_request_setting(request.cascade);
416        let version = self
417            .ddl_controller
418            .run_command(DdlCommand::DropSource(source_id, drop_mode))
419            .await?;
420
421        Ok(Response::new(DropSourceResponse {
422            status: None,
423            version,
424        }))
425    }
426
427    async fn reset_source(
428        &self,
429        request: Request<ResetSourceRequest>,
430    ) -> Result<Response<ResetSourceResponse>, Status> {
431        let request = request.into_inner();
432        let source_id = request.source_id;
433
434        tracing::info!(
435            source_id = %source_id,
436            "Received RESET SOURCE request, routing to DDL controller"
437        );
438
439        // Route to DDL controller
440        let version = self
441            .ddl_controller
442            .run_command(DdlCommand::ResetSource(source_id))
443            .await?;
444
445        Ok(Response::new(ResetSourceResponse {
446            status: None,
447            version,
448        }))
449    }
450
451    async fn create_sink(
452        &self,
453        request: Request<CreateSinkRequest>,
454    ) -> Result<Response<CreateSinkResponse>, Status> {
455        self.env.idle_manager().record_activity();
456
457        let req = request.into_inner();
458
459        let sink = req.get_sink()?.clone();
460        let fragment_graph = req.get_fragment_graph()?.clone();
461        let dependencies = req.get_dependencies().iter().copied().collect();
462        let resource_type = req
463            .resource_type
464            .and_then(|resource_type| resource_type.resource_type)
465            .unwrap_or_else(Self::default_streaming_job_resource_type);
466        let since_timestamp_epoch = req.since_timestamp_epoch;
467
468        let stream_job = StreamingJob::Sink(sink, None);
469
470        let command = DdlCommand::CreateStreamingJob {
471            stream_job,
472            fragment_graph,
473            dependencies,
474            resource_type,
475            if_not_exists: req.if_not_exists,
476            refresh_interval_sec: None,
477            replace_sink: None,
478            since_timestamp_epoch,
479        };
480
481        let version = self.ddl_controller.run_command(command).await?;
482
483        Ok(Response::new(CreateSinkResponse {
484            status: None,
485            version,
486        }))
487    }
488
489    async fn drop_sink(
490        &self,
491        request: Request<DropSinkRequest>,
492    ) -> Result<Response<DropSinkResponse>, Status> {
493        let request = request.into_inner();
494        let sink_id = request.sink_id;
495        let drop_mode = DropMode::from_request_setting(request.cascade);
496
497        let command = DdlCommand::DropStreamingJob {
498            job_id: StreamingJobId::Sink(sink_id),
499            drop_mode,
500        };
501
502        let version = self.ddl_controller.run_command(command).await?;
503
504        self.sink_manager
505            .stop_sink_coordinator(vec![SinkId::from(sink_id)])
506            .await;
507        self.iceberg_compaction_manager
508            .clear_iceberg_maintenance_by_sink_id(SinkId::from(sink_id));
509
510        Ok(Response::new(DropSinkResponse {
511            status: None,
512            version,
513        }))
514    }
515
516    async fn create_subscription(
517        &self,
518        request: Request<CreateSubscriptionRequest>,
519    ) -> Result<Response<CreateSubscriptionResponse>, Status> {
520        self.env.idle_manager().record_activity();
521
522        let req = request.into_inner();
523
524        let subscription = req.get_subscription()?.clone();
525        let command = DdlCommand::CreateSubscription(subscription);
526
527        let version = self.ddl_controller.run_command(command).await?;
528
529        Ok(Response::new(CreateSubscriptionResponse {
530            status: None,
531            version,
532        }))
533    }
534
535    async fn drop_subscription(
536        &self,
537        request: Request<DropSubscriptionRequest>,
538    ) -> Result<Response<DropSubscriptionResponse>, Status> {
539        let request = request.into_inner();
540        let subscription_id = request.subscription_id;
541        let drop_mode = DropMode::from_request_setting(request.cascade);
542
543        let command = DdlCommand::DropSubscription(subscription_id, drop_mode);
544
545        let version = self.ddl_controller.run_command(command).await?;
546
547        Ok(Response::new(DropSubscriptionResponse {
548            status: None,
549            version,
550        }))
551    }
552
553    async fn create_materialized_view(
554        &self,
555        request: Request<CreateMaterializedViewRequest>,
556    ) -> Result<Response<CreateMaterializedViewResponse>, Status> {
557        self.env.idle_manager().record_activity();
558
559        let req = request.into_inner();
560        let mview = req.get_materialized_view()?.clone();
561        let fragment_graph = req.get_fragment_graph()?.clone();
562        let dependencies = req.get_dependencies().iter().copied().collect();
563        let resource_type = req.resource_type.unwrap().resource_type.unwrap();
564
565        let stream_job = StreamingJob::MaterializedView(mview);
566        let version = self
567            .ddl_controller
568            .run_command(DdlCommand::CreateStreamingJob {
569                stream_job,
570                fragment_graph,
571                dependencies,
572                resource_type,
573                if_not_exists: req.if_not_exists,
574                refresh_interval_sec: req.refresh_interval_sec,
575                replace_sink: None,
576                since_timestamp_epoch: None,
577            })
578            .await?;
579
580        Ok(Response::new(CreateMaterializedViewResponse {
581            status: None,
582            version,
583        }))
584    }
585
586    async fn drop_materialized_view(
587        &self,
588        request: Request<DropMaterializedViewRequest>,
589    ) -> Result<Response<DropMaterializedViewResponse>, Status> {
590        self.env.idle_manager().record_activity();
591
592        let request = request.into_inner();
593        let table_id = request.table_id;
594        let drop_mode = DropMode::from_request_setting(request.cascade);
595
596        let version = self
597            .ddl_controller
598            .run_command(DdlCommand::DropStreamingJob {
599                job_id: StreamingJobId::MaterializedView(table_id),
600                drop_mode,
601            })
602            .await?;
603
604        Ok(Response::new(DropMaterializedViewResponse {
605            status: None,
606            version,
607        }))
608    }
609
610    async fn create_index(
611        &self,
612        request: Request<CreateIndexRequest>,
613    ) -> Result<Response<CreateIndexResponse>, Status> {
614        self.env.idle_manager().record_activity();
615
616        let req = request.into_inner();
617        let index = req.get_index()?.clone();
618        let index_table = req.get_index_table()?.clone();
619        let fragment_graph = req.get_fragment_graph()?.clone();
620        let resource_type = req
621            .resource_type
622            .and_then(|resource_type| resource_type.resource_type)
623            .unwrap_or_else(Self::default_streaming_job_resource_type);
624
625        let stream_job = StreamingJob::Index(index, index_table);
626        let version = self
627            .ddl_controller
628            .run_command(DdlCommand::CreateStreamingJob {
629                stream_job,
630                fragment_graph,
631                dependencies: HashSet::new(),
632                resource_type,
633                if_not_exists: req.if_not_exists,
634                replace_sink: None,
635                refresh_interval_sec: None,
636                since_timestamp_epoch: None,
637            })
638            .await?;
639
640        Ok(Response::new(CreateIndexResponse {
641            status: None,
642            version,
643        }))
644    }
645
646    async fn drop_index(
647        &self,
648        request: Request<DropIndexRequest>,
649    ) -> Result<Response<DropIndexResponse>, Status> {
650        self.env.idle_manager().record_activity();
651
652        let request = request.into_inner();
653        let index_id = request.index_id;
654        let drop_mode = DropMode::from_request_setting(request.cascade);
655        let version = self
656            .ddl_controller
657            .run_command(DdlCommand::DropStreamingJob {
658                job_id: StreamingJobId::Index(index_id),
659                drop_mode,
660            })
661            .await?;
662
663        Ok(Response::new(DropIndexResponse {
664            status: None,
665            version,
666        }))
667    }
668
669    async fn create_function(
670        &self,
671        request: Request<CreateFunctionRequest>,
672    ) -> Result<Response<CreateFunctionResponse>, Status> {
673        let req = request.into_inner();
674        let function = req.get_function()?.clone();
675
676        let version = self
677            .ddl_controller
678            .run_command(DdlCommand::CreateFunction(function))
679            .await?;
680
681        Ok(Response::new(CreateFunctionResponse {
682            status: None,
683            version,
684        }))
685    }
686
687    async fn drop_function(
688        &self,
689        request: Request<DropFunctionRequest>,
690    ) -> Result<Response<DropFunctionResponse>, Status> {
691        let request = request.into_inner();
692
693        let version = self
694            .ddl_controller
695            .run_command(DdlCommand::DropFunction(
696                request.function_id,
697                DropMode::from_request_setting(request.cascade),
698            ))
699            .await?;
700
701        Ok(Response::new(DropFunctionResponse {
702            status: None,
703            version,
704        }))
705    }
706
707    async fn create_table(
708        &self,
709        request: Request<CreateTableRequest>,
710    ) -> Result<Response<CreateTableResponse>, Status> {
711        let request = request.into_inner();
712        let job_type = request.get_job_type().unwrap_or_default();
713        let dependencies = request.get_dependencies().iter().copied().collect();
714        let source = request.source;
715        let mview = request.materialized_view.unwrap();
716        let fragment_graph = request.fragment_graph.unwrap();
717
718        let stream_job = StreamingJob::Table(source, mview, job_type);
719        let version = self
720            .ddl_controller
721            .run_command(DdlCommand::CreateStreamingJob {
722                stream_job,
723                fragment_graph,
724                dependencies,
725                resource_type: Self::default_streaming_job_resource_type(),
726                if_not_exists: request.if_not_exists,
727                refresh_interval_sec: None,
728                replace_sink: None,
729                since_timestamp_epoch: None,
730            })
731            .await?;
732
733        Ok(Response::new(CreateTableResponse {
734            status: None,
735            version,
736        }))
737    }
738
739    async fn drop_table(
740        &self,
741        request: Request<DropTableRequest>,
742    ) -> Result<Response<DropTableResponse>, Status> {
743        let request = request.into_inner();
744        let source_id = request.source_id;
745        let table_id = request.table_id;
746
747        let drop_mode = DropMode::from_request_setting(request.cascade);
748        let version = self
749            .ddl_controller
750            .run_command(DdlCommand::DropStreamingJob {
751                job_id: StreamingJobId::Table(source_id.map(|PbSourceId::Id(id)| id), table_id),
752                drop_mode,
753            })
754            .await?;
755
756        Ok(Response::new(DropTableResponse {
757            status: None,
758            version,
759        }))
760    }
761
762    async fn create_view(
763        &self,
764        request: Request<CreateViewRequest>,
765    ) -> Result<Response<CreateViewResponse>, Status> {
766        let req = request.into_inner();
767        let view = req.get_view()?.clone();
768        let dependencies = req
769            .get_dependencies()
770            .iter()
771            .copied()
772            .collect::<HashSet<_>>();
773
774        let version = self
775            .ddl_controller
776            .run_command(DdlCommand::CreateView(view, dependencies))
777            .await?;
778
779        Ok(Response::new(CreateViewResponse {
780            status: None,
781            version,
782        }))
783    }
784
785    async fn drop_view(
786        &self,
787        request: Request<DropViewRequest>,
788    ) -> Result<Response<DropViewResponse>, Status> {
789        let request = request.into_inner();
790        let view_id = request.get_view_id();
791        let drop_mode = DropMode::from_request_setting(request.cascade);
792        let version = self
793            .ddl_controller
794            .run_command(DdlCommand::DropView(view_id, drop_mode))
795            .await?;
796        Ok(Response::new(DropViewResponse {
797            status: None,
798            version,
799        }))
800    }
801
802    async fn risectl_list_state_tables(
803        &self,
804        _request: Request<RisectlListStateTablesRequest>,
805    ) -> Result<Response<RisectlListStateTablesResponse>, Status> {
806        let tables = self
807            .metadata_manager
808            .catalog_controller
809            .list_all_state_tables()
810            .await?;
811        Ok(Response::new(RisectlListStateTablesResponse { tables }))
812    }
813
814    async fn risectl_resume_backfill(
815        &self,
816        request: Request<RisectlResumeBackfillRequest>,
817    ) -> Result<Response<RisectlResumeBackfillResponse>, Status> {
818        let request = request.into_inner();
819        let target = request
820            .target
821            .ok_or_else(|| Status::invalid_argument("missing resume backfill target"))?;
822
823        match target {
824            risectl_resume_backfill_request::Target::JobId(job_id) => {
825                let database_id = self
826                    .metadata_manager
827                    .catalog_controller
828                    .get_object_database_id(ObjectId::new(job_id.as_raw_id()))
829                    .await?;
830                self.barrier_scheduler
831                    .run_command(
832                        database_id,
833                        Command::ResumeBackfill {
834                            target: ResumeBackfillTarget::Job(job_id),
835                        },
836                    )
837                    .await?;
838            }
839            risectl_resume_backfill_request::Target::FragmentId(fragment_id) => {
840                let mut job_ids = self
841                    .metadata_manager
842                    .catalog_controller
843                    .get_fragment_job_id(vec![fragment_id])
844                    .await?;
845                let job_id = job_ids
846                    .pop()
847                    .ok_or_else(|| Status::invalid_argument("fragment not found"))?;
848                let database_id = self
849                    .metadata_manager
850                    .catalog_controller
851                    .get_object_database_id(ObjectId::new(job_id.as_raw_id()))
852                    .await?;
853                self.barrier_scheduler
854                    .run_command(
855                        database_id,
856                        Command::ResumeBackfill {
857                            target: ResumeBackfillTarget::Fragment(fragment_id),
858                        },
859                    )
860                    .await?;
861            }
862        }
863
864        Ok(Response::new(RisectlResumeBackfillResponse {}))
865    }
866
867    async fn replace_job_plan(
868        &self,
869        request: Request<ReplaceJobPlanRequest>,
870    ) -> Result<Response<ReplaceJobPlanResponse>, Status> {
871        let ReplaceJobPlan {
872            fragment_graph,
873            replace_job,
874        } = request.into_inner().get_plan().cloned()?;
875
876        let command = match replace_job {
877            Some(replace_job_plan::ReplaceJob::ReplaceSink(replace_sink)) => {
878                let replace_job_plan::ReplaceSink {
879                    sink,
880                    old_sink_id,
881                    dependencies,
882                    resource_type,
883                } = replace_sink;
884                DdlCommand::CreateStreamingJob {
885                    stream_job: StreamingJob::Sink(sink.unwrap(), None),
886                    fragment_graph: fragment_graph.unwrap(),
887                    dependencies: dependencies.into_iter().collect::<HashSet<_>>(),
888                    resource_type: resource_type
889                        .and_then(|resource_type| resource_type.resource_type)
890                        .unwrap_or(streaming_job_resource_type::ResourceType::Regular(true)),
891                    if_not_exists: false,
892                    refresh_interval_sec: None,
893                    replace_sink: Some(old_sink_id),
894                    since_timestamp_epoch: None,
895                }
896            }
897            replace_job => {
898                DdlCommand::ReplaceStreamJob(Self::extract_replace_table_info(ReplaceJobPlan {
899                    fragment_graph,
900                    replace_job,
901                }))
902            }
903        };
904
905        let version = self.ddl_controller.run_command(command).await?;
906
907        Ok(Response::new(ReplaceJobPlanResponse {
908            status: None,
909            version,
910        }))
911    }
912
913    async fn get_table(
914        &self,
915        request: Request<GetTableRequest>,
916    ) -> Result<Response<GetTableResponse>, Status> {
917        let req = request.into_inner();
918        let table = self
919            .metadata_manager
920            .catalog_controller
921            .get_table_by_name(&req.database_name, &req.table_name)
922            .await?;
923
924        Ok(Response::new(GetTableResponse { table }))
925    }
926
927    async fn alter_name(
928        &self,
929        request: Request<AlterNameRequest>,
930    ) -> Result<Response<AlterNameResponse>, Status> {
931        let AlterNameRequest { object, new_name } = request.into_inner();
932        let version = self
933            .ddl_controller
934            .run_command(DdlCommand::AlterName(object.unwrap(), new_name))
935            .await?;
936        Ok(Response::new(AlterNameResponse {
937            status: None,
938            version,
939        }))
940    }
941
942    /// Only support add column for now.
943    async fn alter_source(
944        &self,
945        request: Request<AlterSourceRequest>,
946    ) -> Result<Response<AlterSourceResponse>, Status> {
947        let AlterSourceRequest { source } = request.into_inner();
948        let version = self
949            .ddl_controller
950            .run_command(DdlCommand::AlterNonSharedSource(source.unwrap()))
951            .await?;
952        Ok(Response::new(AlterSourceResponse {
953            status: None,
954            version,
955        }))
956    }
957
958    async fn alter_owner(
959        &self,
960        request: Request<AlterOwnerRequest>,
961    ) -> Result<Response<AlterOwnerResponse>, Status> {
962        let AlterOwnerRequest { object, owner_id } = request.into_inner();
963        let version = self
964            .ddl_controller
965            .run_command(DdlCommand::AlterObjectOwner(object.unwrap(), owner_id as _))
966            .await?;
967        Ok(Response::new(AlterOwnerResponse {
968            status: None,
969            version,
970        }))
971    }
972
973    async fn alter_subscription_retention(
974        &self,
975        request: Request<AlterSubscriptionRetentionRequest>,
976    ) -> Result<Response<AlterSubscriptionRetentionResponse>, Status> {
977        let AlterSubscriptionRetentionRequest {
978            subscription_id,
979            retention_seconds,
980            definition,
981        } = request.into_inner();
982        let version = self
983            .ddl_controller
984            .run_command(DdlCommand::AlterSubscriptionRetention {
985                subscription_id,
986                retention_seconds,
987                definition,
988            })
989            .await?;
990        Ok(Response::new(AlterSubscriptionRetentionResponse {
991            status: None,
992            version,
993        }))
994    }
995
996    async fn alter_set_schema(
997        &self,
998        request: Request<AlterSetSchemaRequest>,
999    ) -> Result<Response<AlterSetSchemaResponse>, Status> {
1000        let AlterSetSchemaRequest {
1001            object,
1002            new_schema_id,
1003        } = request.into_inner();
1004        let version = self
1005            .ddl_controller
1006            .run_command(DdlCommand::AlterSetSchema(object.unwrap(), new_schema_id))
1007            .await?;
1008        Ok(Response::new(AlterSetSchemaResponse {
1009            status: None,
1010            version,
1011        }))
1012    }
1013
1014    async fn get_ddl_progress(
1015        &self,
1016        _request: Request<GetDdlProgressRequest>,
1017    ) -> Result<Response<GetDdlProgressResponse>, Status> {
1018        Ok(Response::new(GetDdlProgressResponse {
1019            ddl_progress: self.ddl_controller.get_ddl_progress().await?,
1020        }))
1021    }
1022
1023    async fn create_connection(
1024        &self,
1025        request: Request<CreateConnectionRequest>,
1026    ) -> Result<Response<CreateConnectionResponse>, Status> {
1027        let req = request.into_inner();
1028        if req.payload.is_none() {
1029            return Err(Status::invalid_argument("request is empty"));
1030        }
1031
1032        match req.payload.unwrap() {
1033            #[expect(deprecated)]
1034            create_connection_request::Payload::PrivateLink(_) => {
1035                panic!("Private Link Connection has been deprecated")
1036            }
1037            create_connection_request::Payload::ConnectionParams(params) => {
1038                let pb_connection = Connection {
1039                    id: 0.into(),
1040                    schema_id: req.schema_id,
1041                    database_id: req.database_id,
1042                    name: req.name,
1043                    info: Some(ConnectionInfo::ConnectionParams(params)),
1044                    owner: req.owner_id,
1045                };
1046                let version = self
1047                    .ddl_controller
1048                    .run_command(DdlCommand::CreateConnection(pb_connection))
1049                    .await?;
1050                Ok(Response::new(CreateConnectionResponse { version }))
1051            }
1052        }
1053    }
1054
1055    async fn list_connections(
1056        &self,
1057        _request: Request<ListConnectionsRequest>,
1058    ) -> Result<Response<ListConnectionsResponse>, Status> {
1059        let conns = self
1060            .metadata_manager
1061            .catalog_controller
1062            .list_connections()
1063            .await?;
1064
1065        Ok(Response::new(ListConnectionsResponse {
1066            connections: conns,
1067        }))
1068    }
1069
1070    async fn drop_connection(
1071        &self,
1072        request: Request<DropConnectionRequest>,
1073    ) -> Result<Response<DropConnectionResponse>, Status> {
1074        let req = request.into_inner();
1075        let drop_mode = DropMode::from_request_setting(req.cascade);
1076
1077        let version = self
1078            .ddl_controller
1079            .run_command(DdlCommand::DropConnection(req.connection_id, drop_mode))
1080            .await?;
1081
1082        Ok(Response::new(DropConnectionResponse {
1083            status: None,
1084            version,
1085        }))
1086    }
1087
1088    async fn comment_on(
1089        &self,
1090        request: Request<CommentOnRequest>,
1091    ) -> Result<Response<CommentOnResponse>, Status> {
1092        let req = request.into_inner();
1093        let comment = req.get_comment()?.clone();
1094
1095        let version = self
1096            .ddl_controller
1097            .run_command(DdlCommand::CommentOn(Comment {
1098                table_id: comment.table_id,
1099                schema_id: comment.schema_id,
1100                database_id: comment.database_id,
1101                column_index: comment.column_index,
1102                description: comment.description,
1103            }))
1104            .await?;
1105
1106        Ok(Response::new(CommentOnResponse {
1107            status: None,
1108            version,
1109        }))
1110    }
1111
1112    async fn get_tables(
1113        &self,
1114        request: Request<GetTablesRequest>,
1115    ) -> Result<Response<GetTablesResponse>, Status> {
1116        let GetTablesRequest {
1117            table_ids,
1118            include_dropped_tables,
1119        } = request.into_inner();
1120        let ret = self
1121            .metadata_manager
1122            .catalog_controller
1123            .get_table_by_ids(table_ids, include_dropped_tables)
1124            .await?;
1125
1126        let mut tables = HashMap::default();
1127        for table in ret {
1128            tables.insert(table.id, table);
1129        }
1130        Ok(Response::new(GetTablesResponse { tables }))
1131    }
1132
1133    async fn wait(&self, request: Request<WaitRequest>) -> Result<Response<WaitResponse>, Status> {
1134        let req = request.into_inner();
1135        let version = self.ddl_controller.wait(req.job_id).await?;
1136        Ok(Response::new(WaitResponse {
1137            version: Some(version),
1138        }))
1139    }
1140
1141    async fn alter_cdc_table_backfill_parallelism(
1142        &self,
1143        request: Request<AlterCdcTableBackfillParallelismRequest>,
1144    ) -> Result<Response<AlterCdcTableBackfillParallelismResponse>, Status> {
1145        let req = request.into_inner();
1146        let job_id = req.get_table_id();
1147        let parallelism = *req.get_parallelism()?;
1148
1149        let table_parallelism = ModelTableParallelism::from(parallelism);
1150        let streaming_parallelism = match table_parallelism {
1151            ModelTableParallelism::Fixed(n) => StreamingParallelism::Fixed(n),
1152            _ => bail_invalid_parameter!(
1153                "CDC table backfill parallelism must be set to a fixed value"
1154            ),
1155        };
1156
1157        self.ddl_controller
1158            .reschedule_cdc_table_backfill(
1159                job_id,
1160                ReschedulePolicy::Parallelism(ParallelismPolicy {
1161                    parallelism: streaming_parallelism,
1162                    adaptive_parallelism_strategy: None,
1163                }),
1164            )
1165            .await?;
1166        Ok(Response::new(AlterCdcTableBackfillParallelismResponse {}))
1167    }
1168
1169    async fn alter_parallelism(
1170        &self,
1171        request: Request<AlterParallelismRequest>,
1172    ) -> Result<Response<AlterParallelismResponse>, Status> {
1173        let req = request.into_inner();
1174
1175        let job_id = req.get_table_id();
1176        let parallelism = *req.get_parallelism()?;
1177        let deferred = req.get_deferred();
1178
1179        let adaptive_parallelism_strategy = req.adaptive_parallelism_strategy;
1180        let (parallelism, adaptive_parallelism_strategy) = match parallelism.get_parallelism()? {
1181            Parallelism::Fixed(FixedParallelism { parallelism }) => {
1182                (StreamingParallelism::Fixed(*parallelism as _), None)
1183            }
1184            Parallelism::Auto(_) | Parallelism::Adaptive(_) => (
1185                StreamingParallelism::Adaptive,
1186                Some(adaptive_parallelism_strategy.unwrap_or_else(|| "AUTO".to_owned())),
1187            ),
1188            Parallelism::Custom(_) => (
1189                StreamingParallelism::Adaptive,
1190                Some(adaptive_parallelism_strategy.ok_or_else(|| {
1191                    Status::invalid_argument(
1192                        "adaptive_parallelism_strategy is required for custom parallelism",
1193                    )
1194                })?),
1195            ),
1196        };
1197
1198        if let Some(strategy) = adaptive_parallelism_strategy.as_deref() {
1199            parse_strategy(strategy).map_err(|e| {
1200                Status::invalid_argument(format!(
1201                    "invalid adaptive parallelism strategy: {}",
1202                    e.as_report()
1203                ))
1204            })?;
1205        };
1206
1207        self.ddl_controller
1208            .reschedule_streaming_job(
1209                job_id,
1210                ReschedulePolicy::Parallelism(ParallelismPolicy {
1211                    parallelism,
1212                    adaptive_parallelism_strategy,
1213                }),
1214                deferred,
1215            )
1216            .await?;
1217
1218        Ok(Response::new(AlterParallelismResponse {}))
1219    }
1220
1221    async fn alter_backfill_parallelism(
1222        &self,
1223        request: Request<AlterBackfillParallelismRequest>,
1224    ) -> Result<Response<AlterBackfillParallelismResponse>, Status> {
1225        let req = request.into_inner();
1226
1227        let job_id = req.get_table_id();
1228        let deferred = req.get_deferred();
1229        let adaptive_parallelism_strategy = req.adaptive_parallelism_strategy;
1230
1231        let parallelism = match req.parallelism {
1232            None => None,
1233            Some(parallelism) => {
1234                let (parallelism, adaptive_parallelism_strategy) = match parallelism
1235                    .get_parallelism()?
1236                {
1237                    Parallelism::Fixed(FixedParallelism { parallelism }) => {
1238                        (StreamingParallelism::Fixed(*parallelism as _), None)
1239                    }
1240                    Parallelism::Auto(_) | Parallelism::Adaptive(_) => (
1241                        StreamingParallelism::Adaptive,
1242                        Some(adaptive_parallelism_strategy.unwrap_or_else(|| "AUTO".to_owned())),
1243                    ),
1244                    Parallelism::Custom(_) => (
1245                        StreamingParallelism::Adaptive,
1246                        Some(adaptive_parallelism_strategy.ok_or_else(|| {
1247                            Status::invalid_argument(
1248                                "adaptive_parallelism_strategy is required for custom parallelism",
1249                            )
1250                        })?),
1251                    ),
1252                };
1253
1254                if let Some(strategy) = adaptive_parallelism_strategy.as_deref() {
1255                    parse_strategy(strategy).map_err(|e| {
1256                        Status::invalid_argument(format!(
1257                            "invalid adaptive parallelism strategy: {}",
1258                            e.as_report()
1259                        ))
1260                    })?;
1261                };
1262
1263                Some(ParallelismPolicy {
1264                    parallelism,
1265                    adaptive_parallelism_strategy,
1266                })
1267            }
1268        };
1269
1270        self.ddl_controller
1271            .reschedule_streaming_job_backfill_parallelism(job_id, parallelism, deferred)
1272            .await?;
1273
1274        Ok(Response::new(AlterBackfillParallelismResponse {}))
1275    }
1276
1277    async fn alter_fragment_parallelism(
1278        &self,
1279        request: Request<AlterFragmentParallelismRequest>,
1280    ) -> Result<Response<AlterFragmentParallelismResponse>, Status> {
1281        let req = request.into_inner();
1282
1283        let fragment_ids = req.fragment_ids;
1284        if fragment_ids.is_empty() {
1285            return Err(Status::invalid_argument(
1286                "at least one fragment id must be provided",
1287            ));
1288        }
1289
1290        let parallelism = match req.parallelism {
1291            Some(parallelism) => {
1292                let streaming_parallelism = match parallelism.get_parallelism()? {
1293                    Parallelism::Fixed(FixedParallelism { parallelism }) => {
1294                        StreamingParallelism::Fixed(*parallelism as _)
1295                    }
1296                    Parallelism::Auto(_) | Parallelism::Adaptive(_) => {
1297                        StreamingParallelism::Adaptive
1298                    }
1299                    _ => bail_unavailable!(),
1300                };
1301                Some(streaming_parallelism)
1302            }
1303            None => None,
1304        };
1305
1306        let fragment_targets = fragment_ids
1307            .into_iter()
1308            .map(|fragment_id| (fragment_id, parallelism.clone()))
1309            .collect();
1310
1311        self.ddl_controller
1312            .reschedule_fragments(fragment_targets)
1313            .await?;
1314
1315        Ok(Response::new(AlterFragmentParallelismResponse {}))
1316    }
1317
1318    async fn alter_streaming_job_config(
1319        &self,
1320        request: Request<AlterStreamingJobConfigRequest>,
1321    ) -> Result<Response<AlterStreamingJobConfigResponse>, Status> {
1322        let AlterStreamingJobConfigRequest {
1323            job_id,
1324            entries_to_add,
1325            keys_to_remove,
1326        } = request.into_inner();
1327
1328        self.ddl_controller
1329            .run_command(DdlCommand::AlterStreamingJobConfig(
1330                job_id,
1331                entries_to_add,
1332                keys_to_remove,
1333            ))
1334            .await?;
1335
1336        Ok(Response::new(AlterStreamingJobConfigResponse {}))
1337    }
1338
1339    /// Auto schema change for cdc sources,
1340    /// called by the source parser when a schema change is detected.
1341    async fn auto_schema_change(
1342        &self,
1343        request: Request<AutoSchemaChangeRequest>,
1344    ) -> Result<Response<AutoSchemaChangeResponse>, Status> {
1345        let req = request.into_inner();
1346
1347        // randomly select a frontend worker to get the replace table plan
1348        let workers = self
1349            .metadata_manager
1350            .list_worker_node(Some(WorkerType::Frontend), Some(State::Running))
1351            .await?;
1352        let worker = workers
1353            .choose(&mut thread_rng())
1354            .ok_or_else(|| MetaError::from(anyhow!("no frontend worker available")))?;
1355
1356        let client = self
1357            .env
1358            .frontend_client_pool()
1359            .get(worker)
1360            .await
1361            .map_err(MetaError::from)?;
1362
1363        let Some(schema_change) = req.schema_change else {
1364            return Err(Status::invalid_argument(
1365                "schema change message is required",
1366            ));
1367        };
1368
1369        for table_change in schema_change.table_changes {
1370            for c in &table_change.columns {
1371                let c = ColumnCatalog::from(c.clone());
1372
1373                let invalid_col_type = |column_type: &str, c: &ColumnCatalog| {
1374                    tracing::warn!(target: "auto_schema_change",
1375                      cdc_table_id = table_change.cdc_table_id,
1376                    upstraem_ddl = table_change.upstream_ddl,
1377                        "invalid column type from cdc table change");
1378                    Err(Status::invalid_argument(format!(
1379                        "invalid column type: {} from cdc table change, column: {:?}",
1380                        column_type, c
1381                    )))
1382                };
1383                if c.is_generated() {
1384                    return invalid_col_type("generated column", &c);
1385                }
1386                if c.is_rw_sys_column() {
1387                    return invalid_col_type("rw system column", &c);
1388                }
1389                if c.is_hidden {
1390                    return invalid_col_type("hidden column", &c);
1391                }
1392            }
1393
1394            // get the table catalog corresponding to the cdc table
1395            let tables: Vec<Table> = self
1396                .metadata_manager
1397                .get_table_catalog_by_cdc_table_id(&table_change.cdc_table_id)
1398                .await?;
1399
1400            for table in tables {
1401                // Since we only support `ADD` and `DROP` column, we check whether the new columns and the original columns
1402                // is a subset of the other.
1403                let original_columns_by_name: HashMap<String, ColumnCatalog> = table
1404                    .columns
1405                    .iter()
1406                    .filter_map(|col| {
1407                        let col = ColumnCatalog::from(col.clone());
1408                        cdc_auto_schema_change_comparable_column(&col).map(|(name, _)| (name, col))
1409                    })
1410                    .collect();
1411
1412                let original_column_types: HashMap<String, DataType> = original_columns_by_name
1413                    .iter()
1414                    .map(|(name, col)| (name.clone(), col.data_type().clone()))
1415                    .collect();
1416
1417                let original_column_names: HashSet<String> =
1418                    HashSet::from_iter(original_column_types.keys().cloned());
1419
1420                let cdc_table_type =
1421                    PbCdcTableType::try_from(table.cdc_table_type.unwrap_or_default())
1422                        .unwrap_or(PbCdcTableType::Unspecified);
1423                let table_change = normalize_cdc_auto_schema_change_column_names(
1424                    table_change.clone(),
1425                    cdc_table_type,
1426                    &original_columns_by_name,
1427                );
1428
1429                let collect_new_columns = |table_change: &TableSchemaChange| {
1430                    let mut new_columns: HashSet<(String, DataType)> =
1431                        HashSet::from_iter(table_change.columns.iter().filter_map(|col| {
1432                            let col = ColumnCatalog::from(col.clone());
1433                            cdc_auto_schema_change_comparable_column(&col)
1434                        }));
1435
1436                    // For subset/superset check, we need to add visible connector additional columns defined by INCLUDE in the original table to new_columns.
1437                    // This includes both _rw columns and user-defined INCLUDE columns (e.g., INCLUDE TIMESTAMP AS xxx).
1438                    for col in &table.columns {
1439                        let col = ColumnCatalog::from(col.clone());
1440                        if col.is_connector_additional_column()
1441                            && !col.is_hidden()
1442                            && !col.is_generated()
1443                        {
1444                            new_columns
1445                                .insert((col.column_desc.name.clone(), col.data_type().clone()));
1446                        }
1447                    }
1448
1449                    new_columns
1450                };
1451
1452                let new_columns = collect_new_columns(&table_change);
1453
1454                let new_column_names: HashSet<String> =
1455                    HashSet::from_iter(new_columns.iter().map(|(name, _)| name.clone()));
1456                let is_add_or_drop_by_name = original_column_names.is_subset(&new_column_names)
1457                    || original_column_names.is_superset(&new_column_names);
1458
1459                // Debezium schema change events carry the full table schema. Preserve existing
1460                // validator-compatible RW types before both validation and replacement planning.
1461                let table_change =
1462                    if original_column_names != new_column_names && is_add_or_drop_by_name {
1463                        normalize_cdc_auto_schema_change_existing_column_types(
1464                            table_change,
1465                            cdc_table_type,
1466                            &original_columns_by_name,
1467                        )
1468                    } else {
1469                        // Keep the schema change without type normalization for non-add/drop-only
1470                        // cases, such as a mixed add-and-drop change. Existing validation below
1471                        // should reject unsupported schema changes without masking them.
1472                        table_change
1473                    };
1474
1475                let original_columns: HashSet<(String, DataType)> = HashSet::from_iter(
1476                    original_column_types
1477                        .iter()
1478                        .map(|(name, data_type)| (name.clone(), data_type.clone())),
1479                );
1480
1481                let new_columns = collect_new_columns(&table_change);
1482
1483                if !(original_columns.is_subset(&new_columns)
1484                    || original_columns.is_superset(&new_columns))
1485                {
1486                    tracing::warn!(target: "auto_schema_change",
1487                                    table_id = %table.id,
1488                                    cdc_table_id = table.cdc_table_id,
1489                                    upstraem_ddl = table_change.upstream_ddl,
1490                                    original_columns = ?original_columns,
1491                                    new_columns = ?new_columns,
1492                                    "New columns should be a subset or superset of the original columns (including hidden columns), since only `ADD COLUMN` and `DROP COLUMN` is supported");
1493
1494                    let fail_info = "New columns should be a subset or superset of the original columns (including hidden columns), since only `ADD COLUMN` and `DROP COLUMN` is supported".to_owned();
1495                    add_auto_schema_change_fail_event_log(
1496                        &self.meta_metrics,
1497                        table.id,
1498                        table.name.clone(),
1499                        table_change.cdc_table_id.clone(),
1500                        table_change.upstream_ddl.clone(),
1501                        &self.env.event_log_manager_ref(),
1502                        fail_info,
1503                    );
1504
1505                    return Err(Status::invalid_argument(
1506                        "New columns should be a subset or superset of the original columns (including hidden columns)",
1507                    ));
1508                }
1509                // skip the schema change if there is no change to original columns
1510                if original_columns == new_columns {
1511                    tracing::warn!(target: "auto_schema_change",
1512                                   table_id = %table.id,
1513                                   cdc_table_id = table.cdc_table_id,
1514                                   upstraem_ddl = table_change.upstream_ddl,
1515                                    original_columns = ?original_columns,
1516                                    new_columns = ?new_columns,
1517                                   "No change to columns, skipping the schema change");
1518                    continue;
1519                }
1520
1521                let latency_timer = self
1522                    .meta_metrics
1523                    .auto_schema_change_latency
1524                    .with_label_values(&[&table.id.to_string(), &table.name])
1525                    .start_timer();
1526                // send a request to the frontend to get the ReplaceJobPlan
1527                // will retry with exponential backoff if the request fails
1528                let resp = client
1529                    .get_table_replace_plan(GetTableReplacePlanRequest {
1530                        database_id: table.database_id,
1531                        table_id: table.id,
1532                        cdc_table_change: Some(table_change.clone()),
1533                    })
1534                    .await;
1535
1536                match resp {
1537                    Ok(resp) => {
1538                        let resp = resp.into_inner();
1539                        if let Some(plan) = resp.replace_plan {
1540                            let plan = Self::extract_replace_table_info(plan);
1541                            plan.streaming_job.table().inspect(|t| {
1542                                tracing::info!(
1543                                    target: "auto_schema_change",
1544                                    table_id = %t.id,
1545                                    cdc_table_id = t.cdc_table_id,
1546                                    upstraem_ddl = table_change.upstream_ddl,
1547                                    "Start the replace config change")
1548                            });
1549                            // start the schema change procedure
1550                            let replace_res = self
1551                                .ddl_controller
1552                                .run_command(DdlCommand::ReplaceStreamJob(plan))
1553                                .await;
1554
1555                            match replace_res {
1556                                Ok(_) => {
1557                                    tracing::info!(
1558                                        target: "auto_schema_change",
1559                                        table_id = %table.id,
1560                                        cdc_table_id = table.cdc_table_id,
1561                                        "Table replaced success");
1562
1563                                    self.meta_metrics
1564                                        .auto_schema_change_success_cnt
1565                                        .with_label_values(&[&table.id.to_string(), &table.name])
1566                                        .inc();
1567                                    latency_timer.observe_duration();
1568                                }
1569                                Err(e) => {
1570                                    tracing::error!(
1571                                        target: "auto_schema_change",
1572                                        error = %e.as_report(),
1573                                        table_id = %table.id,
1574                                        cdc_table_id = table.cdc_table_id,
1575                                        upstraem_ddl = table_change.upstream_ddl,
1576                                        "failed to replace the table",
1577                                    );
1578                                    let fail_info =
1579                                        format!("failed to replace the table: {}", e.as_report());
1580                                    add_auto_schema_change_fail_event_log(
1581                                        &self.meta_metrics,
1582                                        table.id,
1583                                        table.name.clone(),
1584                                        table_change.cdc_table_id.clone(),
1585                                        table_change.upstream_ddl.clone(),
1586                                        &self.env.event_log_manager_ref(),
1587                                        fail_info,
1588                                    );
1589                                }
1590                            };
1591                        }
1592                    }
1593                    Err(e) => {
1594                        tracing::error!(
1595                            target: "auto_schema_change",
1596                            error = %e.as_report(),
1597                            table_id = %table.id,
1598                            cdc_table_id = table.cdc_table_id,
1599                            "failed to get replace table plan",
1600                        );
1601                        let fail_info =
1602                            format!("failed to get replace table plan: {}", e.as_report());
1603                        add_auto_schema_change_fail_event_log(
1604                            &self.meta_metrics,
1605                            table.id,
1606                            table.name.clone(),
1607                            table_change.cdc_table_id.clone(),
1608                            table_change.upstream_ddl.clone(),
1609                            &self.env.event_log_manager_ref(),
1610                            fail_info,
1611                        );
1612                    }
1613                };
1614            }
1615        }
1616
1617        Ok(Response::new(AutoSchemaChangeResponse {}))
1618    }
1619
1620    async fn wait_iceberg_pk_index_sink_epoch(
1621        &self,
1622        request: Request<WaitIcebergPkIndexSinkEpochRequest>,
1623    ) -> Result<Response<WaitIcebergPkIndexSinkEpochResponse>, Status> {
1624        let req = request.into_inner();
1625        let snapshot_id = self
1626            .iceberg_pk_index_sink_manager
1627            .wait_epoch(req.sink_id, req.epoch)
1628            .await
1629            .map_err(|e| {
1630                Status::internal(format!(
1631                    "Failed to wait for pk-index sink epoch: {}",
1632                    e.as_report()
1633                ))
1634            })?;
1635        Ok(Response::new(WaitIcebergPkIndexSinkEpochResponse {
1636            snapshot_id,
1637        }))
1638    }
1639
1640    async fn alter_swap_rename(
1641        &self,
1642        request: Request<AlterSwapRenameRequest>,
1643    ) -> Result<Response<AlterSwapRenameResponse>, Status> {
1644        let req = request.into_inner();
1645
1646        let version = self
1647            .ddl_controller
1648            .run_command(DdlCommand::AlterSwapRename(req.object.unwrap()))
1649            .await?;
1650
1651        Ok(Response::new(AlterSwapRenameResponse {
1652            status: None,
1653            version,
1654        }))
1655    }
1656
1657    async fn alter_resource_group(
1658        &self,
1659        request: Request<AlterResourceGroupRequest>,
1660    ) -> Result<Response<AlterResourceGroupResponse>, Status> {
1661        let req = request.into_inner();
1662
1663        let job_id = req.get_job_id();
1664        let deferred = req.get_deferred();
1665        let resource_group = req.resource_group;
1666
1667        self.ddl_controller
1668            .reschedule_streaming_job(
1669                job_id,
1670                ReschedulePolicy::ResourceGroup(ResourceGroupPolicy { resource_group }),
1671                deferred,
1672            )
1673            .await?;
1674
1675        Ok(Response::new(AlterResourceGroupResponse {}))
1676    }
1677
1678    async fn alter_database_resource_group(
1679        &self,
1680        request: Request<AlterDatabaseResourceGroupRequest>,
1681    ) -> Result<Response<AlterDatabaseResourceGroupResponse>, Status> {
1682        let req = request.into_inner();
1683
1684        let version = self
1685            .ddl_controller
1686            .run_command(DdlCommand::AlterDatabaseResourceGroup(
1687                req.database_id,
1688                req.resource_group,
1689                req.deferred,
1690            ))
1691            .await?;
1692
1693        Ok(Response::new(AlterDatabaseResourceGroupResponse {
1694            status: None,
1695            version,
1696        }))
1697    }
1698
1699    async fn alter_database_param(
1700        &self,
1701        request: Request<AlterDatabaseParamRequest>,
1702    ) -> Result<Response<AlterDatabaseParamResponse>, Status> {
1703        let req = request.into_inner();
1704        let database_id = req.database_id;
1705
1706        let param = match req.param.unwrap() {
1707            alter_database_param_request::Param::BarrierIntervalMs(value) => {
1708                AlterDatabaseParam::BarrierIntervalMs(value.value)
1709            }
1710            alter_database_param_request::Param::CheckpointFrequency(value) => {
1711                AlterDatabaseParam::CheckpointFrequency(value.value)
1712            }
1713        };
1714        let version = self
1715            .ddl_controller
1716            .run_command(DdlCommand::AlterDatabaseParam(database_id, param))
1717            .await?;
1718
1719        return Ok(Response::new(AlterDatabaseParamResponse {
1720            status: None,
1721            version,
1722        }));
1723    }
1724
1725    async fn compact_iceberg_table(
1726        &self,
1727        request: Request<CompactIcebergTableRequest>,
1728    ) -> Result<Response<CompactIcebergTableResponse>, Status> {
1729        let req = request.into_inner();
1730        let sink_id = req.sink_id;
1731
1732        // Trigger manual compaction directly using the sink ID
1733        let task_id = self
1734            .iceberg_compaction_manager
1735            .trigger_manual_compaction(sink_id)
1736            .await
1737            .map_err(|e| {
1738                Status::internal(format!("Failed to trigger compaction: {}", e.as_report()))
1739            })?;
1740
1741        Ok(Response::new(CompactIcebergTableResponse {
1742            status: None,
1743            task_id,
1744        }))
1745    }
1746
1747    async fn expire_iceberg_table_snapshots(
1748        &self,
1749        request: Request<ExpireIcebergTableSnapshotsRequest>,
1750    ) -> Result<Response<ExpireIcebergTableSnapshotsResponse>, Status> {
1751        let req = request.into_inner();
1752        let sink_id = req.sink_id;
1753
1754        // Trigger manual snapshot expiration directly using the sink ID
1755        self.iceberg_compaction_manager
1756            .check_and_expire_snapshots(sink_id)
1757            .await
1758            .map_err(|e| {
1759                Status::internal(format!("Failed to expire snapshots: {}", e.as_report()))
1760            })?;
1761
1762        Ok(Response::new(ExpireIcebergTableSnapshotsResponse {
1763            status: None,
1764        }))
1765    }
1766
1767    async fn rewrite_iceberg_table_manifests(
1768        &self,
1769        request: Request<RewriteIcebergTableManifestsRequest>,
1770    ) -> Result<Response<RewriteIcebergTableManifestsResponse>, Status> {
1771        let req = request.into_inner();
1772        let sink_id = req.sink_id;
1773
1774        self.iceberg_compaction_manager
1775            .check_and_rewrite_manifests(sink_id)
1776            .await
1777            .map_err(|e| {
1778                Status::internal(format!(
1779                    "Failed to rewrite manifests for sink {}: {}",
1780                    sink_id,
1781                    e.as_report()
1782                ))
1783            })?;
1784
1785        Ok(Response::new(RewriteIcebergTableManifestsResponse {
1786            status: None,
1787        }))
1788    }
1789
1790    async fn create_iceberg_table(
1791        &self,
1792        request: Request<CreateIcebergTableRequest>,
1793    ) -> Result<Response<CreateIcebergTableResponse>, Status> {
1794        let req = request.into_inner();
1795        let CreateIcebergTableRequest {
1796            table_info,
1797            sink_info,
1798            iceberg_source,
1799            if_not_exists,
1800        } = req;
1801
1802        // 1. create table job
1803        let PbTableJobInfo {
1804            source,
1805            table,
1806            fragment_graph,
1807            job_type,
1808        } = table_info.unwrap();
1809        let mut table = table.unwrap();
1810        let mut fragment_graph = fragment_graph.unwrap();
1811        let database_id = table.get_database_id();
1812        let schema_id = table.get_schema_id();
1813        let table_name = table.get_name().to_owned();
1814
1815        // Mark table as background creation, so that it won't block sink creation.
1816        table.create_type = PbCreateType::Background as _;
1817
1818        // Set the source rate limit to 0 and reset it back after the iceberg sink is backfilling.
1819        let source_rate_limit = if let Some(source) = &source {
1820            for fragment in fragment_graph.fragments.values_mut() {
1821                stream_graph_visitor::visit_fragment_mut(fragment, |node| {
1822                    if let NodeBody::Source(source_node) = node
1823                        && let Some(inner) = &mut source_node.source_inner
1824                    {
1825                        inner.rate_limit = Some(0);
1826                    }
1827                });
1828            }
1829            Some(source.rate_limit)
1830        } else {
1831            None
1832        };
1833
1834        let stream_job =
1835            StreamingJob::Table(source, table, PbTableJobType::try_from(job_type).unwrap());
1836        let _ = self
1837            .ddl_controller
1838            .run_command(DdlCommand::CreateStreamingJob {
1839                stream_job,
1840                fragment_graph,
1841                dependencies: HashSet::new(),
1842                resource_type: Self::default_streaming_job_resource_type(),
1843                if_not_exists,
1844                refresh_interval_sec: None,
1845                replace_sink: None,
1846                since_timestamp_epoch: None,
1847            })
1848            .await?;
1849
1850        let table_catalog = self
1851            .metadata_manager
1852            .catalog_controller
1853            .get_table_catalog_by_name(database_id, schema_id, &table_name)
1854            .await?
1855            .ok_or(Status::not_found("Internal error: table not found"))?;
1856
1857        // 2. create iceberg sink job
1858        let PbSinkJobInfo {
1859            sink,
1860            fragment_graph,
1861        } = sink_info.unwrap();
1862        let mut sink = sink.unwrap();
1863
1864        // Mark sink as background creation, so that it won't block source creation.
1865        sink.create_type = PbCreateType::Background as _;
1866
1867        let enable_pk_index = sink
1868            .properties
1869            .get(ENABLE_PK_INDEX)
1870            .is_some_and(|v| v.eq_ignore_ascii_case("true"));
1871        // The internal iceberg sink is planned before the table catalog exists, so this field
1872        // still carries a placeholder table id when the request reaches meta.
1873        sink.auto_refresh_schema_from_table = if enable_pk_index {
1874            None
1875        } else {
1876            Some(table_catalog.id)
1877        };
1878
1879        let mut fragment_graph = fragment_graph.unwrap();
1880
1881        assert_eq!(fragment_graph.dependent_table_ids.len(), 1);
1882        assert!(
1883            risingwave_common::catalog::TableId::from(fragment_graph.dependent_table_ids[0])
1884                .is_placeholder()
1885        );
1886        fragment_graph.dependent_table_ids[0] = table_catalog.id;
1887        for fragment in fragment_graph.fragments.values_mut() {
1888            stream_graph_visitor::visit_fragment_mut(fragment, |node| match node {
1889                NodeBody::StreamScan(scan) => {
1890                    scan.table_id = table_catalog.id;
1891                    if let Some(table_desc) = &mut scan.table_desc {
1892                        assert!(
1893                            risingwave_common::catalog::TableId::from(table_desc.table_id)
1894                                .is_placeholder()
1895                        );
1896                        table_desc.table_id = table_catalog.id;
1897                        table_desc.maybe_vnode_count = table_catalog.maybe_vnode_count;
1898                    }
1899                    if let Some(table) = &mut scan.arrangement_table {
1900                        assert!(
1901                            risingwave_common::catalog::TableId::from(table.id).is_placeholder()
1902                        );
1903                        *table = table_catalog.clone();
1904                    }
1905                }
1906                NodeBody::BatchPlan(plan) => {
1907                    if let Some(table_desc) = &mut plan.table_desc {
1908                        assert!(
1909                            risingwave_common::catalog::TableId::from(table_desc.table_id)
1910                                .is_placeholder()
1911                        );
1912                        table_desc.table_id = table_catalog.id;
1913                        table_desc.maybe_vnode_count = table_catalog.maybe_vnode_count;
1914                    }
1915                }
1916                _ => {}
1917            });
1918        }
1919
1920        let table_id = table_catalog.id;
1921        let dependencies = HashSet::from_iter([table_id.into(), schema_id.into()]);
1922        let stream_job = StreamingJob::Sink(sink, Some(table_id));
1923        let res = self
1924            .ddl_controller
1925            .run_command(DdlCommand::CreateStreamingJob {
1926                stream_job,
1927                fragment_graph,
1928                dependencies,
1929                resource_type: Self::default_streaming_job_resource_type(),
1930                if_not_exists,
1931                refresh_interval_sec: None,
1932                replace_sink: None,
1933                since_timestamp_epoch: None,
1934            })
1935            .await;
1936
1937        if res.is_err() {
1938            let _ = self
1939                .ddl_controller
1940                .run_command(DdlCommand::DropStreamingJob {
1941                    job_id: StreamingJobId::Table(None, table_id),
1942                    drop_mode: DropMode::Cascade,
1943                })
1944                .await
1945                .inspect_err(|err| {
1946                    tracing::error!(error = %err.as_report(),
1947                        "Failed to clean up table after iceberg sink creation failure",
1948                    );
1949                });
1950            res?;
1951        }
1952
1953        // 3. reset source rate limit back to normal after sink creation
1954        if let Some(source_rate_limit) = source_rate_limit
1955            && source_rate_limit != Some(0)
1956        {
1957            let OptionalAssociatedSourceId::AssociatedSourceId(source_id) =
1958                table_catalog.optional_associated_source_id.unwrap();
1959            let fragment_nodes = self
1960                .metadata_manager
1961                .update_source_rate_limit_by_source_id(source_id, source_rate_limit)
1962                .await?;
1963            let throttle_config = ThrottleConfig {
1964                throttle_type: risingwave_pb::common::ThrottleType::Source.into(),
1965                rate_limit: source_rate_limit,
1966            };
1967            let config = fragment_nodes
1968                .into_iter()
1969                .map(|(fragment_id, stream_node)| (fragment_id, (throttle_config, stream_node)))
1970                .collect();
1971            let _ = self
1972                .barrier_scheduler
1973                .run_command(database_id, Command::Throttle { config })
1974                .await?;
1975        }
1976
1977        // 4. create iceberg source
1978        let iceberg_source = iceberg_source.unwrap();
1979        let res = self
1980            .ddl_controller
1981            .run_command(DdlCommand::CreateNonSharedSource(
1982                iceberg_source,
1983                Some(table_id),
1984            ))
1985            .await;
1986        if res.is_err() {
1987            let _ = self
1988                .ddl_controller
1989                .run_command(DdlCommand::DropStreamingJob {
1990                    job_id: StreamingJobId::Table(None, table_id),
1991                    drop_mode: DropMode::Cascade,
1992                })
1993                .await
1994                .inspect_err(|err| {
1995                    tracing::error!(
1996                        error = %err.as_report(),
1997                        "Failed to clean up table after iceberg source creation failure",
1998                    );
1999                });
2000        }
2001
2002        Ok(Response::new(CreateIcebergTableResponse {
2003            status: None,
2004            version: res?,
2005        }))
2006    }
2007}
2008
2009fn add_auto_schema_change_fail_event_log(
2010    meta_metrics: &MetaMetrics,
2011    table_id: TableId,
2012    table_name: String,
2013    cdc_table_id: String,
2014    upstream_ddl: String,
2015    event_log_manager: &EventLogManagerRef,
2016    fail_info: String,
2017) {
2018    meta_metrics
2019        .auto_schema_change_failure_cnt
2020        .with_label_values(&[&table_id.to_string(), &table_name])
2021        .inc();
2022    let event = event_log::EventAutoSchemaChangeFail {
2023        table_id,
2024        table_name,
2025        cdc_table_id,
2026        upstream_ddl,
2027        fail_info,
2028    };
2029    event_log_manager.add_event_logs(vec![event_log::Event::AutoSchemaChangeFail(event)]);
2030}
2031
2032fn cdc_auto_schema_change_comparable_column(column: &ColumnCatalog) -> Option<(String, DataType)> {
2033    if column.is_generated() || column.is_hidden() {
2034        None
2035    } else {
2036        Some((column.column_desc.name.clone(), column.data_type().clone()))
2037    }
2038}
2039
2040fn normalize_cdc_auto_schema_change_column_names(
2041    mut table_change: TableSchemaChange,
2042    cdc_table_type: PbCdcTableType,
2043    original_columns_by_name: &HashMap<String, ColumnCatalog>,
2044) -> TableSchemaChange {
2045    if cdc_table_type != PbCdcTableType::Mysql {
2046        return table_change;
2047    }
2048
2049    // MySQL column names are case-insensitive, while RisingWave catalog names may preserve
2050    // explicitly quoted spelling. Reuse that spelling for existing columns and apply the same
2051    // lowercase convention as MySQL snapshot discovery only to genuinely new columns.
2052    let mut original_names_by_lowercase = HashMap::<String, Option<String>>::new();
2053    for original_name in original_columns_by_name.keys() {
2054        original_names_by_lowercase
2055            .entry(original_name.to_lowercase())
2056            .and_modify(|name| *name = None)
2057            .or_insert_with(|| Some(original_name.clone()));
2058    }
2059
2060    for column in &mut table_change.columns {
2061        let mut column_catalog = ColumnCatalog::from(column.clone());
2062        if column_catalog.is_generated() || column_catalog.is_hidden() {
2063            continue;
2064        }
2065
2066        let incoming_name = &column_catalog.column_desc.name;
2067        let lowercase_name = incoming_name.to_lowercase();
2068        let normalized_name = if original_columns_by_name.contains_key(incoming_name) {
2069            incoming_name.clone()
2070        } else {
2071            original_names_by_lowercase
2072                .get(&lowercase_name)
2073                .and_then(|name| name.clone())
2074                .unwrap_or(lowercase_name)
2075        };
2076
2077        column_catalog.column_desc.name = normalized_name;
2078        *column = column_catalog.to_protobuf();
2079    }
2080
2081    table_change
2082}
2083
2084fn normalize_cdc_auto_schema_change_existing_column_types(
2085    mut table_change: TableSchemaChange,
2086    cdc_table_type: PbCdcTableType,
2087    original_columns_by_name: &HashMap<String, ColumnCatalog>,
2088) -> TableSchemaChange {
2089    for column in &mut table_change.columns {
2090        let mut column_catalog = ColumnCatalog::from(column.clone());
2091        if column_catalog.is_generated() || column_catalog.is_hidden() {
2092            continue;
2093        }
2094
2095        if let Some(original_column) =
2096            original_columns_by_name.get(&column_catalog.column_desc.name)
2097            && cdc_auto_schema_change_existing_type_compatible(
2098                cdc_table_type,
2099                original_column.data_type(),
2100                column_catalog.data_type(),
2101            )
2102        {
2103            column_catalog.column_desc.data_type = original_column.data_type().clone();
2104            column_catalog.column_desc.generated_or_default_column = original_column
2105                .column_desc
2106                .generated_or_default_column
2107                .clone();
2108            *column = column_catalog.to_protobuf();
2109        }
2110    }
2111
2112    table_change
2113}
2114
2115#[cfg(test)]
2116mod tests {
2117    use risingwave_common::catalog::{ColumnDesc, ColumnId};
2118    use risingwave_common::types::ScalarImpl;
2119    use risingwave_pb::ddl_service::table_schema_change::TableChangeType;
2120    use risingwave_pb::plan_common::column_desc::GeneratedOrDefaultColumn;
2121
2122    use super::*;
2123
2124    fn pb_column(name: &str, data_type: DataType) -> risingwave_pb::plan_common::ColumnCatalog {
2125        ColumnCatalog::visible(ColumnDesc::named(name, ColumnId::placeholder(), data_type))
2126            .to_protobuf()
2127    }
2128
2129    fn pb_column_with_default_value(
2130        name: &str,
2131        data_type: DataType,
2132        snapshot_value: ScalarImpl,
2133    ) -> risingwave_pb::plan_common::ColumnCatalog {
2134        ColumnCatalog::visible(ColumnDesc::named_with_default_value(
2135            name,
2136            ColumnId::placeholder(),
2137            data_type,
2138            Some(snapshot_value),
2139        ))
2140        .to_protobuf()
2141    }
2142
2143    fn pb_table_change(
2144        columns: Vec<risingwave_pb::plan_common::ColumnCatalog>,
2145    ) -> TableSchemaChange {
2146        TableSchemaChange {
2147            change_type: TableChangeType::Alter as _,
2148            cdc_table_id: "1.db.t".to_owned(),
2149            columns,
2150            upstream_ddl: "ALTER TABLE t ADD COLUMN note VARCHAR(255)".to_owned(),
2151        }
2152    }
2153
2154    #[test]
2155    fn test_cdc_auto_schema_change_normalizes_compatible_existing_column_types() {
2156        let original_columns_by_name = HashMap::from([
2157            (
2158                "id".to_owned(),
2159                ColumnCatalog::visible(ColumnDesc::named(
2160                    "id",
2161                    ColumnId::placeholder(),
2162                    DataType::Int64,
2163                )),
2164            ),
2165            (
2166                "v".to_owned(),
2167                ColumnCatalog::visible(ColumnDesc::named(
2168                    "v",
2169                    ColumnId::placeholder(),
2170                    DataType::Varchar,
2171                )),
2172            ),
2173        ]);
2174        let table_change = pb_table_change(vec![
2175            pb_column("id", DataType::Int32),
2176            pb_column("v", DataType::Varchar),
2177            pb_column("note", DataType::Varchar),
2178        ]);
2179
2180        let normalized = normalize_cdc_auto_schema_change_existing_column_types(
2181            table_change,
2182            PbCdcTableType::Mysql,
2183            &original_columns_by_name,
2184        );
2185        let columns = normalized
2186            .columns
2187            .into_iter()
2188            .map(ColumnCatalog::from)
2189            .map(|column| (column.column_desc.name, column.column_desc.data_type))
2190            .collect::<HashMap<_, _>>();
2191
2192        assert_eq!(columns["id"], DataType::Int64);
2193        assert_eq!(columns["v"], DataType::Varchar);
2194        assert_eq!(columns["note"], DataType::Varchar);
2195    }
2196
2197    #[test]
2198    fn test_mysql_cdc_auto_schema_change_normalizes_column_names() {
2199        let original_columns_by_name = HashMap::from([
2200            (
2201                "Id".to_owned(),
2202                ColumnCatalog::visible(ColumnDesc::named(
2203                    "Id",
2204                    ColumnId::placeholder(),
2205                    DataType::Int64,
2206                )),
2207            ),
2208            (
2209                "Details".to_owned(),
2210                ColumnCatalog::visible(ColumnDesc::named(
2211                    "Details",
2212                    ColumnId::placeholder(),
2213                    DataType::Varchar,
2214                )),
2215            ),
2216        ]);
2217        let table_change = pb_table_change(vec![
2218            pb_column("ID", DataType::Int32),
2219            pb_column("details", DataType::Varchar),
2220            pb_column("NewCol", DataType::Varchar),
2221        ]);
2222
2223        let normalized = normalize_cdc_auto_schema_change_column_names(
2224            table_change,
2225            PbCdcTableType::Mysql,
2226            &original_columns_by_name,
2227        );
2228        let normalized = normalize_cdc_auto_schema_change_existing_column_types(
2229            normalized,
2230            PbCdcTableType::Mysql,
2231            &original_columns_by_name,
2232        );
2233        let columns = normalized
2234            .columns
2235            .into_iter()
2236            .map(ColumnCatalog::from)
2237            .map(|column| (column.column_desc.name, column.column_desc.data_type))
2238            .collect::<Vec<_>>();
2239
2240        assert_eq!(
2241            columns,
2242            vec![
2243                ("Id".to_owned(), DataType::Int64),
2244                ("Details".to_owned(), DataType::Varchar),
2245                ("newcol".to_owned(), DataType::Varchar),
2246            ]
2247        );
2248    }
2249
2250    #[test]
2251    fn test_non_mysql_cdc_auto_schema_change_preserves_column_names() {
2252        let original_columns_by_name = HashMap::from([(
2253            "Id".to_owned(),
2254            ColumnCatalog::visible(ColumnDesc::named(
2255                "Id",
2256                ColumnId::placeholder(),
2257                DataType::Int64,
2258            )),
2259        )]);
2260        let table_change = pb_table_change(vec![
2261            pb_column("ID", DataType::Int64),
2262            pb_column("NewCol", DataType::Varchar),
2263        ]);
2264
2265        let normalized = normalize_cdc_auto_schema_change_column_names(
2266            table_change,
2267            PbCdcTableType::Postgres,
2268            &original_columns_by_name,
2269        );
2270        let column_names = normalized
2271            .columns
2272            .into_iter()
2273            .map(ColumnCatalog::from)
2274            .map(|column| column.column_desc.name)
2275            .collect::<Vec<_>>();
2276
2277        assert_eq!(column_names, vec!["ID".to_owned(), "NewCol".to_owned()]);
2278    }
2279
2280    #[test]
2281    fn test_cdc_auto_schema_change_preserves_default_when_normalizing_type() {
2282        let original = ColumnCatalog::visible(ColumnDesc::named_with_default_value(
2283            "v",
2284            ColumnId::placeholder(),
2285            DataType::Int64,
2286            Some(ScalarImpl::Int64(1)),
2287        ));
2288        let original_columns_by_name = HashMap::from([("v".to_owned(), original.clone())]);
2289        let table_change = pb_table_change(vec![pb_column_with_default_value(
2290            "v",
2291            DataType::Int32,
2292            ScalarImpl::Int32(1),
2293        )]);
2294
2295        let normalized = normalize_cdc_auto_schema_change_existing_column_types(
2296            table_change,
2297            PbCdcTableType::Mysql,
2298            &original_columns_by_name,
2299        );
2300        let column = ColumnCatalog::from(normalized.columns.into_iter().next().unwrap());
2301
2302        assert_eq!(column.column_desc.data_type, DataType::Int64);
2303        assert_eq!(
2304            column.column_desc.generated_or_default_column,
2305            original.column_desc.generated_or_default_column
2306        );
2307        assert!(matches!(
2308            column.column_desc.generated_or_default_column,
2309            Some(GeneratedOrDefaultColumn::DefaultColumn(_))
2310        ));
2311    }
2312}