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