1use 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 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 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 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 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 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 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 if let Some(table_id) = iceberg_table_id {
314 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 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 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 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 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 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 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 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 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}