Skip to main content

risingwave_meta/controller/catalog/
create_op.rs

1// Copyright 2024 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 risingwave_common::cast::datetime_to_timestamp_millis;
16use risingwave_common::system_param::{OverrideValidate, Validate};
17use risingwave_common::util::epoch::Epoch;
18use risingwave_pb::common::PbObjectType;
19use risingwave_pb::meta::PbObjectDependency;
20
21use super::*;
22use crate::barrier::SnapshotBackfillInfo;
23
24impl CatalogController {
25    pub(crate) async fn create_object(
26        txn: &DatabaseTransaction,
27        obj_type: ObjectType,
28        owner_id: UserId,
29        belong_to_oid: Option<ObjectId>,
30    ) -> MetaResult<object::Model> {
31        let (database_id, schema_id) = match (obj_type, belong_to_oid) {
32            (ObjectType::Database, None) => (None, None),
33            (ObjectType::Database, Some(_)) => {
34                return Err(MetaError::invalid_parameter(
35                    "database object must not have a parent",
36                ));
37            }
38            (ObjectType::Schema, Some(database_oid)) => {
39                let database = Object::find_by_id(database_oid)
40                    .one(txn)
41                    .await?
42                    .ok_or_else(|| MetaError::catalog_id_not_found("database", database_oid))?;
43                if database.obj_type != ObjectType::Database {
44                    return Err(MetaError::invalid_parameter(
45                        "schema object must belong to a database",
46                    ));
47                }
48                (Some(database_oid.as_database_id()), None)
49            }
50            (ObjectType::Schema, None) => {
51                return Err(MetaError::invalid_parameter(
52                    "schema object must belong to a database",
53                ));
54            }
55            (_, Some(parent_oid)) => {
56                let parent = Object::find_by_id(parent_oid)
57                    .one(txn)
58                    .await?
59                    .ok_or_else(|| MetaError::catalog_id_not_found("parent object", parent_oid))?;
60                if parent.obj_type == ObjectType::Schema {
61                    (parent.database_id, Some(parent.oid.as_schema_id()))
62                } else {
63                    (parent.database_id, parent.schema_id)
64                }
65            }
66            (_, None) => {
67                return Err(MetaError::invalid_parameter(
68                    "non-database object must have a parent",
69                ));
70            }
71        };
72        let active_db = object::ActiveModel {
73            oid: Default::default(),
74            obj_type: Set(obj_type),
75            owner_id: Set(owner_id),
76            schema_id: Set(schema_id),
77            database_id: Set(database_id),
78            belong_to_oid: Set(belong_to_oid),
79            initialized_at: Default::default(),
80            created_at: Default::default(),
81            initialized_at_cluster_version: Set(Some(current_cluster_version())),
82            created_at_cluster_version: Set(Some(current_cluster_version())),
83        };
84        Ok(active_db.insert(txn).await?)
85    }
86
87    pub async fn create_database(
88        &self,
89        db: PbDatabase,
90    ) -> MetaResult<(NotificationVersion, risingwave_meta_model::database::Model)> {
91        // validate first
92        if let Some(ref interval) = db.barrier_interval_ms {
93            OverrideValidate::barrier_interval_ms(interval).map_err(|e| anyhow::anyhow!(e))?;
94        }
95        if let Some(ref frequency) = db.checkpoint_frequency {
96            OverrideValidate::checkpoint_frequency(frequency).map_err(|e| anyhow::anyhow!(e))?;
97        }
98
99        let inner = self.inner.write().await;
100        let owner_id = db.owner as _;
101        let txn = inner.db.begin().await?;
102        ensure_user_id(owner_id, &txn).await?;
103        check_database_name_duplicate(&db.name, &txn).await?;
104
105        let db_obj = Self::create_object(&txn, ObjectType::Database, owner_id, None).await?;
106        let mut db: database::ActiveModel = db.into();
107        db.database_id = Set(db_obj.oid.as_database_id());
108        let db = db.insert(&txn).await?;
109
110        let mut schemas = vec![];
111        for schema_name in iter::once(DEFAULT_SCHEMA_NAME).chain(SYSTEM_SCHEMAS) {
112            let schema_obj =
113                Self::create_object(&txn, ObjectType::Schema, owner_id, Some(db_obj.oid)).await?;
114            let schema = schema::ActiveModel {
115                schema_id: Set(schema_obj.oid.as_schema_id()),
116                name: Set(schema_name.into()),
117            };
118            let schema = schema.insert(&txn).await?;
119            schemas.push(ObjectModel(schema, schema_obj, None).into());
120        }
121        txn.commit().await?;
122
123        let mut version = self
124            .notify_frontend(
125                NotificationOperation::Add,
126                NotificationInfo::Database(ObjectModel(db.clone(), db_obj, None).into()),
127            )
128            .await;
129        for schema in schemas {
130            version = self
131                .notify_frontend(NotificationOperation::Add, NotificationInfo::Schema(schema))
132                .await;
133        }
134
135        Ok((version, db))
136    }
137
138    pub async fn create_schema(&self, schema: PbSchema) -> MetaResult<NotificationVersion> {
139        let inner = self.inner.write().await;
140        let owner_id = schema.owner as _;
141        let txn = inner.db.begin().await?;
142        ensure_user_id(owner_id, &txn).await?;
143        ensure_object_id(ObjectType::Database, schema.database_id, &txn).await?;
144        check_schema_name_duplicate(&schema.name, schema.database_id, &txn).await?;
145
146        let schema_obj = Self::create_object(
147            &txn,
148            ObjectType::Schema,
149            owner_id,
150            Some(schema.database_id.as_object_id()),
151        )
152        .await?;
153        let mut schema: schema::ActiveModel = schema.into();
154        schema.schema_id = Set(schema_obj.oid.as_schema_id());
155        let schema = schema.insert(&txn).await?;
156
157        let updated_user_info =
158            grant_default_privileges_automatically(&txn, schema_obj.oid).await?;
159
160        txn.commit().await?;
161
162        let mut version = self
163            .notify_frontend(
164                NotificationOperation::Add,
165                NotificationInfo::Schema(ObjectModel(schema, schema_obj, None).into()),
166            )
167            .await;
168
169        // notify default privileges for schemas
170        if !updated_user_info.is_empty() {
171            version = self.notify_users_update(updated_user_info).await;
172        }
173
174        Ok(version)
175    }
176
177    pub async fn create_subscription_catalog(
178        &self,
179        pb_subscription: &mut PbSubscription,
180    ) -> MetaResult<()> {
181        let inner = self.inner.write().await;
182        let txn = inner.db.begin().await?;
183
184        ensure_user_id(pb_subscription.owner as _, &txn).await?;
185        ensure_object_id(ObjectType::Database, pb_subscription.database_id, &txn).await?;
186        ensure_object_id(ObjectType::Schema, pb_subscription.schema_id, &txn).await?;
187        check_subscription_name_duplicate(pb_subscription, &txn).await?;
188
189        let obj = Self::create_object(
190            &txn,
191            ObjectType::Subscription,
192            pb_subscription.owner as _,
193            Some(pb_subscription.schema_id.as_object_id()),
194        )
195        .await?;
196        pb_subscription.id = obj.oid.as_subscription_id();
197        let subscription: subscription::ActiveModel = pb_subscription.clone().into();
198        Subscription::insert(subscription).exec(&txn).await?;
199
200        // record object dependency.
201        ObjectDependency::insert(object_dependency::ActiveModel {
202            oid: Set(pb_subscription.dependent_table_id.into()),
203            used_by: Set(pb_subscription.id.into()),
204            ..Default::default()
205        })
206        .exec(&txn)
207        .await?;
208        txn.commit().await?;
209        Ok(())
210    }
211
212    pub async fn create_source(
213        &self,
214        mut pb_source: PbSource,
215        iceberg_table_id: Option<TableId>,
216    ) -> MetaResult<(SourceId, NotificationVersion)> {
217        let mut inner = self.inner.write().await;
218        let owner_id = pb_source.owner as _;
219        let txn = inner.db.begin().await?;
220        ensure_user_id(owner_id, &txn).await?;
221        ensure_object_id(ObjectType::Database, pb_source.database_id, &txn).await?;
222        ensure_object_id(ObjectType::Schema, pb_source.schema_id, &txn).await?;
223        check_relation_name_duplicate(
224            &pb_source.name,
225            pb_source.database_id,
226            pb_source.schema_id,
227            &txn,
228        )
229        .await?;
230
231        // handle secret ref
232        let secret_ids = get_referred_secret_ids_from_source(&pb_source)?;
233        let connection_ids = get_referred_connection_ids_from_source(&pb_source);
234        let cdc_source_id = pb_source
235            .info
236            .as_ref()
237            .and_then(|info| info.external_table.as_ref())
238            .map(|table| table.source_id);
239        if let Some(cdc_source_id) = cdc_source_id {
240            ensure_object_id(ObjectType::Source, cdc_source_id, &txn).await?;
241        }
242
243        let source_obj = Self::create_object(
244            &txn,
245            ObjectType::Source,
246            owner_id,
247            iceberg_table_id
248                .map(TableId::as_object_id)
249                .or(Some(pb_source.schema_id.as_object_id())),
250        )
251        .await?;
252        let source_id = source_obj.oid.as_source_id();
253        pb_source.id = source_id;
254        if let Some(external_table) = pb_source
255            .info
256            .as_mut()
257            .and_then(|info| info.external_table.as_mut())
258        {
259            external_table.table_id = source_id.as_cdc_table_id();
260            // The connect properties were only needed to validate the descriptor. Consumers
261            // resolve them from the upstream shared source when they are planned, so persisting
262            // a copy here would only make credentials and secret references stale after
263            // `ALTER SOURCE`.
264            external_table.connect_properties.clear();
265            external_table.secret_refs.clear();
266        }
267        let source: source::ActiveModel = pb_source.clone().into();
268        Source::insert(source).exec(&txn).await?;
269
270        // Add secret, connection, and upstream shared CDC source dependencies.
271        let dep_relation_ids = secret_ids
272            .iter()
273            .copied()
274            .map_into()
275            .chain(connection_ids.iter().copied().map_into())
276            .chain(cdc_source_id.map(|id| id.as_object_id()))
277            .collect_vec();
278        let dependencies = if !dep_relation_ids.is_empty() {
279            ObjectDependency::insert_many(dep_relation_ids.into_iter().map(|id| {
280                object_dependency::ActiveModel {
281                    oid: Set(id),
282                    used_by: Set(source_id.as_object_id()),
283                    ..Default::default()
284                }
285            }))
286            .exec(&txn)
287            .await?;
288            secret_ids
289                .iter()
290                .map(|id| PbObjectDependency {
291                    object_id: source_id.as_object_id(),
292                    referenced_object_id: id.as_object_id(),
293                    referenced_object_type: PbObjectType::Secret as _,
294                })
295                .chain(connection_ids.iter().map(|id| PbObjectDependency {
296                    object_id: source_id.as_object_id(),
297                    referenced_object_id: id.as_object_id(),
298                    referenced_object_type: PbObjectType::Connection as _,
299                }))
300                .chain(cdc_source_id.map(|id| PbObjectDependency {
301                    object_id: source_id.as_object_id(),
302                    referenced_object_id: id.as_object_id(),
303                    referenced_object_type: PbObjectType::Source as _,
304                }))
305                .collect_vec()
306        } else {
307            vec![]
308        };
309
310        let mut job_notifications = vec![];
311        let mut updated_user_info = vec![];
312        // check if it belongs to iceberg table
313        if let Some(table_id) = iceberg_table_id {
314            // 1. finish iceberg table job.
315            let table_notifications = self
316                .finish_streaming_job_inner(&txn, table_id.as_job_id())
317                .await?;
318            job_notifications.push((table_id.as_job_id(), table_notifications));
319
320            // 2. finish iceberg sink job.
321            let sink_id = Sink::find()
322                .select_only()
323                .column(sink::Column::SinkId)
324                .join(JoinType::InnerJoin, sink::Relation::Object.def())
325                .filter(
326                    object::Column::BelongToOid
327                        .eq(table_id)
328                        .and(object::Column::ObjType.eq(ObjectType::Sink)),
329                )
330                .into_tuple::<SinkId>()
331                .one(&txn)
332                .await?
333                .ok_or_else(|| {
334                    MetaError::catalog_id_not_found("iceberg sink for table", table_id)
335                })?;
336            let sink_job_id = sink_id.as_job_id();
337            let sink_notifications = self.finish_streaming_job_inner(&txn, sink_job_id).await?;
338            job_notifications.push((sink_job_id, sink_notifications));
339        } else {
340            updated_user_info = grant_default_privileges_automatically(&txn, source_id).await?;
341        }
342
343        txn.commit().await?;
344
345        for (job_id, (op, objects, user_info, dependencies)) in job_notifications {
346            let mut version = self
347                .notify_frontend(
348                    op,
349                    NotificationInfo::ObjectGroup(PbObjectGroup {
350                        objects,
351                        dependencies,
352                    }),
353                )
354                .await;
355            if !user_info.is_empty() {
356                version = self.notify_users_update(user_info).await;
357            }
358            inner
359                .creating_table_finish_notifier
360                .values_mut()
361                .for_each(|creating_tables| {
362                    if let Some(txs) = creating_tables.remove(&job_id) {
363                        for tx in txs {
364                            let _ = tx.send(Ok(version));
365                        }
366                    }
367                });
368        }
369
370        let mut version = self
371            .notify_frontend(
372                NotificationOperation::Add,
373                NotificationInfo::ObjectGroup(PbObjectGroup {
374                    objects: vec![PbObject {
375                        object_info: Some(PbObjectInfo::Source(pb_source)),
376                    }],
377                    dependencies,
378                }),
379            )
380            .await;
381
382        // notify default privileges for source
383        if !updated_user_info.is_empty() {
384            version = self.notify_users_update(updated_user_info).await;
385        }
386
387        Ok((source_id, version))
388    }
389
390    pub async fn create_function(
391        &self,
392        mut pb_function: PbFunction,
393    ) -> MetaResult<NotificationVersion> {
394        let inner = self.inner.write().await;
395        let owner_id = pb_function.owner as _;
396        let txn = inner.db.begin().await?;
397        ensure_user_id(owner_id, &txn).await?;
398        ensure_object_id(ObjectType::Database, pb_function.database_id, &txn).await?;
399        ensure_object_id(ObjectType::Schema, pb_function.schema_id, &txn).await?;
400        check_function_signature_duplicate(&pb_function, &txn).await?;
401
402        let function_obj = Self::create_object(
403            &txn,
404            ObjectType::Function,
405            owner_id,
406            Some(pb_function.schema_id.as_object_id()),
407        )
408        .await?;
409        pb_function.id = function_obj.oid.as_function_id();
410        pb_function.created_at_epoch = Some(
411            Epoch::from_unix_millis(datetime_to_timestamp_millis(function_obj.created_at) as _).0,
412        );
413        pb_function.created_at_cluster_version = function_obj.created_at_cluster_version;
414        let function: function::ActiveModel = pb_function.clone().into();
415        Function::insert(function).exec(&txn).await?;
416
417        let updated_user_info =
418            grant_default_privileges_automatically(&txn, function_obj.oid).await?;
419
420        txn.commit().await?;
421
422        let mut version = self
423            .notify_frontend(
424                NotificationOperation::Add,
425                NotificationInfo::Function(pb_function),
426            )
427            .await;
428
429        // notify default privileges for functions
430        if !updated_user_info.is_empty() {
431            version = self.notify_users_update(updated_user_info).await;
432        }
433
434        Ok(version)
435    }
436
437    pub async fn create_connection(
438        &self,
439        mut pb_connection: PbConnection,
440    ) -> MetaResult<NotificationVersion> {
441        let inner = self.inner.write().await;
442        let owner_id = pb_connection.owner as _;
443        let txn = inner.db.begin().await?;
444        ensure_user_id(owner_id, &txn).await?;
445        ensure_object_id(ObjectType::Database, pb_connection.database_id, &txn).await?;
446        ensure_object_id(ObjectType::Schema, pb_connection.schema_id, &txn).await?;
447        check_connection_name_duplicate(&pb_connection, &txn).await?;
448
449        let mut dep_secrets: HashSet<SecretId> = HashSet::new();
450        if let Some(ConnectionInfo::ConnectionParams(params)) = &pb_connection.info {
451            dep_secrets.extend(
452                params
453                    .secret_refs
454                    .values()
455                    .map(|secret_ref| secret_ref.secret_id),
456            );
457        }
458
459        let conn_obj = Self::create_object(
460            &txn,
461            ObjectType::Connection,
462            owner_id,
463            Some(pb_connection.schema_id.as_object_id()),
464        )
465        .await?;
466        pb_connection.id = conn_obj.oid.as_connection_id();
467        let connection: connection::ActiveModel = pb_connection.clone().into();
468        Connection::insert(connection).exec(&txn).await?;
469
470        for secret_id in &dep_secrets {
471            ObjectDependency::insert(object_dependency::ActiveModel {
472                oid: Set(secret_id.as_object_id()),
473                used_by: Set(conn_obj.oid),
474                ..Default::default()
475            })
476            .exec(&txn)
477            .await?;
478        }
479
480        let updated_user_info = grant_default_privileges_automatically(&txn, conn_obj.oid).await?;
481
482        txn.commit().await?;
483
484        {
485            // call meta telemetry here to report the connection creation
486            report_event(
487                PbTelemetryEventStage::Unspecified,
488                "connection_create",
489                pb_connection.get_id().as_raw_id() as i64,
490                {
491                    pb_connection.info.as_ref().and_then(|info| match info {
492                        ConnectionInfo::ConnectionParams(params) => {
493                            Some(params.connection_type().as_str_name().to_owned())
494                        }
495                        _ => None,
496                    })
497                },
498                None,
499                None,
500            );
501        }
502
503        let mut version = self
504            .notify_frontend(
505                NotificationOperation::Add,
506                NotificationInfo::ObjectGroup(PbObjectGroup {
507                    objects: vec![PbObject {
508                        object_info: Some(PbObjectInfo::Connection(pb_connection)),
509                    }],
510                    dependencies: dep_secrets
511                        .iter()
512                        .map(|secret_id| PbObjectDependency {
513                            object_id: conn_obj.oid,
514                            referenced_object_id: secret_id.as_object_id(),
515                            referenced_object_type: PbObjectType::Secret as _,
516                        })
517                        .collect(),
518                }),
519            )
520            .await;
521
522        // notify default privileges for connections
523        if !updated_user_info.is_empty() {
524            version = self.notify_users_update(updated_user_info).await;
525        }
526
527        Ok(version)
528    }
529
530    pub async fn create_secret(
531        &self,
532        mut pb_secret: PbSecret,
533        secret_plain_payload: Vec<u8>,
534    ) -> MetaResult<NotificationVersion> {
535        let inner = self.inner.write().await;
536        let owner_id = pb_secret.owner as _;
537        let txn = inner.db.begin().await?;
538        ensure_user_id(owner_id, &txn).await?;
539        ensure_object_id(ObjectType::Database, pb_secret.database_id, &txn).await?;
540        ensure_object_id(ObjectType::Schema, pb_secret.schema_id, &txn).await?;
541        check_secret_name_duplicate(&pb_secret, &txn).await?;
542
543        let secret_obj = Self::create_object(
544            &txn,
545            ObjectType::Secret,
546            owner_id,
547            Some(pb_secret.schema_id.as_object_id()),
548        )
549        .await?;
550        pb_secret.id = secret_obj.oid.as_secret_id();
551        let secret: secret::ActiveModel = pb_secret.clone().into();
552        Secret::insert(secret).exec(&txn).await?;
553
554        let updated_user_info =
555            grant_default_privileges_automatically(&txn, secret_obj.oid).await?;
556
557        txn.commit().await?;
558
559        // Notify the compute and frontend node plain secret
560        let mut secret_plain = pb_secret;
561        secret_plain.value.clone_from(&secret_plain_payload);
562
563        LocalSecretManager::global().add_secret(secret_plain.id, secret_plain_payload);
564        self.env
565            .notification_manager()
566            .notify_compute_without_version(Operation::Add, Info::Secret(secret_plain.clone()));
567
568        let mut version = self
569            .notify_frontend(
570                NotificationOperation::Add,
571                NotificationInfo::Secret(secret_plain),
572            )
573            .await;
574
575        // notify default privileges for secrets
576        if !updated_user_info.is_empty() {
577            version = self.notify_users_update(updated_user_info).await;
578        }
579
580        Ok(version)
581    }
582
583    pub async fn create_view(
584        &self,
585        mut pb_view: PbView,
586        dependencies: HashSet<ObjectId>,
587    ) -> MetaResult<NotificationVersion> {
588        let inner = self.inner.write().await;
589        let owner_id = pb_view.owner as _;
590        let txn = inner.db.begin().await?;
591        ensure_user_id(owner_id, &txn).await?;
592        ensure_object_id(ObjectType::Database, pb_view.database_id, &txn).await?;
593        ensure_object_id(ObjectType::Schema, pb_view.schema_id, &txn).await?;
594        check_relation_name_duplicate(&pb_view.name, pb_view.database_id, pb_view.schema_id, &txn)
595            .await?;
596        ensure_object_id(ObjectType::Schema, pb_view.schema_id, &txn).await?;
597        check_relation_name_duplicate(&pb_view.name, pb_view.database_id, pb_view.schema_id, &txn)
598            .await?;
599
600        let view_obj = Self::create_object(
601            &txn,
602            ObjectType::View,
603            owner_id,
604            Some(pb_view.schema_id.as_object_id()),
605        )
606        .await?;
607        pb_view.id = view_obj.oid.as_view_id();
608        pb_view.created_at_epoch =
609            Some(Epoch::from_unix_millis(datetime_to_timestamp_millis(view_obj.created_at) as _).0);
610        pb_view.created_at_cluster_version = view_obj.created_at_cluster_version;
611
612        let view: view::ActiveModel = pb_view.clone().into();
613        View::insert(view).exec(&txn).await?;
614
615        let dependencies = if !dependencies.is_empty() {
616            ObjectDependency::insert_many(dependencies.into_iter().map(|obj_id| {
617                object_dependency::ActiveModel {
618                    oid: Set(obj_id),
619                    used_by: Set(view_obj.oid),
620                    ..Default::default()
621                }
622            }))
623            .exec(&txn)
624            .await?;
625            list_object_dependencies_by_object_id(&txn, view_obj.oid).await?
626        } else {
627            vec![]
628        };
629
630        let updated_user_info = grant_default_privileges_automatically(&txn, view_obj.oid).await?;
631
632        txn.commit().await?;
633        let mut version = self
634            .notify_frontend(
635                NotificationOperation::Add,
636                NotificationInfo::ObjectGroup(PbObjectGroup {
637                    objects: vec![PbObject {
638                        object_info: Some(PbObjectInfo::View(pb_view)),
639                    }],
640                    dependencies,
641                }),
642            )
643            .await;
644
645        // notify default privileges for views
646        if !updated_user_info.is_empty() {
647            version = self.notify_users_update(updated_user_info).await;
648        }
649
650        Ok(version)
651    }
652
653    pub async fn validate_cross_db_snapshot_backfill(
654        &self,
655        cross_db_snapshot_backfill_info: &SnapshotBackfillInfo,
656    ) -> MetaResult<()> {
657        if cross_db_snapshot_backfill_info
658            .upstream_mv_table_id_to_backfill_epoch
659            .is_empty()
660        {
661            return Ok(());
662        }
663
664        let inner = self.inner.read().await;
665        let table_ids = cross_db_snapshot_backfill_info
666            .upstream_mv_table_id_to_backfill_epoch
667            .keys()
668            .copied()
669            .map_into()
670            .collect_vec();
671        let cnt = Subscription::find()
672            .select_only()
673            .column(subscription::Column::DependentTableId)
674            .distinct()
675            .filter(subscription::Column::DependentTableId.is_in::<TableId, _>(table_ids))
676            .count(&inner.db)
677            .await? as usize;
678
679        if cnt
680            < cross_db_snapshot_backfill_info
681                .upstream_mv_table_id_to_backfill_epoch
682                .keys()
683                .count()
684        {
685            return Err(MetaError::permission_denied(
686                "Some upstream tables are not subscribed".to_owned(),
687            ));
688        }
689
690        Ok(())
691    }
692}