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
235        let source_obj = Self::create_object(
236            &txn,
237            ObjectType::Source,
238            owner_id,
239            iceberg_table_id
240                .map(TableId::as_object_id)
241                .or(Some(pb_source.schema_id.as_object_id())),
242        )
243        .await?;
244        let source_id = source_obj.oid.as_source_id();
245        pb_source.id = source_id;
246        let source: source::ActiveModel = pb_source.clone().into();
247        Source::insert(source).exec(&txn).await?;
248
249        // add secret and connection dependency
250        let dep_relation_ids = secret_ids
251            .iter()
252            .copied()
253            .map_into()
254            .chain(connection_ids.iter().copied().map_into());
255        let dependencies = if !secret_ids.is_empty() || !connection_ids.is_empty() {
256            ObjectDependency::insert_many(dep_relation_ids.map(|id| {
257                object_dependency::ActiveModel {
258                    oid: Set(id),
259                    used_by: Set(source_id.as_object_id()),
260                    ..Default::default()
261                }
262            }))
263            .exec(&txn)
264            .await?;
265            secret_ids
266                .iter()
267                .map(|id| PbObjectDependency {
268                    object_id: source_id.as_object_id(),
269                    referenced_object_id: id.as_object_id(),
270                    referenced_object_type: PbObjectType::Secret as _,
271                })
272                .chain(connection_ids.iter().map(|id| PbObjectDependency {
273                    object_id: source_id.as_object_id(),
274                    referenced_object_id: id.as_object_id(),
275                    referenced_object_type: PbObjectType::Connection as _,
276                }))
277                .collect_vec()
278        } else {
279            vec![]
280        };
281
282        let mut job_notifications = vec![];
283        let mut updated_user_info = vec![];
284        // check if it belongs to iceberg table
285        if let Some(table_id) = iceberg_table_id {
286            // 1. finish iceberg table job.
287            let table_notifications = self
288                .finish_streaming_job_inner(&txn, table_id.as_job_id())
289                .await?;
290            job_notifications.push((table_id.as_job_id(), table_notifications));
291
292            // 2. finish iceberg sink job.
293            let sink_id = Sink::find()
294                .select_only()
295                .column(sink::Column::SinkId)
296                .join(JoinType::InnerJoin, sink::Relation::Object.def())
297                .filter(
298                    object::Column::BelongToOid
299                        .eq(table_id)
300                        .and(object::Column::ObjType.eq(ObjectType::Sink)),
301                )
302                .into_tuple::<SinkId>()
303                .one(&txn)
304                .await?
305                .ok_or_else(|| {
306                    MetaError::catalog_id_not_found("iceberg sink for table", table_id)
307                })?;
308            let sink_job_id = sink_id.as_job_id();
309            let sink_notifications = self.finish_streaming_job_inner(&txn, sink_job_id).await?;
310            job_notifications.push((sink_job_id, sink_notifications));
311        } else {
312            updated_user_info = grant_default_privileges_automatically(&txn, source_id).await?;
313        }
314
315        txn.commit().await?;
316
317        for (job_id, (op, objects, user_info, dependencies)) in job_notifications {
318            let mut version = self
319                .notify_frontend(
320                    op,
321                    NotificationInfo::ObjectGroup(PbObjectGroup {
322                        objects,
323                        dependencies,
324                    }),
325                )
326                .await;
327            if !user_info.is_empty() {
328                version = self.notify_users_update(user_info).await;
329            }
330            inner
331                .creating_table_finish_notifier
332                .values_mut()
333                .for_each(|creating_tables| {
334                    if let Some(txs) = creating_tables.remove(&job_id) {
335                        for tx in txs {
336                            let _ = tx.send(Ok(version));
337                        }
338                    }
339                });
340        }
341
342        let mut version = self
343            .notify_frontend(
344                NotificationOperation::Add,
345                NotificationInfo::ObjectGroup(PbObjectGroup {
346                    objects: vec![PbObject {
347                        object_info: Some(PbObjectInfo::Source(pb_source)),
348                    }],
349                    dependencies,
350                }),
351            )
352            .await;
353
354        // notify default privileges for source
355        if !updated_user_info.is_empty() {
356            version = self.notify_users_update(updated_user_info).await;
357        }
358
359        Ok((source_id, version))
360    }
361
362    pub async fn create_function(
363        &self,
364        mut pb_function: PbFunction,
365    ) -> MetaResult<NotificationVersion> {
366        let inner = self.inner.write().await;
367        let owner_id = pb_function.owner as _;
368        let txn = inner.db.begin().await?;
369        ensure_user_id(owner_id, &txn).await?;
370        ensure_object_id(ObjectType::Database, pb_function.database_id, &txn).await?;
371        ensure_object_id(ObjectType::Schema, pb_function.schema_id, &txn).await?;
372        check_function_signature_duplicate(&pb_function, &txn).await?;
373
374        let function_obj = Self::create_object(
375            &txn,
376            ObjectType::Function,
377            owner_id,
378            Some(pb_function.schema_id.as_object_id()),
379        )
380        .await?;
381        pb_function.id = function_obj.oid.as_function_id();
382        pb_function.created_at_epoch = Some(
383            Epoch::from_unix_millis(datetime_to_timestamp_millis(function_obj.created_at) as _).0,
384        );
385        pb_function.created_at_cluster_version = function_obj.created_at_cluster_version;
386        let function: function::ActiveModel = pb_function.clone().into();
387        Function::insert(function).exec(&txn).await?;
388
389        let updated_user_info =
390            grant_default_privileges_automatically(&txn, function_obj.oid).await?;
391
392        txn.commit().await?;
393
394        let mut version = self
395            .notify_frontend(
396                NotificationOperation::Add,
397                NotificationInfo::Function(pb_function),
398            )
399            .await;
400
401        // notify default privileges for functions
402        if !updated_user_info.is_empty() {
403            version = self.notify_users_update(updated_user_info).await;
404        }
405
406        Ok(version)
407    }
408
409    pub async fn create_connection(
410        &self,
411        mut pb_connection: PbConnection,
412    ) -> MetaResult<NotificationVersion> {
413        let inner = self.inner.write().await;
414        let owner_id = pb_connection.owner as _;
415        let txn = inner.db.begin().await?;
416        ensure_user_id(owner_id, &txn).await?;
417        ensure_object_id(ObjectType::Database, pb_connection.database_id, &txn).await?;
418        ensure_object_id(ObjectType::Schema, pb_connection.schema_id, &txn).await?;
419        check_connection_name_duplicate(&pb_connection, &txn).await?;
420
421        let mut dep_secrets: HashSet<SecretId> = HashSet::new();
422        if let Some(ConnectionInfo::ConnectionParams(params)) = &pb_connection.info {
423            dep_secrets.extend(
424                params
425                    .secret_refs
426                    .values()
427                    .map(|secret_ref| secret_ref.secret_id),
428            );
429        }
430
431        let conn_obj = Self::create_object(
432            &txn,
433            ObjectType::Connection,
434            owner_id,
435            Some(pb_connection.schema_id.as_object_id()),
436        )
437        .await?;
438        pb_connection.id = conn_obj.oid.as_connection_id();
439        let connection: connection::ActiveModel = pb_connection.clone().into();
440        Connection::insert(connection).exec(&txn).await?;
441
442        for secret_id in &dep_secrets {
443            ObjectDependency::insert(object_dependency::ActiveModel {
444                oid: Set(secret_id.as_object_id()),
445                used_by: Set(conn_obj.oid),
446                ..Default::default()
447            })
448            .exec(&txn)
449            .await?;
450        }
451
452        let updated_user_info = grant_default_privileges_automatically(&txn, conn_obj.oid).await?;
453
454        txn.commit().await?;
455
456        {
457            // call meta telemetry here to report the connection creation
458            report_event(
459                PbTelemetryEventStage::Unspecified,
460                "connection_create",
461                pb_connection.get_id().as_raw_id() as i64,
462                {
463                    pb_connection.info.as_ref().and_then(|info| match info {
464                        ConnectionInfo::ConnectionParams(params) => {
465                            Some(params.connection_type().as_str_name().to_owned())
466                        }
467                        _ => None,
468                    })
469                },
470                None,
471                None,
472            );
473        }
474
475        let mut version = self
476            .notify_frontend(
477                NotificationOperation::Add,
478                NotificationInfo::ObjectGroup(PbObjectGroup {
479                    objects: vec![PbObject {
480                        object_info: Some(PbObjectInfo::Connection(pb_connection)),
481                    }],
482                    dependencies: dep_secrets
483                        .iter()
484                        .map(|secret_id| PbObjectDependency {
485                            object_id: conn_obj.oid,
486                            referenced_object_id: secret_id.as_object_id(),
487                            referenced_object_type: PbObjectType::Secret as _,
488                        })
489                        .collect(),
490                }),
491            )
492            .await;
493
494        // notify default privileges for connections
495        if !updated_user_info.is_empty() {
496            version = self.notify_users_update(updated_user_info).await;
497        }
498
499        Ok(version)
500    }
501
502    pub async fn create_secret(
503        &self,
504        mut pb_secret: PbSecret,
505        secret_plain_payload: Vec<u8>,
506    ) -> MetaResult<NotificationVersion> {
507        let inner = self.inner.write().await;
508        let owner_id = pb_secret.owner as _;
509        let txn = inner.db.begin().await?;
510        ensure_user_id(owner_id, &txn).await?;
511        ensure_object_id(ObjectType::Database, pb_secret.database_id, &txn).await?;
512        ensure_object_id(ObjectType::Schema, pb_secret.schema_id, &txn).await?;
513        check_secret_name_duplicate(&pb_secret, &txn).await?;
514
515        let secret_obj = Self::create_object(
516            &txn,
517            ObjectType::Secret,
518            owner_id,
519            Some(pb_secret.schema_id.as_object_id()),
520        )
521        .await?;
522        pb_secret.id = secret_obj.oid.as_secret_id();
523        let secret: secret::ActiveModel = pb_secret.clone().into();
524        Secret::insert(secret).exec(&txn).await?;
525
526        let updated_user_info =
527            grant_default_privileges_automatically(&txn, secret_obj.oid).await?;
528
529        txn.commit().await?;
530
531        // Notify the compute and frontend node plain secret
532        let mut secret_plain = pb_secret;
533        secret_plain.value.clone_from(&secret_plain_payload);
534
535        LocalSecretManager::global().add_secret(secret_plain.id, secret_plain_payload);
536        self.env
537            .notification_manager()
538            .notify_compute_without_version(Operation::Add, Info::Secret(secret_plain.clone()));
539
540        let mut version = self
541            .notify_frontend(
542                NotificationOperation::Add,
543                NotificationInfo::Secret(secret_plain),
544            )
545            .await;
546
547        // notify default privileges for secrets
548        if !updated_user_info.is_empty() {
549            version = self.notify_users_update(updated_user_info).await;
550        }
551
552        Ok(version)
553    }
554
555    pub async fn create_view(
556        &self,
557        mut pb_view: PbView,
558        dependencies: HashSet<ObjectId>,
559    ) -> MetaResult<NotificationVersion> {
560        let inner = self.inner.write().await;
561        let owner_id = pb_view.owner as _;
562        let txn = inner.db.begin().await?;
563        ensure_user_id(owner_id, &txn).await?;
564        ensure_object_id(ObjectType::Database, pb_view.database_id, &txn).await?;
565        ensure_object_id(ObjectType::Schema, pb_view.schema_id, &txn).await?;
566        check_relation_name_duplicate(&pb_view.name, pb_view.database_id, pb_view.schema_id, &txn)
567            .await?;
568        ensure_object_id(ObjectType::Schema, pb_view.schema_id, &txn).await?;
569        check_relation_name_duplicate(&pb_view.name, pb_view.database_id, pb_view.schema_id, &txn)
570            .await?;
571
572        let view_obj = Self::create_object(
573            &txn,
574            ObjectType::View,
575            owner_id,
576            Some(pb_view.schema_id.as_object_id()),
577        )
578        .await?;
579        pb_view.id = view_obj.oid.as_view_id();
580        pb_view.created_at_epoch =
581            Some(Epoch::from_unix_millis(datetime_to_timestamp_millis(view_obj.created_at) as _).0);
582        pb_view.created_at_cluster_version = view_obj.created_at_cluster_version;
583
584        let view: view::ActiveModel = pb_view.clone().into();
585        View::insert(view).exec(&txn).await?;
586
587        let dependencies = if !dependencies.is_empty() {
588            ObjectDependency::insert_many(dependencies.into_iter().map(|obj_id| {
589                object_dependency::ActiveModel {
590                    oid: Set(obj_id),
591                    used_by: Set(view_obj.oid),
592                    ..Default::default()
593                }
594            }))
595            .exec(&txn)
596            .await?;
597            list_object_dependencies_by_object_id(&txn, view_obj.oid).await?
598        } else {
599            vec![]
600        };
601
602        let updated_user_info = grant_default_privileges_automatically(&txn, view_obj.oid).await?;
603
604        txn.commit().await?;
605        let mut version = self
606            .notify_frontend(
607                NotificationOperation::Add,
608                NotificationInfo::ObjectGroup(PbObjectGroup {
609                    objects: vec![PbObject {
610                        object_info: Some(PbObjectInfo::View(pb_view)),
611                    }],
612                    dependencies,
613                }),
614            )
615            .await;
616
617        // notify default privileges for views
618        if !updated_user_info.is_empty() {
619            version = self.notify_users_update(updated_user_info).await;
620        }
621
622        Ok(version)
623    }
624
625    pub async fn validate_cross_db_snapshot_backfill(
626        &self,
627        cross_db_snapshot_backfill_info: &SnapshotBackfillInfo,
628    ) -> MetaResult<()> {
629        if cross_db_snapshot_backfill_info
630            .upstream_mv_table_id_to_backfill_epoch
631            .is_empty()
632        {
633            return Ok(());
634        }
635
636        let inner = self.inner.read().await;
637        let table_ids = cross_db_snapshot_backfill_info
638            .upstream_mv_table_id_to_backfill_epoch
639            .keys()
640            .copied()
641            .map_into()
642            .collect_vec();
643        let cnt = Subscription::find()
644            .select_only()
645            .column(subscription::Column::DependentTableId)
646            .distinct()
647            .filter(subscription::Column::DependentTableId.is_in::<TableId, _>(table_ids))
648            .count(&inner.db)
649            .await? as usize;
650
651        if cnt
652            < cross_db_snapshot_backfill_info
653                .upstream_mv_table_id_to_backfill_epoch
654                .keys()
655                .count()
656        {
657            return Err(MetaError::permission_denied(
658                "Some upstream tables are not subscribed".to_owned(),
659            ));
660        }
661
662        Ok(())
663    }
664}