Skip to main content

risingwave_meta/rpc/
ddl_controller.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::cmp::Ordering;
16use std::collections::{HashMap, HashSet};
17use std::num::NonZeroUsize;
18use std::sync::Arc;
19use std::sync::atomic::AtomicU64;
20use std::time::Duration;
21
22use anyhow::{Context, anyhow};
23use await_tree::InstrumentAwait;
24use either::Either;
25use itertools::Itertools;
26use risingwave_common::catalog::{
27    AlterDatabaseParam, ColumnCatalog, ColumnId, Field, FragmentTypeFlag,
28};
29use risingwave_common::hash::VnodeCountCompat;
30use risingwave_common::id::{JobId, TableId};
31use risingwave_common::secret::{LocalSecretManager, SecretEncryption};
32use risingwave_common::system_param::adaptive_parallelism_strategy::parse_strategy;
33use risingwave_common::system_param::reader::SystemParamsRead;
34use risingwave_common::util::stream_graph_visitor::visit_stream_node_cont_mut;
35use risingwave_common::{bail, bail_not_implemented};
36use risingwave_connector::WithOptionsSecResolved;
37use risingwave_connector::connector_common::validate_connection;
38use risingwave_connector::sink::SinkParam;
39use risingwave_connector::sink::iceberg::IcebergSink;
40use risingwave_connector::source::cdc::CdcScanOptions;
41use risingwave_connector::source::{
42    ConnectorProperties, SourceEnumeratorContext, UPSTREAM_SOURCE_KEY,
43};
44use risingwave_meta_model::object::ObjectType;
45use risingwave_meta_model::refresh_job::RefreshState;
46use risingwave_meta_model::{
47    ConnectionId, DatabaseId, DispatcherType, FragmentId, FunctionId, IndexId, JobStatus, ObjectId,
48    SchemaId, SecretId, SinkId, SourceId, StreamingParallelism, SubscriptionId, UserId, ViewId,
49    streaming_job,
50};
51use risingwave_pb::catalog::{
52    Comment, Connection, CreateType, Database, Function, PbTable, Schema, Secret, Source,
53    Subscription, Table, View,
54};
55use risingwave_pb::ddl_service::alter_owner_request::Object;
56use risingwave_pb::ddl_service::{
57    DdlProgress, TableJobType, WaitVersion, alter_name_request, alter_set_schema_request,
58    alter_swap_rename_request, streaming_job_resource_type,
59};
60use risingwave_pb::meta::table_fragments::fragment::FragmentDistributionType as PbFragmentDistributionType;
61use risingwave_pb::plan_common::{PbColumnCatalog, PbExternalTableDesc};
62use risingwave_pb::stream_plan::stream_node::NodeBody;
63use risingwave_pb::stream_plan::{
64    PbDispatchOutputMapping, PbStreamFragmentGraph, PbStreamNode, PbUpstreamSinkInfo,
65    StreamFragmentGraph as StreamFragmentGraphProto,
66};
67use risingwave_pb::telemetry::{PbTelemetryDatabaseObject, PbTelemetryEventStage};
68use strum::Display;
69use thiserror_ext::AsReport;
70use tokio::sync::Semaphore;
71use tokio::time::sleep;
72use tracing::Instrument;
73
74use crate::barrier::{BarrierManagerRef, Command};
75use crate::controller::catalog::{DropTableConnectorContext, ReleaseContext};
76use crate::controller::streaming_job::{FinishAutoRefreshSchemaSinkContext, SinkIntoTableContext};
77use crate::controller::utils::build_select_node_list;
78use crate::error::{MetaErrorInner, bail_invalid_parameter};
79use crate::manager::iceberg_compaction::IcebergCompactionManagerRef;
80use crate::manager::iceberg_pk_index_sink::IcebergPkIndexSinkManager;
81use crate::manager::sink_coordination::SinkCoordinatorManager;
82use crate::manager::{
83    IGNORED_NOTIFICATION_VERSION, LocalNotification, MetaSrvEnv, MetadataManager,
84    NotificationVersion, StreamingJob, StreamingJobType,
85};
86use crate::model::{
87    DownstreamFragmentRelation, FragmentDownstreamRelation, FragmentId as CatalogFragmentId,
88    StreamContext, StreamJobFragments, StreamJobFragmentsToCreate,
89};
90use crate::stream::cdc::{
91    parallel_cdc_table_backfill_fragment, try_init_parallel_cdc_table_snapshot_splits,
92};
93use crate::stream::{
94    ActorGraphBuildResult, ActorGraphBuilder, AutoRefreshSchemaSinkContext,
95    CompleteStreamFragmentGraph, CreateStreamingJobContext, CreateStreamingJobOption,
96    FragmentGraphDownstreamContext, FragmentGraphUpstreamContext, GlobalStreamManagerRef,
97    ParallelismPolicy, ReplaceStreamJobContext, ReschedulePolicy, SourceChange, SourceManagerRef,
98    StreamFragmentGraph, UpstreamSinkInfo, check_sink_fragments_support_refresh_schema,
99    cleanup_dropped_streaming_jobs, create_source_worker, first_variant_column,
100    rewrite_refresh_schema_sink_fragment, state_match, validate_sink,
101};
102use crate::telemetry::report_event;
103use crate::{MetaError, MetaResult};
104
105#[derive(PartialEq)]
106pub enum DropMode {
107    Restrict,
108    Cascade,
109}
110
111impl DropMode {
112    pub fn from_request_setting(cascade: bool) -> DropMode {
113        if cascade {
114            DropMode::Cascade
115        } else {
116            DropMode::Restrict
117        }
118    }
119}
120
121#[derive(strum::AsRefStr)]
122pub enum StreamingJobId {
123    MaterializedView(TableId),
124    Sink(SinkId),
125    Table(Option<SourceId>, TableId),
126    Index(IndexId),
127}
128
129impl std::fmt::Display for StreamingJobId {
130    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
131        write!(f, "{}", self.as_ref())?;
132        write!(f, "({})", self.id())
133    }
134}
135
136impl StreamingJobId {
137    fn id(&self) -> JobId {
138        match self {
139            StreamingJobId::MaterializedView(id) | StreamingJobId::Table(_, id) => id.as_job_id(),
140            StreamingJobId::Index(id) => id.as_job_id(),
141            StreamingJobId::Sink(id) => id.as_job_id(),
142        }
143    }
144}
145
146/// It’s used to describe the information of the job that needs to be replaced
147/// and it will be used during replacing table and creating sink into table operations.
148pub struct ReplaceStreamJobInfo {
149    pub streaming_job: StreamingJob,
150    pub fragment_graph: StreamFragmentGraphProto,
151}
152
153#[derive(Display)]
154pub enum DdlCommand {
155    CreateDatabase(Database),
156    DropDatabase(DatabaseId),
157    CreateSchema(Schema),
158    DropSchema(SchemaId, DropMode),
159    CreateNonSharedSource(Source, Option<TableId>),
160    DropSource(SourceId, DropMode),
161    ResetSource(SourceId),
162    CreateFunction(Function),
163    DropFunction(FunctionId, DropMode),
164    CreateView(View, HashSet<ObjectId>),
165    DropView(ViewId, DropMode),
166    CreateStreamingJob {
167        stream_job: StreamingJob,
168        fragment_graph: StreamFragmentGraphProto,
169        dependencies: HashSet<ObjectId>,
170        resource_type: streaming_job_resource_type::ResourceType,
171        if_not_exists: bool,
172        refresh_interval_sec: Option<u64>,
173        replace_sink: Option<SinkId>,
174        since_timestamp_epoch: Option<u64>,
175    },
176    DropStreamingJob {
177        job_id: StreamingJobId,
178        drop_mode: DropMode,
179    },
180    AlterName(alter_name_request::Object, String),
181    AlterSwapRename(alter_swap_rename_request::Object),
182    ReplaceStreamJob(ReplaceStreamJobInfo),
183    AlterNonSharedSource(Source),
184    AlterObjectOwner(Object, UserId),
185    AlterSetSchema(alter_set_schema_request::Object, SchemaId),
186    CreateConnection(Connection),
187    DropConnection(ConnectionId, DropMode),
188    CreateSecret(Secret),
189    AlterSecret(Secret),
190    DropSecret(SecretId, DropMode),
191    CommentOn(Comment),
192    CreateSubscription(Subscription),
193    DropSubscription(SubscriptionId, DropMode),
194    AlterSubscriptionRetention {
195        subscription_id: SubscriptionId,
196        retention_seconds: u64,
197        definition: String,
198    },
199    AlterDatabaseParam(DatabaseId, AlterDatabaseParam),
200    AlterDatabaseResourceGroup(DatabaseId, Option<String>, bool),
201    AlterStreamingJobConfig(JobId, HashMap<String, String>, Vec<String>),
202}
203
204impl DdlCommand {
205    /// Returns the name or ID of the object that this command operates on, for observability and debugging.
206    fn object(&self) -> Either<String, ObjectId> {
207        use Either::*;
208        match self {
209            DdlCommand::CreateDatabase(database) => Left(database.name.clone()),
210            DdlCommand::DropDatabase(id) => Right(id.as_object_id()),
211            DdlCommand::CreateSchema(schema) => Left(schema.name.clone()),
212            DdlCommand::DropSchema(id, _) => Right(id.as_object_id()),
213            DdlCommand::CreateNonSharedSource(source, _) => Left(source.name.clone()),
214            DdlCommand::DropSource(id, _) => Right(id.as_object_id()),
215            DdlCommand::ResetSource(id) => Right(id.as_object_id()),
216            DdlCommand::CreateFunction(function) => Left(function.name.clone()),
217            DdlCommand::DropFunction(id, _) => Right(id.as_object_id()),
218            DdlCommand::CreateView(view, _) => Left(view.name.clone()),
219            DdlCommand::DropView(id, _) => Right(id.as_object_id()),
220            DdlCommand::CreateStreamingJob { stream_job, .. } => Left(stream_job.name()),
221            DdlCommand::DropStreamingJob { job_id, .. } => Right(job_id.id().as_object_id()),
222            DdlCommand::AlterName(object, _) => Left(format!("{object:?}")),
223            DdlCommand::AlterSwapRename(object) => Left(format!("{object:?}")),
224            DdlCommand::ReplaceStreamJob(info) => Left(info.streaming_job.name()),
225            DdlCommand::AlterNonSharedSource(source) => Left(source.name.clone()),
226            DdlCommand::AlterObjectOwner(object, _) => Left(format!("{object:?}")),
227            DdlCommand::AlterSetSchema(object, _) => Left(format!("{object:?}")),
228            DdlCommand::CreateConnection(connection) => Left(connection.name.clone()),
229            DdlCommand::DropConnection(id, _) => Right(id.as_object_id()),
230            DdlCommand::CreateSecret(secret) => Left(secret.name.clone()),
231            DdlCommand::AlterSecret(secret) => Left(secret.name.clone()),
232            DdlCommand::DropSecret(id, _) => Right(id.as_object_id()),
233            DdlCommand::CommentOn(comment) => Right(comment.table_id.into()),
234            DdlCommand::CreateSubscription(subscription) => Left(subscription.name.clone()),
235            DdlCommand::DropSubscription(id, _) => Right(id.as_object_id()),
236            DdlCommand::AlterSubscriptionRetention {
237                subscription_id, ..
238            } => Right(subscription_id.as_object_id()),
239            DdlCommand::AlterDatabaseParam(id, _) => Right(id.as_object_id()),
240            DdlCommand::AlterDatabaseResourceGroup(id, _, _) => Right(id.as_object_id()),
241            DdlCommand::AlterStreamingJobConfig(job_id, _, _) => Right(job_id.as_object_id()),
242        }
243    }
244
245    fn allow_in_recovery(&self) -> bool {
246        match self {
247            DdlCommand::DropDatabase(_)
248            | DdlCommand::DropSchema(_, _)
249            | DdlCommand::DropSource(_, _)
250            | DdlCommand::DropFunction(_, _)
251            | DdlCommand::DropView(_, _)
252            | DdlCommand::DropStreamingJob { .. }
253            | DdlCommand::DropConnection(_, _)
254            | DdlCommand::DropSecret(_, _)
255            | DdlCommand::DropSubscription(_, _)
256            | DdlCommand::AlterName(_, _)
257            | DdlCommand::AlterObjectOwner(_, _)
258            | DdlCommand::AlterSetSchema(_, _)
259            | DdlCommand::CreateDatabase(_)
260            | DdlCommand::CreateSchema(_)
261            | DdlCommand::CreateFunction(_)
262            | DdlCommand::CreateView(_, _)
263            | DdlCommand::CreateConnection(_)
264            | DdlCommand::CommentOn(_)
265            | DdlCommand::CreateSecret(_)
266            | DdlCommand::AlterSecret(_)
267            | DdlCommand::AlterSwapRename(_)
268            | DdlCommand::AlterDatabaseParam(_, _)
269            | DdlCommand::AlterDatabaseResourceGroup(_, _, _)
270            | DdlCommand::AlterStreamingJobConfig(_, _, _)
271            | DdlCommand::AlterSubscriptionRetention { .. } => true,
272            DdlCommand::CreateStreamingJob { .. }
273            | DdlCommand::CreateNonSharedSource(_, _)
274            | DdlCommand::ReplaceStreamJob(_)
275            | DdlCommand::AlterNonSharedSource(_)
276            | DdlCommand::ResetSource(_)
277            | DdlCommand::CreateSubscription(_) => false,
278        }
279    }
280}
281
282#[derive(Clone)]
283pub struct DdlController {
284    pub(crate) env: MetaSrvEnv,
285
286    pub(crate) metadata_manager: MetadataManager,
287    pub(crate) stream_manager: GlobalStreamManagerRef,
288    pub(crate) source_manager: SourceManagerRef,
289    barrier_manager: BarrierManagerRef,
290    sink_manager: SinkCoordinatorManager,
291    iceberg_compaction_manager: IcebergCompactionManagerRef,
292    iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
293
294    // The semaphore is used to limit the number of concurrent streaming job creation.
295    pub(crate) creating_streaming_job_permits: Arc<CreatingStreamingJobPermit>,
296
297    /// Sequence number for DDL commands, used for observability and debugging.
298    seq: Arc<AtomicU64>,
299}
300
301#[derive(Clone)]
302pub struct CreatingStreamingJobPermit {
303    pub(crate) semaphore: Arc<Semaphore>,
304}
305
306impl CreatingStreamingJobPermit {
307    async fn new(env: &MetaSrvEnv) -> Self {
308        let mut permits = env
309            .system_params_reader()
310            .await
311            .max_concurrent_creating_streaming_jobs() as usize;
312        if permits == 0 {
313            // if the system parameter is set to zero, use the max permitted value.
314            permits = Semaphore::MAX_PERMITS;
315        }
316        let semaphore = Arc::new(Semaphore::new(permits));
317
318        let (local_notification_tx, mut local_notification_rx) =
319            tokio::sync::mpsc::unbounded_channel();
320        env.notification_manager()
321            .insert_local_sender(local_notification_tx);
322        let semaphore_clone = semaphore.clone();
323        tokio::spawn(async move {
324            while let Some(notification) = local_notification_rx.recv().await {
325                let LocalNotification::SystemParamsChange(p) = &notification else {
326                    continue;
327                };
328                let mut new_permits = p.max_concurrent_creating_streaming_jobs() as usize;
329                if new_permits == 0 {
330                    new_permits = Semaphore::MAX_PERMITS;
331                }
332                match permits.cmp(&new_permits) {
333                    Ordering::Less => {
334                        semaphore_clone.add_permits(new_permits - permits);
335                    }
336                    Ordering::Equal => continue,
337                    Ordering::Greater => {
338                        let to_release = permits - new_permits;
339                        let reduced = semaphore_clone.forget_permits(to_release);
340                        // TODO: implement dynamic semaphore with limits by ourself.
341                        if reduced != to_release {
342                            tracing::warn!(
343                                "no enough permits to release, expected {}, but reduced {}",
344                                to_release,
345                                reduced
346                            );
347                        }
348                    }
349                }
350                tracing::info!(
351                    "max_concurrent_creating_streaming_jobs changed from {} to {}",
352                    permits,
353                    new_permits
354                );
355                permits = new_permits;
356            }
357        });
358
359        Self { semaphore }
360    }
361}
362
363impl DdlController {
364    fn validate_specified_parallelism(
365        specified_parallelism: Option<NonZeroUsize>,
366        specified_backfill_parallelism: Option<NonZeroUsize>,
367        max_parallelism: NonZeroUsize,
368    ) -> MetaResult<()> {
369        if let Some(parallelism) = specified_parallelism
370            && parallelism > max_parallelism
371        {
372            bail_invalid_parameter!(
373                "specified parallelism {} should not exceed max parallelism {}",
374                parallelism,
375                max_parallelism,
376            );
377        }
378        if let Some(backfill_parallelism) = specified_backfill_parallelism
379            && backfill_parallelism > max_parallelism
380        {
381            bail_invalid_parameter!(
382                "specified backfill parallelism {} should not exceed max parallelism {}",
383                backfill_parallelism,
384                max_parallelism,
385            );
386        }
387        Ok(())
388    }
389
390    fn validate_serverless_backfill_enabled(
391        &self,
392        resource_type: &streaming_job_resource_type::ResourceType,
393    ) -> MetaResult<()> {
394        if matches!(
395            resource_type,
396            streaming_job_resource_type::ResourceType::ServerlessBackfill(true)
397        ) && self.env.opts.serverless_backfill_controller_addr.is_empty()
398        {
399            bail_invalid_parameter!(
400                "Serverless Backfill is disabled. Use RisingWave cloud at https://cloud.risingwave.com/auth/signup to try this feature"
401            );
402        }
403
404        Ok(())
405    }
406
407    pub async fn new(
408        env: MetaSrvEnv,
409        metadata_manager: MetadataManager,
410        stream_manager: GlobalStreamManagerRef,
411        source_manager: SourceManagerRef,
412        barrier_manager: BarrierManagerRef,
413        sink_manager: SinkCoordinatorManager,
414        iceberg_compaction_manager: IcebergCompactionManagerRef,
415        iceberg_pk_index_sink_manager: IcebergPkIndexSinkManager,
416    ) -> Self {
417        let creating_streaming_job_permits = Arc::new(CreatingStreamingJobPermit::new(&env).await);
418        Self {
419            env,
420            metadata_manager,
421            stream_manager,
422            source_manager,
423            barrier_manager,
424            sink_manager,
425            iceberg_compaction_manager,
426            iceberg_pk_index_sink_manager,
427            creating_streaming_job_permits,
428            seq: Arc::new(AtomicU64::new(0)),
429        }
430    }
431
432    /// Obtains the next sequence number for DDL commands, for observability and debugging purposes.
433    pub fn next_seq(&self) -> u64 {
434        // This is a simple atomic increment operation.
435        self.seq.fetch_add(1, std::sync::atomic::Ordering::Relaxed)
436    }
437
438    /// `run_command` spawns a tokio coroutine to execute the target ddl command. When the client
439    /// has been interrupted during executing, the request will be cancelled by tonic. Since we have
440    /// a lot of logic for revert, status management, notification and so on, ensuring consistency
441    /// would be a huge hassle and pain if we don't spawn here.
442    ///
443    /// Though returning `Option`, it's always `Some`, to simplify the handling logic
444    pub async fn run_command(&self, command: DdlCommand) -> MetaResult<Option<WaitVersion>> {
445        if !command.allow_in_recovery() {
446            self.barrier_manager.check_status_running()?;
447        }
448
449        let await_tree_key = format!("DDL Command {}", self.next_seq());
450        let await_tree_span = await_tree::span!("{command}({})", command.object());
451
452        let ctrl = self.clone();
453        let fut = Box::pin(async move {
454            match command {
455                DdlCommand::CreateDatabase(database) => ctrl.create_database(database).await,
456                DdlCommand::DropDatabase(database_id) => ctrl.drop_database(database_id).await,
457                DdlCommand::CreateSchema(schema) => ctrl.create_schema(schema).await,
458                DdlCommand::DropSchema(schema_id, drop_mode) => {
459                    ctrl.drop_schema(schema_id, drop_mode).await
460                }
461                DdlCommand::CreateNonSharedSource(source, iceberg_table_id) => {
462                    ctrl.create_non_shared_source(source, iceberg_table_id)
463                        .await
464                }
465                DdlCommand::DropSource(source_id, drop_mode) => {
466                    ctrl.drop_source(source_id, drop_mode).await
467                }
468                DdlCommand::ResetSource(source_id) => ctrl.reset_source(source_id).await,
469                DdlCommand::CreateFunction(function) => ctrl.create_function(function).await,
470                DdlCommand::DropFunction(function_id, drop_mode) => {
471                    ctrl.drop_function(function_id, drop_mode).await
472                }
473                DdlCommand::CreateView(view, dependencies) => {
474                    ctrl.create_view(view, dependencies).await
475                }
476                DdlCommand::DropView(view_id, drop_mode) => {
477                    ctrl.drop_view(view_id, drop_mode).await
478                }
479                DdlCommand::CreateStreamingJob {
480                    stream_job,
481                    fragment_graph,
482                    dependencies,
483                    resource_type,
484                    if_not_exists,
485                    refresh_interval_sec,
486                    replace_sink,
487                    since_timestamp_epoch,
488                } => {
489                    ctrl.create_streaming_job(
490                        stream_job,
491                        fragment_graph,
492                        dependencies,
493                        resource_type,
494                        if_not_exists,
495                        refresh_interval_sec,
496                        replace_sink,
497                        since_timestamp_epoch,
498                    )
499                    .await
500                }
501                DdlCommand::DropStreamingJob { job_id, drop_mode } => {
502                    ctrl.drop_streaming_job(job_id, drop_mode).await
503                }
504                DdlCommand::ReplaceStreamJob(ReplaceStreamJobInfo {
505                    streaming_job,
506                    fragment_graph,
507                }) => ctrl.replace_job(streaming_job, fragment_graph).await,
508                DdlCommand::AlterName(relation, name) => ctrl.alter_name(relation, &name).await,
509                DdlCommand::AlterObjectOwner(object, owner_id) => {
510                    ctrl.alter_owner(object, owner_id).await
511                }
512                DdlCommand::AlterSetSchema(object, new_schema_id) => {
513                    ctrl.alter_set_schema(object, new_schema_id).await
514                }
515                DdlCommand::CreateConnection(connection) => {
516                    ctrl.create_connection(connection).await
517                }
518                DdlCommand::DropConnection(connection_id, drop_mode) => {
519                    ctrl.drop_connection(connection_id, drop_mode).await
520                }
521                DdlCommand::CreateSecret(secret) => ctrl.create_secret(secret).await,
522                DdlCommand::DropSecret(secret_id, drop_mode) => {
523                    ctrl.drop_secret(secret_id, drop_mode).await
524                }
525                DdlCommand::AlterSecret(secret) => ctrl.alter_secret(secret).await,
526                DdlCommand::AlterNonSharedSource(source) => {
527                    ctrl.alter_non_shared_source(source).await
528                }
529                DdlCommand::CommentOn(comment) => ctrl.comment_on(comment).await,
530                DdlCommand::CreateSubscription(subscription) => {
531                    ctrl.create_subscription(subscription).await
532                }
533                DdlCommand::DropSubscription(subscription_id, drop_mode) => {
534                    ctrl.drop_subscription(subscription_id, drop_mode).await
535                }
536                DdlCommand::AlterSubscriptionRetention {
537                    subscription_id,
538                    retention_seconds,
539                    definition,
540                } => {
541                    ctrl.alter_subscription_retention(
542                        subscription_id,
543                        retention_seconds,
544                        definition,
545                    )
546                    .await
547                }
548                DdlCommand::AlterSwapRename(objects) => ctrl.alter_swap_rename(objects).await,
549                DdlCommand::AlterDatabaseParam(database_id, param) => {
550                    ctrl.alter_database_param(database_id, param).await
551                }
552                DdlCommand::AlterDatabaseResourceGroup(database_id, resource_group, deferred) => {
553                    ctrl.alter_database_resource_group(database_id, resource_group, deferred)
554                        .await
555                }
556                DdlCommand::AlterStreamingJobConfig(job_id, entries_to_add, keys_to_remove) => {
557                    ctrl.alter_streaming_job_config(job_id, entries_to_add, keys_to_remove)
558                        .await
559                }
560            }
561        })
562        .in_current_span();
563        let fut = (self.env.await_tree_reg())
564            .register(await_tree_key, await_tree_span)
565            .instrument(Box::pin(fut));
566        let notification_version = tokio::spawn(fut).await.map_err(|e| anyhow!(e))??;
567        Ok(Some(WaitVersion {
568            catalog_version: notification_version,
569            hummock_version_id: self.barrier_manager.get_hummock_version_id().await,
570        }))
571    }
572
573    pub async fn get_ddl_progress(&self) -> MetaResult<Vec<DdlProgress>> {
574        self.barrier_manager.get_ddl_progress().await
575    }
576
577    async fn create_database(&self, database: Database) -> MetaResult<NotificationVersion> {
578        let (version, updated_db) = self
579            .metadata_manager
580            .catalog_controller
581            .create_database(database)
582            .await?;
583        // If persistent successfully, notify `GlobalBarrierManager` to create database asynchronously.
584        self.barrier_manager
585            .update_database_barrier(
586                updated_db.database_id,
587                updated_db.barrier_interval_ms.map(|v| v as u32),
588                updated_db.checkpoint_frequency.map(|v| v as u64),
589            )
590            .await?;
591        Ok(version)
592    }
593
594    #[tracing::instrument(skip(self), level = "debug")]
595    pub async fn reschedule_streaming_job(
596        &self,
597        job_id: JobId,
598        target: ReschedulePolicy,
599        mut deferred: bool,
600    ) -> MetaResult<()> {
601        tracing::info!("altering parallelism for job {}", job_id);
602        if self.barrier_manager.check_status_running().is_err() {
603            tracing::info!(
604                "alter parallelism is set to deferred mode because the system is in recovery state"
605            );
606            deferred = true;
607        }
608
609        self.stream_manager
610            .reschedule_streaming_job(job_id, target, deferred)
611            .await
612    }
613
614    pub async fn reschedule_streaming_job_backfill_parallelism(
615        &self,
616        job_id: JobId,
617        parallelism: Option<ParallelismPolicy>,
618        mut deferred: bool,
619    ) -> MetaResult<()> {
620        tracing::info!("altering backfill parallelism for job {}", job_id);
621        if self.barrier_manager.check_status_running().is_err() {
622            tracing::info!(
623                "alter backfill parallelism is set to deferred mode because the system is in recovery state"
624            );
625            deferred = true;
626        }
627
628        self.stream_manager
629            .reschedule_streaming_job_backfill_parallelism(job_id, parallelism, deferred)
630            .await
631    }
632
633    pub async fn reschedule_cdc_table_backfill(
634        &self,
635        job_id: JobId,
636        target: ReschedulePolicy,
637    ) -> MetaResult<()> {
638        tracing::info!("alter CDC table backfill parallelism");
639        if self.barrier_manager.check_status_running().is_err() {
640            return Err(anyhow::anyhow!("CDC table backfill reschedule is unavailable because the system is in recovery state").into());
641        }
642        self.stream_manager
643            .reschedule_cdc_table_backfill(job_id, target)
644            .await
645    }
646
647    pub async fn reschedule_fragments(
648        &self,
649        fragment_targets: HashMap<FragmentId, Option<StreamingParallelism>>,
650    ) -> MetaResult<()> {
651        tracing::info!(
652            "altering parallelism for fragments {:?}",
653            fragment_targets.keys()
654        );
655        let fragment_targets = fragment_targets
656            .into_iter()
657            .map(|(fragment_id, parallelism)| (fragment_id as CatalogFragmentId, parallelism))
658            .collect();
659
660        self.stream_manager
661            .reschedule_fragments(fragment_targets)
662            .await
663    }
664
665    async fn drop_database(&self, database_id: DatabaseId) -> MetaResult<NotificationVersion> {
666        self.drop_object(ObjectType::Database, database_id, DropMode::Cascade)
667            .await
668    }
669
670    async fn create_schema(&self, schema: Schema) -> MetaResult<NotificationVersion> {
671        self.metadata_manager
672            .catalog_controller
673            .create_schema(schema)
674            .await
675    }
676
677    async fn drop_schema(
678        &self,
679        schema_id: SchemaId,
680        drop_mode: DropMode,
681    ) -> MetaResult<NotificationVersion> {
682        self.drop_object(ObjectType::Schema, schema_id, drop_mode)
683            .await
684    }
685
686    /// Shared source is handled in [`Self::create_streaming_job`]
687    async fn create_non_shared_source(
688        &self,
689        source: Source,
690        iceberg_table_id: Option<TableId>,
691    ) -> MetaResult<NotificationVersion> {
692        if let Some(cdc_table_desc) = source
693            .info
694            .as_ref()
695            .and_then(|info| info.external_table.as_ref())
696        {
697            assert!(iceberg_table_id.is_none());
698            // A CDC table source has no streaming job, so validate the declared schema against
699            // the upstream table here, the same way `CREATE TABLE ... FROM <cdc source>` does
700            // while creating its job.
701            self.validate_cdc_table_desc(cdc_table_desc).await?;
702            let (_, version) = self
703                .metadata_manager
704                .catalog_controller
705                .create_source(source, iceberg_table_id)
706                .await?;
707            return Ok(version);
708        }
709
710        let handle = create_source_worker(
711            &source,
712            self.source_manager.metrics.clone(),
713            self.env.await_tree_reg().clone(),
714        )
715        .await
716        .context("failed to create source worker")?;
717
718        let (source_id, version) = self
719            .metadata_manager
720            .catalog_controller
721            .create_source(source, iceberg_table_id)
722            .await?;
723        self.source_manager
724            .register_source_with_handle(source_id, handle)
725            .await;
726        Ok(version)
727    }
728
729    async fn drop_source(
730        &self,
731        source_id: SourceId,
732        drop_mode: DropMode,
733    ) -> MetaResult<NotificationVersion> {
734        self.drop_object(ObjectType::Source, source_id, drop_mode)
735            .await
736    }
737
738    async fn reset_source(&self, source_id: SourceId) -> MetaResult<NotificationVersion> {
739        tracing::info!(source_id = %source_id, "resetting CDC source offset to latest");
740
741        // Get database_id for the source
742        let database_id = self
743            .metadata_manager
744            .catalog_controller
745            .get_object_database_id(source_id)
746            .await?;
747
748        self.stream_manager
749            .barrier_scheduler
750            .run_command(database_id, Command::ResetSource { source_id })
751            .await?;
752
753        // RESET SOURCE doesn't modify catalog, so return the current catalog version
754        let version = self
755            .metadata_manager
756            .catalog_controller
757            .notify_frontend_trivial()
758            .await;
759        Ok(version)
760    }
761
762    /// This replaces the source in the catalog.
763    /// Note: `StreamSourceInfo` in downstream MVs' `SourceExecutor`s are not updated.
764    async fn alter_non_shared_source(&self, source: Source) -> MetaResult<NotificationVersion> {
765        self.metadata_manager
766            .catalog_controller
767            .alter_non_shared_source(source)
768            .await
769    }
770
771    async fn create_function(&self, function: Function) -> MetaResult<NotificationVersion> {
772        self.metadata_manager
773            .catalog_controller
774            .create_function(function)
775            .await
776    }
777
778    async fn drop_function(
779        &self,
780        function_id: FunctionId,
781        drop_mode: DropMode,
782    ) -> MetaResult<NotificationVersion> {
783        self.drop_object(ObjectType::Function, function_id, drop_mode)
784            .await
785    }
786
787    async fn create_view(
788        &self,
789        view: View,
790        dependencies: HashSet<ObjectId>,
791    ) -> MetaResult<NotificationVersion> {
792        self.metadata_manager
793            .catalog_controller
794            .create_view(view, dependencies)
795            .await
796    }
797
798    async fn drop_view(
799        &self,
800        view_id: ViewId,
801        drop_mode: DropMode,
802    ) -> MetaResult<NotificationVersion> {
803        self.drop_object(ObjectType::View, view_id, drop_mode).await
804    }
805
806    async fn create_connection(&self, connection: Connection) -> MetaResult<NotificationVersion> {
807        validate_connection(&connection).await?;
808        self.metadata_manager
809            .catalog_controller
810            .create_connection(connection)
811            .await
812    }
813
814    async fn drop_connection(
815        &self,
816        connection_id: ConnectionId,
817        drop_mode: DropMode,
818    ) -> MetaResult<NotificationVersion> {
819        self.drop_object(ObjectType::Connection, connection_id, drop_mode)
820            .await
821    }
822
823    async fn alter_database_param(
824        &self,
825        database_id: DatabaseId,
826        param: AlterDatabaseParam,
827    ) -> MetaResult<NotificationVersion> {
828        let (version, updated_db) = self
829            .metadata_manager
830            .catalog_controller
831            .alter_database_param(database_id, param)
832            .await?;
833        // If persistent successfully, notify `GlobalBarrierManager` to update param asynchronously.
834        self.barrier_manager
835            .update_database_barrier(
836                database_id,
837                updated_db.barrier_interval_ms.map(|v| v as u32),
838                updated_db.checkpoint_frequency.map(|v| v as u64),
839            )
840            .await?;
841        Ok(version)
842    }
843
844    async fn alter_database_resource_group(
845        &self,
846        database_id: DatabaseId,
847        resource_group: Option<String>,
848        _deferred: bool,
849    ) -> MetaResult<NotificationVersion> {
850        let version = self
851            .metadata_manager
852            .catalog_controller
853            .alter_database_resource_group(database_id, resource_group)
854            .await?;
855
856        Ok(version)
857    }
858
859    // The 'secret' part of the request we receive from the frontend is in plaintext;
860    // here, we need to encrypt it before storing it in the catalog.
861    fn get_encrypted_payload(&self, secret: &Secret) -> MetaResult<Vec<u8>> {
862        let secret_store_private_key = self
863            .env
864            .opts
865            .secret_store_private_key
866            .clone()
867            .ok_or_else(|| anyhow!("secret_store_private_key is not configured"))?;
868
869        let encrypted_payload = SecretEncryption::encrypt(
870            secret_store_private_key.as_slice(),
871            secret.get_value().as_slice(),
872        )
873        .context(format!("failed to encrypt secret {}", secret.name))?;
874        Ok(encrypted_payload
875            .serialize()
876            .context(format!("failed to serialize secret {}", secret.name))?)
877    }
878
879    async fn create_secret(&self, mut secret: Secret) -> MetaResult<NotificationVersion> {
880        // The 'secret' part of the request we receive from the frontend is in plaintext;
881        // here, we need to encrypt it before storing it in the catalog.
882        let secret_plain_payload = secret.value.clone();
883        let encrypted_payload = self.get_encrypted_payload(&secret)?;
884        secret.value = encrypted_payload;
885
886        self.metadata_manager
887            .catalog_controller
888            .create_secret(secret, secret_plain_payload)
889            .await
890    }
891
892    async fn drop_secret(
893        &self,
894        secret_id: SecretId,
895        drop_mode: DropMode,
896    ) -> MetaResult<NotificationVersion> {
897        self.drop_object(ObjectType::Secret, secret_id, drop_mode)
898            .await
899    }
900
901    async fn alter_secret(&self, mut secret: Secret) -> MetaResult<NotificationVersion> {
902        let secret_plain_payload = secret.value.clone();
903        let encrypted_payload = self.get_encrypted_payload(&secret)?;
904        secret.value = encrypted_payload;
905        self.metadata_manager
906            .catalog_controller
907            .alter_secret(secret, secret_plain_payload)
908            .await
909    }
910
911    async fn create_subscription(
912        &self,
913        mut subscription: Subscription,
914    ) -> MetaResult<NotificationVersion> {
915        tracing::debug!("create subscription");
916        let _permit = self
917            .creating_streaming_job_permits
918            .semaphore
919            .acquire()
920            .await
921            .unwrap();
922        let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
923        self.metadata_manager
924            .catalog_controller
925            .create_subscription_catalog(&mut subscription)
926            .await?;
927        if let Err(err) = self.stream_manager.create_subscription(&subscription).await {
928            tracing::debug!(error = %err.as_report(), "failed to create subscription");
929            let _ = self
930                .metadata_manager
931                .catalog_controller
932                .try_abort_creating_subscription(subscription.id)
933                .await
934                .inspect_err(|e| {
935                    tracing::error!(
936                        error = %e.as_report(),
937                        "failed to abort create subscription after failure"
938                    );
939                });
940            return Err(err);
941        }
942
943        let version = self
944            .metadata_manager
945            .catalog_controller
946            .notify_create_subscription(subscription.id)
947            .await?;
948        tracing::debug!("finish create subscription");
949        Ok(version)
950    }
951
952    async fn drop_subscription(
953        &self,
954        subscription_id: SubscriptionId,
955        drop_mode: DropMode,
956    ) -> MetaResult<NotificationVersion> {
957        tracing::debug!("preparing drop subscription");
958        let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
959        let subscription = self
960            .metadata_manager
961            .catalog_controller
962            .get_subscription_by_id(subscription_id)
963            .await?;
964        let table_id = subscription.dependent_table_id;
965        let database_id = subscription.database_id;
966        let (_, version) = self
967            .metadata_manager
968            .catalog_controller
969            .drop_object(ObjectType::Subscription, subscription_id, drop_mode)
970            .await?;
971        self.stream_manager
972            .drop_subscription(database_id, subscription_id, table_id)
973            .await;
974        tracing::debug!("finish drop subscription");
975        Ok(version)
976    }
977
978    async fn alter_subscription_retention(
979        &self,
980        subscription_id: SubscriptionId,
981        retention_seconds: u64,
982        definition: String,
983    ) -> MetaResult<NotificationVersion> {
984        tracing::debug!("alter subscription retention");
985        let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
986        let (version, subscription) = self
987            .metadata_manager
988            .catalog_controller
989            .alter_subscription_retention(subscription_id, retention_seconds, definition)
990            .await?;
991        self.stream_manager
992            .alter_subscription_retention(
993                subscription.database_id,
994                subscription.id,
995                subscription.dependent_table_id,
996                subscription.retention_seconds,
997            )
998            .await?;
999        tracing::debug!("finish alter subscription retention");
1000        Ok(version)
1001    }
1002
1003    /// Validates the connect properties in the `cdc_table_desc` stored in the `StreamCdcScan` node
1004    #[await_tree::instrument]
1005    pub(crate) async fn validate_cdc_table(
1006        &self,
1007        table: &Table,
1008        table_fragments: &StreamJobFragments,
1009    ) -> MetaResult<()> {
1010        let stream_scan_fragment =
1011            Itertools::exactly_one(table_fragments.fragments.values().filter(|f| {
1012                f.fragment_type_mask.contains(FragmentTypeFlag::StreamScan)
1013                    || f.fragment_type_mask
1014                        .contains(FragmentTypeFlag::StreamCdcScan)
1015            }))
1016            .ok()
1017            .with_context(|| {
1018                format!(
1019                    "expect exactly one stream scan fragment, got: {:?}",
1020                    table_fragments.fragments
1021                )
1022            })?;
1023        fn assert_parallelism(
1024            distribution_type: PbFragmentDistributionType,
1025            node_body: &Option<NodeBody>,
1026        ) {
1027            if let Some(NodeBody::StreamCdcScan(node)) = node_body {
1028                if let Some(o) = node.options
1029                    && CdcScanOptions::from_proto(&o).is_parallelized_backfill()
1030                {
1031                    // Use parallel CDC backfill.
1032                } else {
1033                    assert_eq!(
1034                        distribution_type,
1035                        PbFragmentDistributionType::Single,
1036                        "Non-parallelized CDC scan fragment should have Single distribution"
1037                    );
1038                }
1039            }
1040        }
1041        let mut found_cdc_scan = false;
1042        match &stream_scan_fragment.nodes.node_body {
1043            Some(NodeBody::StreamCdcScan(_)) => {
1044                assert_parallelism(
1045                    stream_scan_fragment.distribution_type,
1046                    &stream_scan_fragment.nodes.node_body,
1047                );
1048                if self
1049                    .validate_cdc_table_inner(&stream_scan_fragment.nodes.node_body, table.id)
1050                    .await?
1051                {
1052                    found_cdc_scan = true;
1053                }
1054            }
1055            // When there's generated columns, the cdc scan node is wrapped in a project node
1056            Some(NodeBody::Project(_)) => {
1057                for input in &stream_scan_fragment.nodes.input {
1058                    assert_parallelism(stream_scan_fragment.distribution_type, &input.node_body);
1059                    if self
1060                        .validate_cdc_table_inner(&input.node_body, table.id)
1061                        .await?
1062                    {
1063                        found_cdc_scan = true;
1064                    }
1065                }
1066            }
1067            _ => {
1068                bail!("Unexpected node body for stream cdc scan");
1069            }
1070        };
1071        if !found_cdc_scan {
1072            bail!("No stream cdc scan node found in stream scan fragment");
1073        }
1074        Ok(())
1075    }
1076
1077    async fn validate_cdc_table_inner(
1078        &self,
1079        node_body: &Option<NodeBody>,
1080        table_id: TableId,
1081    ) -> MetaResult<bool> {
1082        if let Some(NodeBody::StreamCdcScan(stream_cdc_scan)) = node_body
1083            && let Some(ref cdc_table_desc) = stream_cdc_scan.cdc_table_desc
1084        {
1085            self.validate_cdc_table_desc(cdc_table_desc).await?;
1086            tracing::debug!(?table_id, "validate cdc table success");
1087            Ok(true)
1088        } else {
1089            Ok(false)
1090        }
1091    }
1092
1093    /// Validates a CDC table descriptor against its upstream table, by creating a throw-away
1094    /// split enumerator: the connector validator checks that the table exists and that the
1095    /// declared columns and primary key match the upstream ones.
1096    pub(crate) async fn validate_cdc_table_desc(
1097        &self,
1098        cdc_table_desc: &PbExternalTableDesc,
1099    ) -> MetaResult<()> {
1100        let options_with_secret = WithOptionsSecResolved::new(
1101            cdc_table_desc.connect_properties.clone(),
1102            cdc_table_desc.secret_refs.clone(),
1103        );
1104
1105        let mut props = ConnectorProperties::extract(options_with_secret, true)?;
1106        props.init_from_pb_cdc_table_desc(cdc_table_desc);
1107
1108        let _enumerator = props
1109            .create_split_enumerator(SourceEnumeratorContext::dummy().into())
1110            .await?;
1111
1112        Ok(())
1113    }
1114
1115    pub async fn validate_table_for_sink(&self, table_id: TableId) -> MetaResult<()> {
1116        let migrated = self
1117            .metadata_manager
1118            .catalog_controller
1119            .has_table_been_migrated(table_id)
1120            .await?;
1121        if !migrated {
1122            Err(anyhow::anyhow!("Creating sink into table is not allowed for unmigrated table {}. Please migrate it first.", table_id).into())
1123        } else {
1124            Ok(())
1125        }
1126    }
1127
1128    /// For [`CreateType::Foreground`], the function will only return after backfilling finishes
1129    /// ([`crate::manager::MetadataManager::wait_streaming_job_finished`]).
1130    #[await_tree::instrument(boxed, "create_streaming_job({streaming_job})")]
1131    pub async fn create_streaming_job(
1132        &self,
1133        mut streaming_job: StreamingJob,
1134        fragment_graph: StreamFragmentGraphProto,
1135        dependencies: HashSet<ObjectId>,
1136        resource_type: streaming_job_resource_type::ResourceType,
1137        if_not_exists: bool,
1138        refresh_interval_sec: Option<u64>,
1139        replace_sink: Option<SinkId>,
1140        since_timestamp_epoch: Option<u64>,
1141    ) -> MetaResult<NotificationVersion> {
1142        let replace_sink_info = if let Some(old_sink_id) = replace_sink {
1143            let StreamingJob::Sink(sink, _) = &streaming_job else {
1144                bail!("replace sink requires a sink job")
1145            };
1146            if sink.target_table.is_some() {
1147                bail_not_implemented!("replace sink into table")
1148            }
1149
1150            Some(old_sink_id)
1151        } else {
1152            if let StreamingJob::Sink(sink, _) = &streaming_job
1153                && let Some(target_table) = sink.target_table
1154            {
1155                self.validate_table_for_sink(target_table).await?;
1156            }
1157            None
1158        };
1159        self.validate_serverless_backfill_enabled(&resource_type)?;
1160        let ctx = StreamContext::from_protobuf(fragment_graph.get_ctx().unwrap());
1161        let adaptive_parallelism_strategy =
1162            (!fragment_graph.adaptive_parallelism_strategy.is_empty()).then(|| {
1163                parse_strategy(&fragment_graph.adaptive_parallelism_strategy)
1164                    .expect("adaptive parallelism strategy should be validated in frontend")
1165            });
1166        let backfill_adaptive_parallelism_strategy = (!fragment_graph
1167            .backfill_adaptive_parallelism_strategy
1168            .is_empty())
1169        .then(|| {
1170            parse_strategy(&fragment_graph.backfill_adaptive_parallelism_strategy)
1171                .expect("backfill adaptive parallelism strategy should be validated in frontend")
1172        });
1173
1174        let streaming_job_model = match self
1175            .metadata_manager
1176            .catalog_controller
1177            .create_job_catalog(
1178                &mut streaming_job,
1179                &ctx,
1180                &fragment_graph.parallelism,
1181                fragment_graph.max_parallelism as _,
1182                dependencies,
1183                resource_type.clone(),
1184                &fragment_graph.backfill_parallelism,
1185                adaptive_parallelism_strategy,
1186                backfill_adaptive_parallelism_strategy,
1187                replace_sink_info.as_ref(),
1188                refresh_interval_sec,
1189            )
1190            .await
1191        {
1192            Ok(model) => model,
1193            Err(meta_err) => {
1194                if !if_not_exists {
1195                    return Err(meta_err);
1196                }
1197                return if let MetaErrorInner::Duplicated(_, _, Some(job_id)) = meta_err.inner() {
1198                    if streaming_job.create_type() == CreateType::Foreground {
1199                        let database_id = streaming_job.database_id();
1200                        self.metadata_manager
1201                            .wait_streaming_job_finished(database_id, *job_id)
1202                            .await
1203                    } else {
1204                        Ok(IGNORED_NOTIFICATION_VERSION)
1205                    }
1206                } else {
1207                    Err(meta_err)
1208                };
1209            }
1210        };
1211        let job_id = streaming_job.id();
1212        if let Some(old_sink_id) = replace_sink_info.as_ref() {
1213            tracing::debug!(
1214                old_sink_id = %old_sink_id,
1215                new_sink_id = %job_id,
1216                definition = streaming_job.definition(),
1217                create_type = streaming_job.create_type().as_str_name(),
1218                "starting replacement sink",
1219            );
1220        } else {
1221            tracing::debug!(
1222                id = %job_id,
1223                definition = streaming_job.definition(),
1224                create_type = streaming_job.create_type().as_str_name(),
1225                job_type = ?streaming_job.job_type(),
1226                "starting streaming job",
1227            );
1228        }
1229        // TODO: acquire permits for recovered background DDLs.
1230        let permit = self
1231            .creating_streaming_job_permits
1232            .semaphore
1233            .clone()
1234            .acquire_owned()
1235            .instrument_await("acquire_creating_streaming_job_permit")
1236            .await
1237            .unwrap();
1238        let reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
1239
1240        let name = streaming_job.name();
1241        let definition = streaming_job.definition();
1242        let database_id = streaming_job.database_id();
1243        let source_id = match &streaming_job {
1244            StreamingJob::Table(Some(src), _, _) | StreamingJob::Source(src) => Some(src.id),
1245            _ => None,
1246        };
1247        let create_result = match self
1248            .generate_streaming_job(
1249                ctx,
1250                streaming_job,
1251                fragment_graph,
1252                resource_type.clone(),
1253                streaming_job_model,
1254                replace_sink_info,
1255                since_timestamp_epoch,
1256            )
1257            .await
1258        {
1259            Ok((stream_job_fragments, ctx)) => {
1260                self.stream_manager
1261                    .create_streaming_job(stream_job_fragments, ctx, permit, reschedule_job_lock)
1262                    .await
1263            }
1264            Err(err) => Err((err, false, None)),
1265        };
1266
1267        match create_result {
1268            Ok(version) => Ok(version),
1269            Err((err, is_cancelled, cancel_notifier)) => {
1270                tracing::error!(id = %job_id, error = %err.as_report(), "failed to create streaming job");
1271                let event = risingwave_pb::meta::event_log::EventCreateStreamJobFail {
1272                    id: job_id,
1273                    name,
1274                    definition,
1275                    error: err.as_report().to_string(),
1276                };
1277                self.env.event_log_manager_ref().add_event_logs(vec![
1278                    risingwave_pb::meta::event_log::Event::CreateStreamJobFail(event),
1279                ]);
1280                let abort_result = self
1281                    .metadata_manager
1282                    .catalog_controller
1283                    .try_abort_creating_streaming_job(job_id, is_cancelled)
1284                    .await?;
1285                self.iceberg_compaction_manager
1286                    .clear_maintenance_for_aborted_job(&abort_result);
1287                if let Some(cancel_info) = abort_result.cancel_info {
1288                    self.stream_manager
1289                        .barrier_scheduler
1290                        .run_command(database_id, cancel_info.command)
1291                        .await?;
1292                    cleanup_dropped_streaming_jobs(
1293                        &self.stream_manager.refresh_manager,
1294                        &self.stream_manager.hummock_manager,
1295                        &self.stream_manager.metadata_manager,
1296                        cancel_info.streaming_job_ids,
1297                        cancel_info.state_table_ids,
1298                        "cancel_streaming_job",
1299                    )
1300                    .await?;
1301                }
1302                if let Some(cancel_notifier) = cancel_notifier {
1303                    let _ = cancel_notifier.send(true).inspect_err(|err| {
1304                        tracing::warn!("failed to notify cancellation result: {err}")
1305                    });
1306                }
1307                if abort_result.aborted {
1308                    tracing::warn!(id = %job_id, is_cancelled, "aborted streaming job");
1309                    // FIXME: might also need other cleanup here
1310                    if let Some(source_id) = source_id {
1311                        self.source_manager
1312                            .apply_source_change(SourceChange::DropSource {
1313                                dropped_source_ids: vec![source_id],
1314                            })
1315                            .await;
1316                    }
1317                }
1318                Err(err)
1319            }
1320        }
1321    }
1322
1323    #[await_tree::instrument(boxed)]
1324    async fn generate_streaming_job(
1325        &self,
1326        ctx: StreamContext,
1327        mut streaming_job: StreamingJob,
1328        fragment_graph: StreamFragmentGraphProto,
1329        resource_type: streaming_job_resource_type::ResourceType,
1330        streaming_job_model: streaming_job::Model,
1331        replace_sink: Option<SinkId>,
1332        since_timestamp_epoch: Option<u64>,
1333    ) -> MetaResult<(StreamJobFragmentsToCreate, CreateStreamingJobContext)> {
1334        let mut fragment_graph =
1335            StreamFragmentGraph::new(&self.env, fragment_graph, &streaming_job)?;
1336        streaming_job.set_info_from_graph(&fragment_graph);
1337
1338        // create internal table catalogs and refill table id.
1339        let incomplete_internal_tables = fragment_graph
1340            .incomplete_internal_tables()
1341            .into_values()
1342            .collect_vec();
1343        let table_id_map = self
1344            .metadata_manager
1345            .catalog_controller
1346            .create_internal_table_catalog(&streaming_job, incomplete_internal_tables)
1347            .await?;
1348        fragment_graph.refill_internal_table_ids(table_id_map);
1349
1350        // create fragment and actor catalogs.
1351        tracing::debug!(id = %streaming_job.id(), "building streaming job");
1352        let (mut ctx, stream_job_fragments) = self
1353            .build_stream_job(
1354                ctx,
1355                streaming_job,
1356                fragment_graph,
1357                resource_type,
1358                streaming_job_model,
1359                since_timestamp_epoch,
1360            )
1361            .await?;
1362        ctx.replace_sink = replace_sink;
1363
1364        let streaming_job = &ctx.streaming_job;
1365
1366        match streaming_job {
1367            StreamingJob::Table(None, table, TableJobType::SharedCdcSource) => {
1368                self.validate_cdc_table(table, &stream_job_fragments)
1369                    .await?;
1370            }
1371            StreamingJob::Table(Some(source), ..) => {
1372                // Register the source on the connector node.
1373                self.source_manager.register_source(source).await?;
1374                let connector_name = source
1375                    .get_with_properties()
1376                    .get(UPSTREAM_SOURCE_KEY)
1377                    .cloned();
1378                let attr = source.info.as_ref().map(|source_info| {
1379                    jsonbb::json!({
1380                            "format": source_info.format().as_str_name(),
1381                            "encode": source_info.row_encode().as_str_name(),
1382                    })
1383                });
1384                report_create_object(
1385                    streaming_job.id(),
1386                    "source",
1387                    PbTelemetryDatabaseObject::Source,
1388                    connector_name,
1389                    attr,
1390                );
1391            }
1392            StreamingJob::Sink(sink, _) => {
1393                if sink.auto_refresh_schema_from_table.is_some() {
1394                    check_sink_fragments_support_refresh_schema(&stream_job_fragments.fragments)?;
1395                }
1396                // Validate the sink on the connector node.
1397                validate_sink(sink).await?;
1398                // For Iceberg pk-index sinks, spawn the per-sink commit worker now
1399                // so it's ready to receive epoch reports from the very first
1400                // barrier instead of relying on lazy registration on every
1401                // commit.
1402                if crate::manager::iceberg_pk_index_sink::is_iceberg_pk_index_sink(&sink.properties)
1403                {
1404                    let iceberg_config =
1405                        crate::manager::iceberg_pk_index_sink::build_iceberg_config(sink)?;
1406                    self.iceberg_pk_index_sink_manager
1407                        .register_sink(
1408                            sink.id,
1409                            crate::barrier::to_partial_graph_id(sink.database_id, None),
1410                            iceberg_config,
1411                        )
1412                        .await
1413                        .map_err(|e| anyhow!(e).context("register v3 sink worker"))?;
1414                }
1415                let connector_name = sink.get_properties().get(UPSTREAM_SOURCE_KEY).cloned();
1416                let attr = sink.format_desc.as_ref().map(|sink_info| {
1417                    jsonbb::json!({
1418                        "format": sink_info.format().as_str_name(),
1419                        "encode": sink_info.encode().as_str_name(),
1420                    })
1421                });
1422                report_create_object(
1423                    streaming_job.id(),
1424                    "sink",
1425                    PbTelemetryDatabaseObject::Sink,
1426                    connector_name,
1427                    attr,
1428                );
1429            }
1430            StreamingJob::Source(source) => {
1431                // Register the source on the connector node.
1432                self.source_manager.register_source(source).await?;
1433                let connector_name = source
1434                    .get_with_properties()
1435                    .get(UPSTREAM_SOURCE_KEY)
1436                    .cloned();
1437                let attr = source.info.as_ref().map(|source_info| {
1438                    jsonbb::json!({
1439                            "format": source_info.format().as_str_name(),
1440                            "encode": source_info.row_encode().as_str_name(),
1441                    })
1442                });
1443                report_create_object(
1444                    streaming_job.id(),
1445                    "source",
1446                    PbTelemetryDatabaseObject::Source,
1447                    connector_name,
1448                    attr,
1449                );
1450            }
1451            _ => {}
1452        }
1453
1454        let backfill_orders = ctx.fragment_backfill_ordering.to_meta_model();
1455        self.metadata_manager
1456            .catalog_controller
1457            .prepare_stream_job_fragments(
1458                &stream_job_fragments,
1459                streaming_job,
1460                false,
1461                Some(backfill_orders),
1462            )
1463            .await?;
1464
1465        Ok((stream_job_fragments, ctx))
1466    }
1467
1468    /// `target_replace_info`: when dropping a sink into table, we need to replace the table.
1469    pub async fn drop_object(
1470        &self,
1471        object_type: ObjectType,
1472        object_id: impl Into<ObjectId>,
1473        drop_mode: DropMode,
1474    ) -> MetaResult<NotificationVersion> {
1475        let object_id = object_id.into();
1476        // Fence reschedule and source tick before catalog deletion so post-collect split updates
1477        // cannot race with dropped fragments.
1478        let _reschedule_job_lock = self.stream_manager.reschedule_lock_read_guard().await;
1479        let _source_tick_pause_guard = self.source_manager.pause_tick().await;
1480
1481        let (release_ctx, version) = self
1482            .metadata_manager
1483            .catalog_controller
1484            .drop_object(object_type, object_id, drop_mode)
1485            .await?;
1486
1487        if object_type == ObjectType::Source {
1488            self.env
1489                .notification_manager_ref()
1490                .notify_local_subscribers(LocalNotification::SourceDropped(object_id));
1491        }
1492
1493        let ReleaseContext {
1494            database_id,
1495            removed_streaming_job_ids,
1496            removed_state_table_ids,
1497            removed_source_ids,
1498            removed_secret_ids: secret_ids,
1499            removed_source_fragments,
1500            removed_fragments,
1501            removed_sink_fragment_by_targets,
1502            removed_iceberg_table_sinks,
1503            removed_iceberg_sink_ids,
1504            removed_iceberg_pk_index_sink_ids,
1505        } = release_ctx;
1506        let removed_job_ids_for_sink_coordinators = removed_streaming_job_ids.clone();
1507
1508        // Notify serving module about deleted fragments so it can clean up serving vnode mappings.
1509        // This is driven by the fragment model deletion (cascade from Object::delete_many),
1510        // decoupled from the barrier-driven streaming mapping notifications.
1511        self.env
1512            .notification_manager_ref()
1513            .notify_serving_fragment_mapping_delete(
1514                removed_fragments.iter().map(|id| *id as _).collect(),
1515            );
1516
1517        self.stream_manager
1518            .drop_streaming_jobs(
1519                database_id,
1520                removed_streaming_job_ids,
1521                removed_state_table_ids,
1522                removed_sink_fragment_by_targets
1523                    .into_iter()
1524                    .map(|(target, sinks)| {
1525                        (target as _, sinks.into_iter().map(|id| id as _).collect())
1526                    })
1527                    .collect(),
1528            )
1529            .await;
1530
1531        // clean up sources after dropping streaming jobs.
1532        // Otherwise, e.g., Kafka consumer groups might be recreated after deleted.
1533        self.source_manager
1534            .apply_source_change(SourceChange::DropSource {
1535                dropped_source_ids: removed_source_ids.into_iter().map(|id| id as _).collect(),
1536            })
1537            .await;
1538
1539        // unregister fragments and actors from source manager.
1540        // FIXME: need also unregister source backfill fragments.
1541        let dropped_source_fragments = removed_source_fragments;
1542        self.source_manager
1543            .apply_source_change(SourceChange::DropMv {
1544                dropped_source_fragments,
1545            })
1546            .await;
1547
1548        // clean up iceberg table sinks
1549        for sink in removed_iceberg_table_sinks {
1550            let sink_param = SinkParam::try_from_sink_catalog(sink.into())
1551                .expect("Iceberg sink should be valid");
1552            let iceberg_sink =
1553                IcebergSink::try_from(sink_param).expect("Iceberg sink should be valid");
1554            if let Ok(iceberg_catalog) = iceberg_sink.config.create_catalog().await {
1555                let table_identifier = iceberg_sink.config.full_table_name().unwrap();
1556                tracing::info!(
1557                    "dropping iceberg table {} for dropped sink",
1558                    table_identifier
1559                );
1560
1561                let _ = iceberg_catalog
1562                    .drop_table(&table_identifier)
1563                    .await
1564                    .inspect_err(|err| {
1565                        tracing::error!(
1566                            "failed to drop iceberg table {} during cleanup: {}",
1567                            table_identifier,
1568                            err.as_report()
1569                        );
1570                    });
1571            }
1572        }
1573
1574        // stop sink coordinators for dropped streaming jobs
1575        if !removed_job_ids_for_sink_coordinators.is_empty() {
1576            self.sink_manager
1577                .stop_sink_coordinators_for_jobs(removed_job_ids_for_sink_coordinators)
1578                .await;
1579        }
1580
1581        // Covers user-created iceberg sinks dropped via CASCADE, which are not in
1582        // `removed_iceberg_table_sinks` above.
1583        for sink_id in removed_iceberg_sink_ids {
1584            self.iceberg_compaction_manager
1585                .clear_iceberg_maintenance_by_sink_id(sink_id);
1586        }
1587
1588        // Unregister per-sink commit coordinators for any dropped pk-index iceberg sink,
1589        // including user-created sinks with arbitrary names (not just the
1590        // `__iceberg_sink_%` auto-created ones above).
1591        if !removed_iceberg_pk_index_sink_ids.is_empty() {
1592            self.iceberg_pk_index_sink_manager.unregister_jobs(
1593                removed_iceberg_pk_index_sink_ids
1594                    .into_iter()
1595                    .map(|sink_id| sink_id.as_job_id()),
1596            );
1597        }
1598
1599        // remove secrets.
1600        for secret in secret_ids {
1601            LocalSecretManager::global().remove_secret(secret);
1602        }
1603        Ok(version)
1604    }
1605
1606    /// This is used for `ALTER TABLE ADD/DROP COLUMN` / `ALTER SOURCE ADD COLUMN`.
1607    #[await_tree::instrument(boxed, "replace_streaming_job({streaming_job})")]
1608    pub async fn replace_job(
1609        &self,
1610        mut streaming_job: StreamingJob,
1611        fragment_graph: StreamFragmentGraphProto,
1612    ) -> MetaResult<NotificationVersion> {
1613        match &streaming_job {
1614            StreamingJob::Table(..)
1615            | StreamingJob::Source(..)
1616            | StreamingJob::MaterializedView(..) => {}
1617            StreamingJob::Sink(..) | StreamingJob::Index(..) => {
1618                bail_not_implemented!("schema change for {}", streaming_job.job_type_str())
1619            }
1620        }
1621
1622        let job_id = streaming_job.id();
1623
1624        let _reschedule_job_lock = self.stream_manager.reschedule_lock_write_guard().await;
1625        if let StreamingJob::Table(_, table, _) = &streaming_job
1626            && self
1627                .metadata_manager
1628                .catalog_controller
1629                .get_refresh_job_state(table.id)
1630                .await?
1631                .is_some_and(|state| state != RefreshState::Idle)
1632        {
1633            bail!(
1634                "Cannot alter table {} because it is being refreshed",
1635                table.name
1636            );
1637        }
1638        let ctx = StreamContext::from_protobuf(fragment_graph.get_ctx().unwrap());
1639
1640        // Ensure the max parallelism unchanged before replacing table.
1641        let original_max_parallelism = self
1642            .metadata_manager
1643            .get_job_max_parallelism(streaming_job.id())
1644            .await?;
1645        let fragment_graph = PbStreamFragmentGraph {
1646            max_parallelism: original_max_parallelism as _,
1647            ..fragment_graph
1648        };
1649
1650        // 1. build fragment graph.
1651        let fragment_graph = StreamFragmentGraph::new(&self.env, fragment_graph, &streaming_job)?;
1652        streaming_job.set_info_from_graph(&fragment_graph);
1653
1654        // make it immutable
1655        let streaming_job = streaming_job;
1656
1657        let auto_refresh_schema_sinks = if let StreamingJob::Table(_, table, _) = &streaming_job {
1658            let auto_refresh_schema_sinks = self
1659                .metadata_manager
1660                .catalog_controller
1661                .get_sink_auto_refresh_schema_from(table.id)
1662                .await?;
1663            if !auto_refresh_schema_sinks.is_empty() {
1664                let original_table_columns = self
1665                    .metadata_manager
1666                    .catalog_controller
1667                    .get_table_columns(table.id)
1668                    .await?;
1669                // compare column id to find newly added and removed columns
1670                let original_table_column_ids: HashSet<_> = original_table_columns
1671                    .iter()
1672                    .map(|col| col.column_id())
1673                    .collect();
1674                let new_table_column_ids: HashSet<_> = table
1675                    .columns
1676                    .iter()
1677                    .map(|col| ColumnId::new(col.column_desc.as_ref().unwrap().column_id as _))
1678                    .collect();
1679                let newly_added_columns = table
1680                    .columns
1681                    .iter()
1682                    .filter(|col| {
1683                        !original_table_column_ids.contains(&ColumnId::new(
1684                            col.column_desc.as_ref().unwrap().column_id as _,
1685                        ))
1686                    })
1687                    .map(|col| ColumnCatalog::from(col.clone()))
1688                    .collect_vec();
1689                let removed_columns = original_table_columns
1690                    .iter()
1691                    .filter(|col| !new_table_column_ids.contains(&col.column_id()))
1692                    .cloned()
1693                    .collect_vec();
1694                // Fail before any fragment rewrite or catalog persistence so a rejected ALTER
1695                // leaves no partial state.
1696                if let Some(variant_column) = first_variant_column(&newly_added_columns) {
1697                    let sink_names = auto_refresh_schema_sinks
1698                        .iter()
1699                        .map(|sink| format!("`{}`", sink.name))
1700                        .join(", ");
1701                    return Err(MetaError::invalid_parameter(format!(
1702                        "cannot add VARIANT column `{}` because sink(s) {} with auto schema refresh do not support VARIANT",
1703                        variant_column.name_with_hidden(),
1704                        sink_names,
1705                    )));
1706                }
1707                let mut sinks = Vec::with_capacity(auto_refresh_schema_sinks.len());
1708                for sink in auto_refresh_schema_sinks {
1709                    let sink_job_fragments = self
1710                        .metadata_manager
1711                        .get_job_fragments_by_id(sink.id.as_job_id())
1712                        .await?;
1713                    if sink_job_fragments.fragments.len() != 1 {
1714                        return Err(anyhow!(
1715                            "auto schema refresh sink must have only one fragment, but got {}",
1716                            sink_job_fragments.fragments.len()
1717                        )
1718                        .into());
1719                    }
1720                    let sink_ctx = sink_job_fragments.ctx;
1721                    let original_sink_fragment =
1722                        sink_job_fragments.fragments.into_values().next().unwrap();
1723                    let (new_sink_fragment, new_schema, new_log_store_table) =
1724                        rewrite_refresh_schema_sink_fragment(
1725                            &original_sink_fragment,
1726                            &sink,
1727                            &newly_added_columns,
1728                            &removed_columns,
1729                            table,
1730                            fragment_graph.table_fragment_id(),
1731                            self.env.id_gen_manager(),
1732                        )?;
1733
1734                    let streaming_job = StreamingJob::Sink(sink, None);
1735
1736                    let tmp_sink_model = self
1737                        .metadata_manager
1738                        .catalog_controller
1739                        .create_job_catalog_for_replace(&streaming_job, None, None, None)
1740                        .await?;
1741                    let tmp_sink_id = tmp_sink_model.job_id.as_sink_id();
1742                    let StreamingJob::Sink(sink, _) = streaming_job else {
1743                        unreachable!()
1744                    };
1745
1746                    sinks.push(AutoRefreshSchemaSinkContext {
1747                        tmp_sink_id,
1748                        original_sink: sink,
1749                        original_fragment: original_sink_fragment,
1750                        new_schema,
1751                        newly_add_fields: newly_added_columns
1752                            .iter()
1753                            .map(|col| Field::from(&col.column_desc))
1754                            .collect(),
1755                        removed_column_names: removed_columns
1756                            .iter()
1757                            .map(|col| col.name.clone())
1758                            .collect(),
1759                        new_fragment: new_sink_fragment,
1760                        new_log_store_table: new_log_store_table.map(Box::new),
1761                        ctx: sink_ctx,
1762                    });
1763                }
1764                Some(sinks)
1765            } else {
1766                None
1767            }
1768        } else {
1769            None
1770        };
1771
1772        let streaming_job_model = self
1773            .metadata_manager
1774            .catalog_controller
1775            .create_job_catalog_for_replace(
1776                &streaming_job,
1777                Some(&ctx),
1778                fragment_graph.specified_parallelism().as_ref(),
1779                Some(fragment_graph.max_parallelism()),
1780            )
1781            .await?;
1782        let tmp_id = streaming_job_model.job_id;
1783
1784        let tmp_sink_ids = auto_refresh_schema_sinks.as_ref().map(|sinks| {
1785            sinks
1786                .iter()
1787                .map(|sink| sink.tmp_sink_id.as_object_id())
1788                .collect_vec()
1789        });
1790
1791        tracing::debug!(id = %job_id, "building replace streaming job");
1792        let mut updated_sink_catalogs = vec![];
1793
1794        let mut drop_table_connector_ctx = None;
1795        let result: MetaResult<_> = try {
1796            let (mut ctx, mut stream_job_fragments) = self
1797                .build_replace_job(
1798                    ctx,
1799                    &streaming_job,
1800                    fragment_graph,
1801                    tmp_id,
1802                    auto_refresh_schema_sinks,
1803                    streaming_job_model,
1804                )
1805                .await?;
1806            drop_table_connector_ctx = ctx.drop_table_connector_ctx.clone();
1807            let auto_refresh_schema_sink_finish_ctx =
1808                ctx.auto_refresh_schema_sinks.as_ref().map(|sinks| {
1809                    sinks
1810                        .iter()
1811                        .map(|sink| FinishAutoRefreshSchemaSinkContext {
1812                            tmp_sink_id: sink.tmp_sink_id,
1813                            original_sink_id: sink.original_sink.id,
1814                            columns: sink.new_schema.clone(),
1815                            new_log_store_table: sink.new_log_store_table.clone(),
1816                        })
1817                        .collect()
1818                });
1819
1820            // Handle table that has incoming sinks.
1821            if let StreamingJob::Table(_, table, ..) = &streaming_job {
1822                let union_fragment = stream_job_fragments.inner.union_fragment_for_table();
1823                let upstream_infos = self
1824                    .metadata_manager
1825                    .catalog_controller
1826                    .get_all_upstream_sink_infos(table, union_fragment.fragment_id as _)
1827                    .await?;
1828                refill_upstream_sink_union_in_table(&mut union_fragment.nodes, &upstream_infos);
1829
1830                for upstream_info in &upstream_infos {
1831                    let upstream_fragment_id = upstream_info.sink_fragment_id;
1832                    ctx.upstream_fragment_downstreams
1833                        .entry(upstream_fragment_id)
1834                        .or_default()
1835                        .push(upstream_info.new_sink_downstream.clone());
1836                    if upstream_info.sink_original_target_columns.is_empty() {
1837                        updated_sink_catalogs.push(upstream_info.sink_id);
1838                    }
1839                }
1840            }
1841
1842            let replace_upstream = ctx.replace_upstream.clone();
1843
1844            if let Some(sinks) = &ctx.auto_refresh_schema_sinks {
1845                let empty_downstreams = FragmentDownstreamRelation::default();
1846                for sink in sinks {
1847                    self.metadata_manager
1848                        .catalog_controller
1849                        .prepare_streaming_job(
1850                            sink.tmp_sink_id.as_job_id(),
1851                            || [&sink.new_fragment].into_iter(),
1852                            &empty_downstreams,
1853                            true,
1854                            None,
1855                            None,
1856                        )
1857                        .await?;
1858                }
1859            }
1860
1861            self.metadata_manager
1862                .catalog_controller
1863                .prepare_stream_job_fragments(&stream_job_fragments, &streaming_job, true, None)
1864                .await?;
1865
1866            self.stream_manager
1867                .replace_stream_job(stream_job_fragments, ctx)
1868                .await?;
1869            (replace_upstream, auto_refresh_schema_sink_finish_ctx)
1870        };
1871
1872        match result {
1873            Ok((replace_upstream, auto_refresh_schema_sink_finish_ctx)) => {
1874                let version = self
1875                    .metadata_manager
1876                    .catalog_controller
1877                    .finish_replace_streaming_job(
1878                        tmp_id,
1879                        streaming_job,
1880                        replace_upstream,
1881                        SinkIntoTableContext {
1882                            updated_sink_catalogs,
1883                        },
1884                        drop_table_connector_ctx.as_ref(),
1885                        auto_refresh_schema_sink_finish_ctx,
1886                    )
1887                    .await?;
1888                if let Some(drop_table_connector_ctx) = &drop_table_connector_ctx {
1889                    self.source_manager
1890                        .apply_source_change(SourceChange::DropSource {
1891                            dropped_source_ids: vec![drop_table_connector_ctx.to_remove_source_id],
1892                        })
1893                        .await;
1894                }
1895                Ok(version)
1896            }
1897            Err(err) => {
1898                tracing::error!(id = %job_id, error = ?err.as_report(), "failed to replace job");
1899                let _ = self.metadata_manager
1900                    .catalog_controller
1901                    .try_abort_replacing_streaming_job(tmp_id, tmp_sink_ids)
1902                    .await.inspect_err(|err| {
1903                    tracing::error!(id = %job_id, error = ?err.as_report(), "failed to abort replacing job");
1904                });
1905                Err(err)
1906            }
1907        }
1908    }
1909
1910    #[await_tree::instrument(boxed, "drop_streaming_job{}({job_id})", if let DropMode::Cascade = drop_mode { "_cascade" } else { "" }
1911    )]
1912    async fn drop_streaming_job(
1913        &self,
1914        job_id: StreamingJobId,
1915        drop_mode: DropMode,
1916    ) -> MetaResult<NotificationVersion> {
1917        let (object_id, object_type) = match job_id {
1918            StreamingJobId::MaterializedView(id) => (id.as_object_id(), ObjectType::Table),
1919            StreamingJobId::Sink(id) => (id.as_object_id(), ObjectType::Sink),
1920            StreamingJobId::Table(_, id) => (id.as_object_id(), ObjectType::Table),
1921            StreamingJobId::Index(idx) => (idx.as_object_id(), ObjectType::Index),
1922        };
1923
1924        let job_status = self
1925            .metadata_manager
1926            .catalog_controller
1927            .get_streaming_job_status(job_id.id())
1928            .await?;
1929        let version = match job_status {
1930            JobStatus::Initial => {
1931                let abort_result = self
1932                    .metadata_manager
1933                    .catalog_controller
1934                    .try_abort_creating_streaming_job(job_id.id(), true)
1935                    .await?;
1936                self.iceberg_compaction_manager
1937                    .clear_maintenance_for_aborted_job(&abort_result);
1938                IGNORED_NOTIFICATION_VERSION
1939            }
1940            JobStatus::Creating => {
1941                self.stream_manager
1942                    .cancel_streaming_jobs(vec![job_id.id()])
1943                    .await?;
1944                IGNORED_NOTIFICATION_VERSION
1945            }
1946            JobStatus::Created => self.drop_object(object_type, object_id, drop_mode).await?,
1947        };
1948
1949        Ok(version)
1950    }
1951
1952    /// Builds the actor graph:
1953    /// - Add the upstream fragments to the fragment graph
1954    /// - Schedule the fragments based on their distribution
1955    /// - Expand each fragment into one or several actors
1956    /// - Construct the fragment level backfill order control.
1957    #[await_tree::instrument]
1958    pub(crate) async fn build_stream_job(
1959        &self,
1960        stream_ctx: StreamContext,
1961        mut stream_job: StreamingJob,
1962        fragment_graph: StreamFragmentGraph,
1963        resource_type: streaming_job_resource_type::ResourceType,
1964        streaming_job_model: streaming_job::Model,
1965        since_timestamp_epoch: Option<u64>,
1966    ) -> MetaResult<(CreateStreamingJobContext, StreamJobFragmentsToCreate)> {
1967        let id = stream_job.id();
1968        let max_parallelism = NonZeroUsize::new(fragment_graph.max_parallelism()).unwrap();
1969        Self::validate_specified_parallelism(
1970            fragment_graph.specified_parallelism(),
1971            fragment_graph.specified_backfill_parallelism(),
1972            max_parallelism,
1973        )?;
1974
1975        // 1. Fragment Level ordering graph
1976        let fragment_backfill_ordering = fragment_graph.create_fragment_backfill_ordering();
1977
1978        // 2. Resolve the upstream fragments, extend the fragment graph to a complete graph that
1979        // contains all information needed for building the actor graph.
1980
1981        let (snapshot_backfill_info, cross_db_snapshot_backfill_info) =
1982            fragment_graph.collect_snapshot_backfill_info()?;
1983        assert!(
1984            snapshot_backfill_info
1985                .iter()
1986                .chain([&cross_db_snapshot_backfill_info])
1987                .flat_map(|info| info.upstream_mv_table_id_to_backfill_epoch.values())
1988                .all(|backfill_epoch| backfill_epoch.is_none()),
1989            "should not set backfill epoch when initially build the job: {:?} {:?}",
1990            snapshot_backfill_info,
1991            cross_db_snapshot_backfill_info
1992        );
1993
1994        let locality_fragment_state_table_mapping =
1995            fragment_graph.find_locality_provider_fragment_state_table_mapping();
1996
1997        // check if log store exists for all cross-db upstreams
1998        self.metadata_manager
1999            .catalog_controller
2000            .validate_cross_db_snapshot_backfill(&cross_db_snapshot_backfill_info)
2001            .await?;
2002
2003        let upstream_table_ids = fragment_graph
2004            .dependent_table_ids()
2005            .iter()
2006            .filter(|id| {
2007                !cross_db_snapshot_backfill_info
2008                    .upstream_mv_table_id_to_backfill_epoch
2009                    .contains_key(*id)
2010            })
2011            .cloned()
2012            .collect();
2013
2014        let upstream_root_fragments = self
2015            .metadata_manager
2016            .get_upstream_root_fragments(&upstream_table_ids)
2017            .await?;
2018
2019        if snapshot_backfill_info.is_some() {
2020            match stream_job {
2021                StreamingJob::MaterializedView(_)
2022                | StreamingJob::Sink(..)
2023                | StreamingJob::Index(_, _) => {}
2024                StreamingJob::Table(_, _, _) | StreamingJob::Source(_) => {
2025                    return Err(
2026                        anyhow!("snapshot_backfill not enabled for table and source").into(),
2027                    );
2028                }
2029            }
2030        }
2031
2032        let complete_graph = CompleteStreamFragmentGraph::with_upstreams(
2033            fragment_graph,
2034            FragmentGraphUpstreamContext {
2035                upstream_root_fragments,
2036            },
2037            (&stream_job).into(),
2038        )?;
2039        let database_resource_group = self
2040            .metadata_manager
2041            .get_database_resource_group(stream_job.database_id())
2042            .await?;
2043        let is_serverless_backfill = matches!(
2044            &resource_type,
2045            streaming_job_resource_type::ResourceType::ServerlessBackfill(true)
2046        );
2047
2048        // 3. Build the actor graph.
2049        let actor_graph_builder = ActorGraphBuilder::new(complete_graph)?;
2050
2051        let ActorGraphBuildResult {
2052            graph,
2053            downstream_fragment_relations,
2054            upstream_fragment_downstreams,
2055            replace_upstream,
2056        } = actor_graph_builder.generate_graph()?;
2057        assert!(replace_upstream.is_empty());
2058
2059        // 4. Build the table fragments structure that will be persisted in the stream manager,
2060        // and the context that contains all information needed for building the
2061        // actors on the compute nodes.
2062
2063        let stream_job_fragments =
2064            StreamJobFragments::new(id, graph, stream_ctx.clone(), max_parallelism.get());
2065
2066        if let Some(mview_fragment) = stream_job_fragments.mview_fragment() {
2067            stream_job.set_table_vnode_count(mview_fragment.vnode_count());
2068        }
2069
2070        let new_upstream_sink = if let StreamingJob::Sink(sink, _) = &stream_job
2071            && let Ok(table_id) = sink.get_target_table()
2072        {
2073            let tables = self
2074                .metadata_manager
2075                .get_table_catalog_by_ids(&[*table_id])
2076                .await?;
2077            let target_table = tables
2078                .first()
2079                .ok_or_else(|| MetaError::catalog_id_not_found("table", *table_id))?;
2080            let sink_fragment = stream_job_fragments
2081                .sink_fragment()
2082                .ok_or_else(|| anyhow::anyhow!("sink fragment not found for sink {}", sink.id))?;
2083            let mview_fragment_id = self
2084                .metadata_manager
2085                .catalog_controller
2086                .get_mview_fragment_by_id(table_id.as_job_id())
2087                .await?;
2088            let upstream_sink_info = build_upstream_sink_info(
2089                sink.id,
2090                sink.original_target_columns.clone(),
2091                sink_fragment.fragment_id as _,
2092                target_table,
2093                mview_fragment_id,
2094            )?;
2095            Some(upstream_sink_info)
2096        } else {
2097            None
2098        };
2099
2100        let mut cdc_table_snapshot_splits = None;
2101        if let StreamingJob::Table(None, table, TableJobType::SharedCdcSource) = &stream_job
2102            && let Some((_, stream_cdc_scan)) =
2103                parallel_cdc_table_backfill_fragment(stream_job_fragments.fragments.values())
2104        {
2105            {
2106                // Create parallel splits for a CDC table. The resulted split assignments are persisted and immutable.
2107                let splits = try_init_parallel_cdc_table_snapshot_splits(
2108                    table.id,
2109                    stream_cdc_scan.cdc_table_desc.as_ref().unwrap(),
2110                    self.env.meta_store_ref(),
2111                    stream_cdc_scan.options.as_ref().unwrap(),
2112                    self.env.opts.cdc_table_split_init_insert_batch_size,
2113                    self.env.opts.cdc_table_split_init_sleep_interval_splits,
2114                    self.env.opts.cdc_table_split_init_sleep_duration_millis,
2115                )
2116                .await?;
2117                cdc_table_snapshot_splits = Some(splits);
2118            }
2119        }
2120
2121        let ctx = CreateStreamingJobContext {
2122            upstream_fragment_downstreams,
2123            database_resource_group,
2124            definition: stream_job.definition(),
2125            create_type: stream_job.create_type(),
2126            job_type: (&stream_job).into(),
2127            streaming_job: stream_job,
2128            new_upstream_sink,
2129            option: CreateStreamingJobOption {},
2130            snapshot_backfill_info,
2131            cross_db_snapshot_backfill_info,
2132            fragment_backfill_ordering,
2133            locality_fragment_state_table_mapping,
2134            cdc_table_snapshot_splits,
2135            is_serverless_backfill,
2136            resource_type,
2137            streaming_job_model: streaming_job_model.clone(),
2138            replace_sink: None,
2139            refresh_interval_sec: streaming_job_model.refresh_interval_sec.map(|s| s as u64),
2140            since_timestamp_epoch,
2141        };
2142
2143        Ok((
2144            ctx,
2145            StreamJobFragmentsToCreate {
2146                inner: stream_job_fragments,
2147                downstreams: downstream_fragment_relations,
2148            },
2149        ))
2150    }
2151
2152    /// `build_replace_table` builds a job replacement and returns the context and new job
2153    /// fragments.
2154    ///
2155    /// Note that we use a dummy ID for the new job fragments and replace it with the real one after
2156    /// replacement is finished.
2157    pub(crate) async fn build_replace_job(
2158        &self,
2159        stream_ctx: StreamContext,
2160        stream_job: &StreamingJob,
2161        mut fragment_graph: StreamFragmentGraph,
2162        tmp_job_id: JobId,
2163        auto_refresh_schema_sinks: Option<Vec<AutoRefreshSchemaSinkContext>>,
2164        streaming_job_model: streaming_job::Model,
2165    ) -> MetaResult<(ReplaceStreamJobContext, StreamJobFragmentsToCreate)> {
2166        match &stream_job {
2167            StreamingJob::Table(..)
2168            | StreamingJob::Source(..)
2169            | StreamingJob::MaterializedView(..) => {}
2170            StreamingJob::Sink(..) | StreamingJob::Index(..) => {
2171                bail_not_implemented!("schema change for {}", stream_job.job_type_str())
2172            }
2173        }
2174
2175        let id = stream_job.id();
2176
2177        // check if performing drop table connector
2178        let mut drop_table_associated_source_id = None;
2179        if let StreamingJob::Table(None, _, _) = &stream_job {
2180            drop_table_associated_source_id = self
2181                .metadata_manager
2182                .get_table_associated_source_id(id.as_mv_table_id())
2183                .await?;
2184        }
2185
2186        let old_fragments = self.metadata_manager.get_job_fragments_by_id(id).await?;
2187        let old_internal_table_ids = old_fragments.internal_table_ids();
2188
2189        // handle drop table's associated source
2190        let mut drop_table_connector_ctx = None;
2191        if let Some(to_remove_source_id) = drop_table_associated_source_id {
2192            // drop table's associated source means the fragment containing the table has just one internal table (associated source's state table)
2193            debug_assert!(old_internal_table_ids.len() == 1);
2194
2195            drop_table_connector_ctx = Some(DropTableConnectorContext {
2196                // we do not remove the original table catalog as it's still needed for the streaming job
2197                // just need to remove the ref to the state table
2198                to_change_streaming_job_id: id,
2199                to_remove_state_table_id: old_internal_table_ids[0], // asserted before
2200                to_remove_source_id,
2201            });
2202        } else if stream_job.is_materialized_view() {
2203            // If it's ALTER MV, use `state::match` to match the internal tables, which is more complicated
2204            // but more robust.
2205            let old_fragments_upstreams = self
2206                .metadata_manager
2207                .catalog_controller
2208                .upstream_fragments(old_fragments.fragment_ids())
2209                .await?;
2210
2211            let old_state_graph =
2212                state_match::Graph::from_existing(&old_fragments, &old_fragments_upstreams);
2213            let new_state_graph = state_match::Graph::from_building(&fragment_graph);
2214            let result = state_match::match_graph(&new_state_graph, &old_state_graph)
2215                .context("incompatible altering on the streaming job states")?;
2216
2217            fragment_graph.fit_internal_table_ids_with_mapping(result.table_matches);
2218            fragment_graph.fit_snapshot_backfill_epochs(result.snapshot_backfill_epochs);
2219        } else {
2220            // If it's ALTER TABLE or SOURCE, use a trivial table id matching algorithm to keep the original behavior.
2221            // TODO(alter-mv): this is actually a special case of ALTER MV, can we merge the two branches?
2222            let old_internal_tables = self
2223                .metadata_manager
2224                .get_table_catalog_by_ids(&old_internal_table_ids)
2225                .await?;
2226            fragment_graph.fit_internal_tables_trivial(old_internal_tables)?;
2227        }
2228
2229        // 1. Resolve the edges to the downstream fragments, extend the fragment graph to a complete
2230        // graph that contains all information needed for building the actor graph.
2231        let original_root_fragment = old_fragments
2232            .root_fragment()
2233            .expect("root fragment not found");
2234
2235        let job_type = StreamingJobType::from(stream_job);
2236
2237        // Extract the downstream fragments from the fragment graph.
2238        let mut downstream_fragments = self.metadata_manager.get_downstream_fragments(id).await?;
2239
2240        if let Some(auto_refresh_schema_sinks) = &auto_refresh_schema_sinks {
2241            let mut remaining_fragment: HashSet<_> = auto_refresh_schema_sinks
2242                .iter()
2243                .map(|sink| sink.original_fragment.fragment_id)
2244                .collect();
2245            for (_, downstream_fragment) in &mut downstream_fragments {
2246                if let Some(sink) = auto_refresh_schema_sinks.iter().find(|sink| {
2247                    sink.original_fragment.fragment_id == downstream_fragment.fragment_id
2248                }) {
2249                    assert!(remaining_fragment.remove(&downstream_fragment.fragment_id));
2250                    // Actor locations will be resolved inside barrier worker during rendering.
2251                    // For now, just replace fragment info and nodes.
2252                    *downstream_fragment = sink.new_fragment.clone();
2253                }
2254            }
2255            assert!(remaining_fragment.is_empty());
2256        }
2257
2258        // build complete graph based on the table job type
2259        let complete_graph = match &job_type {
2260            StreamingJobType::Table(TableJobType::General) | StreamingJobType::Source => {
2261                CompleteStreamFragmentGraph::with_downstreams(
2262                    fragment_graph,
2263                    FragmentGraphDownstreamContext {
2264                        original_root_fragment_id: original_root_fragment.fragment_id,
2265                        downstream_fragments,
2266                    },
2267                    job_type,
2268                )?
2269            }
2270            StreamingJobType::Table(TableJobType::SharedCdcSource)
2271            | StreamingJobType::MaterializedView => {
2272                // CDC tables or materialized views can have upstream jobs as well.
2273                let upstream_root_fragments = self
2274                    .metadata_manager
2275                    .get_upstream_root_fragments(fragment_graph.dependent_table_ids())
2276                    .await?;
2277
2278                CompleteStreamFragmentGraph::with_upstreams_and_downstreams(
2279                    fragment_graph,
2280                    FragmentGraphUpstreamContext {
2281                        upstream_root_fragments,
2282                    },
2283                    FragmentGraphDownstreamContext {
2284                        original_root_fragment_id: original_root_fragment.fragment_id,
2285                        downstream_fragments,
2286                    },
2287                    job_type,
2288                )?
2289            }
2290            _ => unreachable!(),
2291        };
2292
2293        let resource_group = self
2294            .metadata_manager
2295            .get_database_resource_group(stream_job.database_id())
2296            .await?;
2297
2298        let actor_graph_builder = ActorGraphBuilder::new(complete_graph)?;
2299
2300        let ActorGraphBuildResult {
2301            graph,
2302            downstream_fragment_relations,
2303            upstream_fragment_downstreams,
2304            mut replace_upstream,
2305        } = actor_graph_builder.generate_graph()?;
2306
2307        // general table & source does not have upstream job, so the dispatchers should be empty
2308        if matches!(
2309            job_type,
2310            StreamingJobType::Source | StreamingJobType::Table(TableJobType::General)
2311        ) {
2312            assert!(upstream_fragment_downstreams.is_empty());
2313        }
2314
2315        // 3. Build the table fragments structure that will be persisted in the stream manager, and
2316        // the context that contains all information needed for building the actors on the compute
2317        // nodes.
2318        let stream_job_fragments =
2319            StreamJobFragments::new(tmp_job_id, graph, stream_ctx, old_fragments.max_parallelism);
2320
2321        if let Some(sinks) = &auto_refresh_schema_sinks {
2322            for sink in sinks {
2323                replace_upstream
2324                    .remove(&sink.new_fragment.fragment_id)
2325                    .expect("should exist");
2326            }
2327        }
2328
2329        // Note: no need to set `vnode_count` as it's already set by the frontend.
2330        // See `get_replace_table_plan`.
2331
2332        let ctx = ReplaceStreamJobContext {
2333            old_fragments,
2334            replace_upstream,
2335            upstream_fragment_downstreams,
2336            streaming_job: stream_job.clone(),
2337            database_resource_group: resource_group,
2338            tmp_id: tmp_job_id,
2339            drop_table_connector_ctx,
2340            auto_refresh_schema_sinks,
2341            streaming_job_model,
2342        };
2343
2344        Ok((
2345            ctx,
2346            StreamJobFragmentsToCreate {
2347                inner: stream_job_fragments,
2348                downstreams: downstream_fragment_relations,
2349            },
2350        ))
2351    }
2352
2353    async fn alter_name(
2354        &self,
2355        relation: alter_name_request::Object,
2356        new_name: &str,
2357    ) -> MetaResult<NotificationVersion> {
2358        let (obj_type, id): (ObjectType, ObjectId) = match relation {
2359            alter_name_request::Object::TableId(id) => (ObjectType::Table, id.into()),
2360            alter_name_request::Object::ViewId(id) => (ObjectType::View, id.into()),
2361            alter_name_request::Object::IndexId(id) => (ObjectType::Index, id.into()),
2362            alter_name_request::Object::SinkId(id) => (ObjectType::Sink, id.into()),
2363            alter_name_request::Object::SourceId(id) => (ObjectType::Source, id.into()),
2364            alter_name_request::Object::SchemaId(id) => (ObjectType::Schema, id.into()),
2365            alter_name_request::Object::DatabaseId(id) => (ObjectType::Database, id.into()),
2366            alter_name_request::Object::SubscriptionId(id) => (ObjectType::Subscription, id.into()),
2367        };
2368        self.metadata_manager
2369            .catalog_controller
2370            .alter_name(obj_type, id, new_name)
2371            .await
2372    }
2373
2374    async fn alter_swap_rename(
2375        &self,
2376        object: alter_swap_rename_request::Object,
2377    ) -> MetaResult<NotificationVersion> {
2378        let (obj_type, src_id, dst_id) = match object {
2379            alter_swap_rename_request::Object::Schema(_) => unimplemented!("schema swap"),
2380            alter_swap_rename_request::Object::Table(objs) => {
2381                let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2382                (ObjectType::Table, src_id, dst_id)
2383            }
2384            alter_swap_rename_request::Object::View(objs) => {
2385                let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2386                (ObjectType::View, src_id, dst_id)
2387            }
2388            alter_swap_rename_request::Object::Source(objs) => {
2389                let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2390                (ObjectType::Source, src_id, dst_id)
2391            }
2392            alter_swap_rename_request::Object::Sink(objs) => {
2393                let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2394                (ObjectType::Sink, src_id, dst_id)
2395            }
2396            alter_swap_rename_request::Object::Subscription(objs) => {
2397                let (src_id, dst_id) = (objs.src_object_id, objs.dst_object_id);
2398                (ObjectType::Subscription, src_id, dst_id)
2399            }
2400        };
2401
2402        self.metadata_manager
2403            .catalog_controller
2404            .alter_swap_rename(obj_type, src_id, dst_id)
2405            .await
2406    }
2407
2408    async fn alter_owner(
2409        &self,
2410        object: Object,
2411        owner_id: UserId,
2412    ) -> MetaResult<NotificationVersion> {
2413        let (obj_type, id): (ObjectType, ObjectId) = match object {
2414            Object::TableId(id) => (ObjectType::Table, id.into()),
2415            Object::ViewId(id) => (ObjectType::View, id.into()),
2416            Object::SourceId(id) => (ObjectType::Source, id.into()),
2417            Object::SinkId(id) => (ObjectType::Sink, id.into()),
2418            Object::SchemaId(id) => (ObjectType::Schema, id.into()),
2419            Object::DatabaseId(id) => (ObjectType::Database, id.into()),
2420            Object::SubscriptionId(id) => (ObjectType::Subscription, id.into()),
2421            Object::ConnectionId(id) => (ObjectType::Connection, id.into()),
2422            Object::FunctionId(id) => (ObjectType::Function, id.into()),
2423            Object::SecretId(id) => (ObjectType::Secret, id.into()),
2424        };
2425        self.metadata_manager
2426            .catalog_controller
2427            .alter_owner(obj_type, id, owner_id as _)
2428            .await
2429    }
2430
2431    async fn alter_set_schema(
2432        &self,
2433        object: alter_set_schema_request::Object,
2434        new_schema_id: SchemaId,
2435    ) -> MetaResult<NotificationVersion> {
2436        let (obj_type, id): (ObjectType, ObjectId) = match object {
2437            alter_set_schema_request::Object::TableId(id) => (ObjectType::Table, id.into()),
2438            alter_set_schema_request::Object::ViewId(id) => (ObjectType::View, id.into()),
2439            alter_set_schema_request::Object::SourceId(id) => (ObjectType::Source, id.into()),
2440            alter_set_schema_request::Object::SinkId(id) => (ObjectType::Sink, id.into()),
2441            alter_set_schema_request::Object::FunctionId(id) => (ObjectType::Function, id.into()),
2442            alter_set_schema_request::Object::ConnectionId(id) => {
2443                (ObjectType::Connection, id.into())
2444            }
2445            alter_set_schema_request::Object::SubscriptionId(id) => {
2446                (ObjectType::Subscription, id.into())
2447            }
2448        };
2449        self.metadata_manager
2450            .catalog_controller
2451            .alter_schema(obj_type, id, new_schema_id)
2452            .await
2453    }
2454
2455    pub async fn wait(&self, job_id: Option<JobId>) -> MetaResult<WaitVersion> {
2456        if let Some(job_id) = job_id {
2457            let database_id = self
2458                .metadata_manager
2459                .catalog_controller
2460                .get_object_database_id(job_id)
2461                .await?;
2462            let catalog_version = self
2463                .metadata_manager
2464                .wait_streaming_job_finished(database_id, job_id)
2465                .await?;
2466            let hummock_version_id = self.barrier_manager.get_hummock_version_id().await;
2467            return Ok(WaitVersion {
2468                catalog_version,
2469                hummock_version_id,
2470            });
2471        }
2472
2473        let timeout_ms = 2 * 60 * 60 * 1000;
2474        let poll_interval = Duration::from_millis(100);
2475        for _ in 0..(timeout_ms / poll_interval.as_millis() as usize) {
2476            let creating_jobs = self
2477                .metadata_manager
2478                .catalog_controller
2479                .list_creating_jobs(true, None)
2480                .await?;
2481            if creating_jobs.is_empty() {
2482                let catalog_version = self
2483                    .metadata_manager
2484                    .catalog_controller
2485                    .notify_frontend_trivial()
2486                    .await;
2487                let hummock_version_id = self.barrier_manager.get_hummock_version_id().await;
2488                return Ok(WaitVersion {
2489                    catalog_version,
2490                    hummock_version_id,
2491                });
2492            }
2493
2494            sleep(poll_interval).await;
2495        }
2496        Err(MetaError::cancelled(format!(
2497            "timeout after {timeout_ms}ms"
2498        )))
2499    }
2500
2501    async fn comment_on(&self, comment: Comment) -> MetaResult<NotificationVersion> {
2502        self.metadata_manager
2503            .catalog_controller
2504            .comment_on(comment)
2505            .await
2506    }
2507
2508    async fn alter_streaming_job_config(
2509        &self,
2510        job_id: JobId,
2511        entries_to_add: HashMap<String, String>,
2512        keys_to_remove: Vec<String>,
2513    ) -> MetaResult<NotificationVersion> {
2514        self.metadata_manager
2515            .catalog_controller
2516            .alter_streaming_job_config(job_id, entries_to_add, keys_to_remove)
2517            .await
2518    }
2519}
2520
2521fn report_create_object(
2522    job_id: JobId,
2523    event_name: &str,
2524    obj_type: PbTelemetryDatabaseObject,
2525    connector_name: Option<String>,
2526    attr_info: Option<jsonbb::Value>,
2527) {
2528    report_event(
2529        PbTelemetryEventStage::CreateStreamJob,
2530        event_name,
2531        job_id.as_raw_id() as _,
2532        connector_name,
2533        Some(obj_type),
2534        attr_info,
2535    );
2536}
2537
2538pub fn build_upstream_sink_info(
2539    sink_id: SinkId,
2540    original_target_columns: Vec<PbColumnCatalog>,
2541    sink_fragment_id: FragmentId,
2542    target_table: &PbTable,
2543    target_fragment_id: FragmentId,
2544) -> MetaResult<UpstreamSinkInfo> {
2545    let sink_columns = if !original_target_columns.is_empty() {
2546        original_target_columns.clone()
2547    } else {
2548        // This is due to the fact that the value did not exist in earlier versions,
2549        // which means no schema changes such as `ADD/DROP COLUMN` have been made to the table.
2550        // Therefore the columns of the table at this point are `original_target_columns`.
2551        // This value of sink will be filled on the meta.
2552        target_table.columns.clone()
2553    };
2554
2555    let sink_output_fields = sink_columns
2556        .iter()
2557        .map(|col| Field::from(col.column_desc.as_ref().unwrap()).to_prost())
2558        .collect_vec();
2559    let output_indices = (0..sink_output_fields.len())
2560        .map(|i| i as u32)
2561        .collect_vec();
2562
2563    let dist_key_indices: anyhow::Result<Vec<u32>> = try {
2564        let sink_idx_by_col_id = sink_columns
2565            .iter()
2566            .enumerate()
2567            .map(|(idx, col)| {
2568                let column_id = col.column_desc.as_ref().unwrap().column_id;
2569                (column_id, idx as u32)
2570            })
2571            .collect::<HashMap<_, _>>();
2572        target_table
2573            .distribution_key
2574            .iter()
2575            .map(|dist_idx| {
2576                let column_id = target_table.columns[*dist_idx as usize]
2577                    .column_desc
2578                    .as_ref()
2579                    .unwrap()
2580                    .column_id;
2581                let sink_idx = sink_idx_by_col_id
2582                    .get(&column_id)
2583                    .ok_or_else(|| anyhow::anyhow!("column id {} not found in sink", column_id))?;
2584                Ok(*sink_idx)
2585            })
2586            .collect::<anyhow::Result<Vec<_>>>()?
2587    };
2588    let dist_key_indices =
2589        dist_key_indices.map_err(|e| e.context("failed to get distribution key indices"))?;
2590    let downstream_fragment_id = target_fragment_id as _;
2591    let new_downstream_relation = DownstreamFragmentRelation {
2592        downstream_fragment_id,
2593        dispatcher_type: DispatcherType::Hash,
2594        dist_key_indices,
2595        output_mapping: PbDispatchOutputMapping::simple(output_indices),
2596    };
2597    let current_target_columns = target_table.get_columns();
2598    let project_exprs = build_select_node_list(&sink_columns, current_target_columns)?;
2599    Ok(UpstreamSinkInfo {
2600        sink_id,
2601        sink_fragment_id: sink_fragment_id as _,
2602        sink_output_fields,
2603        sink_original_target_columns: original_target_columns,
2604        project_exprs,
2605        new_sink_downstream: new_downstream_relation,
2606    })
2607}
2608
2609pub fn refill_upstream_sink_union_in_table(
2610    union_fragment_root: &mut PbStreamNode,
2611    upstream_sink_infos: &Vec<UpstreamSinkInfo>,
2612) {
2613    visit_stream_node_cont_mut(union_fragment_root, |node| {
2614        if let Some(NodeBody::UpstreamSinkUnion(upstream_sink_union)) = &mut node.node_body {
2615            let init_upstreams = upstream_sink_infos
2616                .iter()
2617                .map(|info| PbUpstreamSinkInfo {
2618                    upstream_fragment_id: info.sink_fragment_id,
2619                    sink_output_schema: info.sink_output_fields.clone(),
2620                    project_exprs: info.project_exprs.clone(),
2621                })
2622                .collect();
2623            upstream_sink_union.init_upstreams = init_upstreams;
2624            false
2625        } else {
2626            true
2627        }
2628    });
2629}
2630
2631#[cfg(test)]
2632mod tests {
2633    use std::num::NonZeroUsize;
2634
2635    use super::*;
2636
2637    #[test]
2638    fn test_validate_specified_parallelism_accepts_within_max() {
2639        DdlController::validate_specified_parallelism(
2640            Some(NonZeroUsize::new(4).unwrap()),
2641            Some(NonZeroUsize::new(8).unwrap()),
2642            NonZeroUsize::new(8).unwrap(),
2643        )
2644        .unwrap();
2645    }
2646
2647    #[test]
2648    fn test_validate_specified_parallelism_rejects_parallelism_over_max() {
2649        let result = DdlController::validate_specified_parallelism(
2650            Some(NonZeroUsize::new(9).unwrap()),
2651            None,
2652            NonZeroUsize::new(8).unwrap(),
2653        );
2654        assert!(matches!(
2655            result,
2656            Err(ref e) if matches!(e.inner(), MetaErrorInner::InvalidParameter(_))
2657        ));
2658    }
2659
2660    #[test]
2661    fn test_validate_specified_parallelism_rejects_backfill_parallelism_over_max() {
2662        let result = DdlController::validate_specified_parallelism(
2663            None,
2664            Some(NonZeroUsize::new(9).unwrap()),
2665            NonZeroUsize::new(8).unwrap(),
2666        );
2667        assert!(matches!(
2668            result,
2669            Err(ref e) if matches!(e.inner(), MetaErrorInner::InvalidParameter(_))
2670        ));
2671    }
2672}