1#[cfg(test)]
16mod tests {
17 use risingwave_common::catalog::{FragmentTypeFlag, FragmentTypeMask};
18 use risingwave_common::hash::VirtualNode;
19 use risingwave_meta_model::FragmentId;
20 use risingwave_meta_model::fragment::DistributionType;
21 use risingwave_meta_model::table::HandleConflictBehavior;
22 use risingwave_pb::catalog::subscription::SubscriptionState;
23 use risingwave_pb::catalog::{PbSinkType, StreamSourceInfo};
24 use risingwave_pb::common::{HostAddress, WorkerNode, WorkerType, worker_node};
25 use risingwave_pb::meta::SubscribeType;
26 use risingwave_pb::meta::table_fragments::fragment::PbFragmentDistributionType;
27 use risingwave_pb::stream_plan::stream_node::PbNodeBody;
28 use risingwave_pb::stream_plan::{PbStreamNode, StreamScanNode, StreamScanType};
29 use tokio::sync::{mpsc, oneshot};
30
31 use crate::barrier::Command;
32 use crate::controller::catalog::*;
33 use crate::manager::{LocalNotification, MetaOpts, WorkerKey};
34 use crate::model::{Fragment, FragmentDownstreamRelation};
35 use crate::serving::ServingVnodeMapping;
36
37 const TEST_DATABASE_ID: DatabaseId = DatabaseId::new(1);
38 const TEST_SCHEMA_ID: SchemaId = SchemaId::new(2);
39 const TEST_OWNER_ID: UserId = UserId::new(1);
40
41 async fn insert_test_table(
42 txn: &DatabaseTransaction,
43 table_id: TableId,
44 name: &str,
45 table_type: TableType,
46 belongs_to_job_id: Option<JobId>,
47 definition: &str,
48 ) -> MetaResult<()> {
49 table::ActiveModel {
50 table_id: Set(table_id),
51 name: Set(name.to_owned()),
52 optional_associated_source_id: Set(None),
53 table_type: Set(table_type),
54 belongs_to_job_id: Set(belongs_to_job_id),
55 columns: Set(vec![].into()),
56 pk: Set(vec![].into()),
57 distribution_key: Set(Vec::<i32>::new().into()),
58 stream_key: Set(Vec::<i32>::new().into()),
59 append_only: Set(false),
60 fragment_id: Set(None),
61 vnode_col_index: Set(None),
62 row_id_index: Set(None),
63 value_indices: Set(Vec::<i32>::new().into()),
64 definition: Set(definition.to_owned()),
65 handle_pk_conflict_behavior: Set(HandleConflictBehavior::NoCheck),
66 version_column_indices: Set(None),
67 read_prefix_len_hint: Set(0),
68 watermark_indices: Set(Vec::<i32>::new().into()),
69 dist_key_in_pk: Set(Vec::<i32>::new().into()),
70 dml_fragment_id: Set(None),
71 cardinality: Set(None),
72 cleaned_by_watermark: Set(false),
73 description: Set(None),
74 version: Set(None),
75 retention_seconds: Set(None),
76 cdc_table_id: Set(None),
77 vnode_count: Set(1),
78 webhook_info: Set(None),
79 engine: Set(None),
80 clean_watermark_index_in_pk: Set(None),
81 clean_watermark_indices: Set(None),
82 refreshable: Set(false),
83 vector_index_info: Set(None),
84 cdc_table_type: Set(None),
85 }
86 .insert(txn)
87 .await?;
88 Ok(())
89 }
90
91 async fn insert_test_fragment(
92 txn: &DatabaseTransaction,
93 fragment_id: FragmentId,
94 job_id: JobId,
95 state_table_ids: TableIdArray,
96 ) -> MetaResult<()> {
97 fragment::ActiveModel {
98 fragment_id: Set(fragment_id),
99 job_id: Set(job_id),
100 fragment_type_mask: Set(0),
101 distribution_type: Set(fragment::DistributionType::Hash),
102 stream_node: Set(StreamNode::from(&PbStreamNode::default())),
103 state_table_ids: Set(state_table_ids),
104 upstream_fragment_id: Set(I32Array::default()),
105 vnode_count: Set(1),
106 parallelism: Set(None),
107 }
108 .insert(txn)
109 .await?;
110 Ok(())
111 }
112
113 async fn insert_test_streaming_job(
114 txn: &DatabaseTransaction,
115 name: &str,
116 has_result_table: bool,
117 policy: Option<CacheRefillPolicy>,
118 ) -> MetaResult<(JobId, Option<TableId>, TableId)> {
119 let object_type = if has_result_table {
120 ObjectType::Table
121 } else {
122 ObjectType::Sink
123 };
124 let job_id = CatalogController::create_object(
125 txn,
126 object_type,
127 TEST_OWNER_ID,
128 Some(TEST_SCHEMA_ID.as_object_id()),
129 )
130 .await?
131 .oid
132 .as_job_id();
133 let result_table_id = has_result_table.then_some(job_id.as_mv_table_id());
134 if let Some(table_id) = result_table_id {
135 insert_test_table(txn, table_id, name, TableType::MaterializedView, None, "").await?;
136 }
137
138 let internal_table_id = CatalogController::create_object(
139 txn,
140 ObjectType::Table,
141 TEST_OWNER_ID,
142 Some(job_id.as_object_id()),
143 )
144 .await?
145 .oid
146 .as_table_id();
147 insert_test_table(
148 txn,
149 internal_table_id,
150 &format!("__internal_{name}"),
151 TableType::Internal,
152 Some(job_id),
153 "",
154 )
155 .await?;
156
157 insert_test_streaming_job_model(txn, job_id, policy).await?;
158
159 Ok((job_id, result_table_id, internal_table_id))
160 }
161
162 async fn insert_test_streaming_job_model(
163 txn: &DatabaseTransaction,
164 job_id: JobId,
165 policy: Option<CacheRefillPolicy>,
166 ) -> MetaResult<()> {
167 streaming_job::ActiveModel {
168 job_id: Set(job_id),
169 job_status: Set(JobStatus::Created),
170 create_type: Set(CreateType::Foreground),
171 timezone: Set(None),
172 config_override: Set(policy.map(|policy| {
173 format!(
174 "[streaming.developer]\ncache_refill_policy = \"{}\"\n",
175 policy
176 )
177 })),
178 adaptive_parallelism_strategy: Set(None),
179 parallelism: Set(StreamingParallelism::Adaptive),
180 backfill_parallelism: Set(None),
181 backfill_adaptive_parallelism_strategy: Set(None),
182 backfill_orders: Set(None),
183 max_parallelism: Set(1),
184 specific_resource_group: Set(None),
185 is_serverless_backfill: Set(false),
186 refresh_interval_sec: Set(None),
187 }
188 .insert(txn)
189 .await?;
190
191 Ok(())
192 }
193
194 #[tokio::test]
195 async fn test_cancel_creating_job_includes_belonging_streaming_jobs() -> MetaResult<()> {
196 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
197 let mut inner = mgr.inner.write().await;
198 let txn = inner.db.begin().await?;
199
200 let table_job_id = CatalogController::create_object(
201 &txn,
202 ObjectType::Table,
203 TEST_OWNER_ID,
204 Some(TEST_SCHEMA_ID.as_object_id()),
205 )
206 .await?
207 .oid
208 .as_job_id();
209 insert_test_table(
210 &txn,
211 table_job_id.as_mv_table_id(),
212 "cancel_table",
213 TableType::Table,
214 None,
215 "",
216 )
217 .await?;
218 let sink_job_id = CatalogController::create_object(
219 &txn,
220 ObjectType::Sink,
221 TEST_OWNER_ID,
222 Some(table_job_id.as_object_id()),
223 )
224 .await?
225 .oid
226 .as_job_id();
227 Sink::insert(sink::ActiveModel::from(PbSink {
228 id: sink_job_id.as_sink_id(),
229 schema_id: TEST_SCHEMA_ID,
230 database_id: TEST_DATABASE_ID,
231 name: "cancel_sink".to_owned(),
232 owner: TEST_OWNER_ID as _,
233 sink_type: PbSinkType::AppendOnly as i32,
234 ..Default::default()
235 }))
236 .exec(&txn)
237 .await?;
238 for job_id in [table_job_id, sink_job_id] {
239 insert_test_streaming_job_model(&txn, job_id, None).await?;
240 StreamingJob::update(streaming_job::ActiveModel {
241 job_id: Set(job_id),
242 job_status: Set(JobStatus::Creating),
243 ..Default::default()
244 })
245 .exec(&txn)
246 .await?;
247 }
248
249 let table_state_id = TableId::new(1000);
250 let sink_state_id = TableId::new(1001);
251 insert_test_fragment(
252 &txn,
253 FragmentId::new(100),
254 table_job_id,
255 TableIdArray(vec![table_state_id]),
256 )
257 .await?;
258 insert_test_fragment(
259 &txn,
260 FragmentId::new(101),
261 sink_job_id,
262 TableIdArray(vec![sink_state_id]),
263 )
264 .await?;
265 let (table_finish_tx, table_finish_rx) = oneshot::channel();
266 inner.register_finish_notifier(TEST_DATABASE_ID, table_job_id, table_finish_tx);
267 let (sink_finish_tx, sink_finish_rx) = oneshot::channel();
268 inner.register_finish_notifier(TEST_DATABASE_ID, sink_job_id, sink_finish_tx);
269 txn.commit().await?;
270 drop(inner);
271
272 let abort_result = mgr
273 .try_abort_creating_streaming_job(table_job_id, true)
274 .await?;
275 assert!(abort_result.aborted);
276 assert_eq!(
277 abort_result.aborted_sink_ids,
278 vec![sink_job_id.as_sink_id()]
279 );
280 let cancel_info = abort_result
281 .cancel_info
282 .expect("cancelled table job should have cleanup information");
283 assert_eq!(
284 cancel_info
285 .streaming_job_ids
286 .iter()
287 .copied()
288 .collect::<HashSet<_>>(),
289 HashSet::from([table_job_id, sink_job_id])
290 );
291 assert_eq!(
292 cancel_info
293 .state_table_ids
294 .iter()
295 .copied()
296 .collect::<HashSet<_>>(),
297 HashSet::from([table_state_id, sink_state_id])
298 );
299 let Command::DropStreamingJobs {
300 streaming_job_ids,
301 unregistered_state_table_ids,
302 ..
303 } = cancel_info.command
304 else {
305 unreachable!()
306 };
307 assert_eq!(
308 streaming_job_ids,
309 HashSet::from([table_job_id, sink_job_id])
310 );
311 assert_eq!(
312 unregistered_state_table_ids,
313 HashSet::from([table_state_id, sink_state_id])
314 );
315
316 for finish_rx in [table_finish_rx, sink_finish_rx] {
317 let err = finish_rx
318 .await
319 .expect("aborted job should notify its finish waiter")
320 .expect_err("aborted job should not finish successfully");
321 assert!(err.contains("cancelled"));
322 }
323 let db = &mgr.inner.read().await.db;
324 assert!(Object::find_by_id(table_job_id).one(db).await?.is_none());
325 assert!(Object::find_by_id(sink_job_id).one(db).await?.is_none());
326
327 Ok(())
328 }
329
330 #[tokio::test]
331 async fn test_create_multiple_sinks_into_same_table_and_drop_table() -> MetaResult<()> {
332 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
333 let inner = mgr.inner.write().await;
334 let txn = inner.db.begin().await?;
335 let (_, Some(target_table_id), _) =
336 insert_test_streaming_job(&txn, "mvt", true, None).await?
337 else {
338 unreachable!()
339 };
340 let (mv1_id, Some(_), _) = insert_test_streaming_job(&txn, "mv1", true, None).await? else {
341 unreachable!()
342 };
343 let (mv2_id, Some(_), _) = insert_test_streaming_job(&txn, "mv2", true, None).await? else {
344 unreachable!()
345 };
346 txn.commit().await?;
347 drop(inner);
348
349 let mut sink_ids = Vec::new();
350 let test_sink_tuples = [("s1", mv1_id), ("s2", mv2_id)];
351
352 fn assert_incoming_sink_drop_error<T>(error: &MetaError, test_sink_tuples: &[(&str, T)]) {
353 let message = error.to_string();
354
355 assert!(
356 message.contains("sink") && message.contains("depends on it"),
357 "expected an incoming-sink dependency error, got: {message}"
358 );
359
360 assert!(
361 test_sink_tuples
362 .iter()
363 .all(|(sink_name, _)| message.contains(*sink_name)),
364 "expected the error to mention all incoming sinks, got: {message}"
365 );
366 }
367
368 for (name, source) in test_sink_tuples {
369 let mut job = crate::manager::StreamingJob::Sink(
370 PbSink {
371 name: name.to_owned(),
372 database_id: TEST_DATABASE_ID,
373 schema_id: TEST_SCHEMA_ID,
374 owner: TEST_OWNER_ID as _,
375 target_table: Some(target_table_id),
376 sink_type: PbSinkType::AppendOnly as i32,
377 ..Default::default()
378 },
379 None,
380 );
381 tokio::time::timeout(
385 std::time::Duration::from_secs(3),
386 mgr.create_job_catalog(
387 &mut job,
388 &crate::model::StreamContext::default(),
389 &None,
390 1,
391 HashSet::from([source.as_object_id()]),
392 risingwave_pb::ddl_service::streaming_job_resource_type::ResourceType::Regular(
393 true,
394 ),
395 &None,
396 None,
397 None,
398 None,
399 None,
400 ),
401 )
402 .await
403 .expect("creating a second sink into the same table should not hang")?;
404 sink_ids.push(job.id().as_object_id());
405 }
406 let owned_sink_name = "owned_sink_without_iceberg_prefix";
407 let mut owned_sink = crate::manager::StreamingJob::Sink(
408 PbSink {
409 name: owned_sink_name.to_owned(),
410 database_id: TEST_DATABASE_ID,
411 schema_id: TEST_SCHEMA_ID,
412 owner: TEST_OWNER_ID as _,
413 target_table: Some(target_table_id),
414 sink_type: PbSinkType::AppendOnly as i32,
415 ..Default::default()
416 },
417 Some(target_table_id),
418 );
419 mgr.create_job_catalog(
420 &mut owned_sink,
421 &crate::model::StreamContext::default(),
422 &None,
423 1,
424 HashSet::from([mv1_id.as_object_id()]),
425 risingwave_pb::ddl_service::streaming_job_resource_type::ResourceType::Regular(true),
426 &None,
427 None,
428 None,
429 None,
430 None,
431 )
432 .await?;
433 let owned_sink_id = owned_sink.id().as_object_id();
434 sink_ids.push(owned_sink_id);
435
436 assert_eq!(sink_ids.len(), test_sink_tuples.len() + 1);
438
439 let inner = mgr.inner.read().await;
440 for sink_id in &sink_ids {
441 streaming_job::ActiveModel {
442 job_id: Set(sink_id.as_job_id()),
443 job_status: Set(JobStatus::Created),
444 ..Default::default()
445 }
446 .update(&inner.db)
447 .await?;
448 }
449 let object_dependency_count = ObjectDependency::find()
452 .filter(object_dependency::Column::Oid.eq(target_table_id.as_object_id()))
453 .filter(object_dependency::Column::UsedBy.is_in(sink_ids.clone()))
454 .count(&inner.db)
455 .await?;
456 assert_eq!(object_dependency_count, 0);
457 assert_eq!(
458 Object::find_by_id(owned_sink_id)
459 .one(&inner.db)
460 .await?
461 .unwrap()
462 .belong_to_oid,
463 Some(target_table_id.as_object_id())
464 );
465 drop(inner);
466
467 let error = mgr
468 .drop_object(ObjectType::Table, target_table_id, DropMode::Restrict)
469 .await
470 .expect_err("RESTRICT drop should fail for a table with incoming sinks");
471 assert_incoming_sink_drop_error(&error, &test_sink_tuples);
472 assert!(
473 !error.to_string().contains(owned_sink_name),
474 "an owned incoming sink should not prevent a RESTRICT drop"
475 );
476 mgr.drop_object(ObjectType::Table, target_table_id, DropMode::Cascade)
477 .await
478 .expect("CASCADE drop should succeed");
479
480 let inner = mgr.inner.read().await;
481 let db = &inner.db;
482 assert!(Object::find_by_id(target_table_id).one(db).await?.is_none());
484 assert_eq!(
485 Object::find()
486 .filter(object::Column::Oid.is_in(sink_ids))
487 .count(db)
488 .await?,
489 0
490 );
491 assert!(Object::find_by_id(mv1_id).one(db).await?.is_some());
493 assert!(Object::find_by_id(mv2_id).one(db).await?.is_some());
494
495 Ok(())
496 }
497
498 #[tokio::test]
499 async fn test_replace_upstream_object_rejects_creating_incoming_sink() -> MetaResult<()> {
500 fn assert_replace_concurrency_error(error: &MetaError) {
501 let message = error.to_string();
502 assert!(
504 message.contains("referenced by some creating jobs"),
505 "expected a replace concurrency error, got: {message}"
506 );
507 }
508
509 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
510 let inner = mgr.inner.write().await;
511 let txn = inner.db.begin().await?;
512 let (target_mv_id, Some(target_mv_table_id), _) =
513 insert_test_streaming_job(&txn, "target_mv", true, None).await?
514 else {
515 unreachable!()
516 };
517 let (source_mv_id, Some(_), _) =
518 insert_test_streaming_job(&txn, "source_mv", true, None).await?
519 else {
520 unreachable!()
521 };
522 txn.commit().await?;
523 drop(inner);
524
525 let mut sink = crate::manager::StreamingJob::Sink(
526 PbSink {
527 name: "creating_sink".to_owned(),
528 database_id: TEST_DATABASE_ID,
529 schema_id: TEST_SCHEMA_ID,
530 owner: TEST_OWNER_ID as _,
531 target_table: Some(target_mv_table_id),
532 sink_type: PbSinkType::AppendOnly as i32,
533 ..Default::default()
534 },
535 None,
536 );
537 let creating_sink = mgr
538 .create_job_catalog(
539 &mut sink,
540 &crate::model::StreamContext::default(),
541 &None,
542 1,
543 HashSet::from([source_mv_id.as_object_id()]),
544 risingwave_pb::ddl_service::streaming_job_resource_type::ResourceType::Regular(
545 true,
546 ),
547 &None,
548 None,
549 None,
550 None,
551 None,
552 )
553 .await?;
554 assert_ne!(creating_sink.job_status, JobStatus::Created);
556
557 let replacement = crate::manager::StreamingJob::MaterializedView(PbTable {
558 id: target_mv_table_id,
559 name: "target_mv".to_owned(),
560 database_id: TEST_DATABASE_ID,
561 schema_id: TEST_SCHEMA_ID,
562 owner: TEST_OWNER_ID as _,
563 ..Default::default()
564 });
565 assert_eq!(replacement.id(), target_mv_id);
567
568 let error = mgr
570 .create_job_catalog_for_replace(&replacement, None, None, None)
571 .await
572 .expect_err("replacement should reject a creating incoming sink");
573 assert_replace_concurrency_error(&error);
575
576 Ok(())
577 }
578
579 #[tokio::test]
580 async fn test_replace_upstream_object_with_created_incoming_sink() -> MetaResult<()> {
581 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
582 let inner = mgr.inner.write().await;
583 let txn = inner.db.begin().await?;
584 let (_, Some(target_table_id), _) =
585 insert_test_streaming_job(&txn, "target_table", true, None).await?
586 else {
587 unreachable!()
588 };
589 let (source_mv_id, Some(source_mv_table_id), _) =
590 insert_test_streaming_job(&txn, "source_mv", true, None).await?
591 else {
592 unreachable!()
593 };
594 txn.commit().await?;
595 drop(inner);
596
597 let mut sink = crate::manager::StreamingJob::Sink(
598 PbSink {
599 name: "created_sink".to_owned(),
600 database_id: TEST_DATABASE_ID,
601 schema_id: TEST_SCHEMA_ID,
602 owner: TEST_OWNER_ID as _,
603 target_table: Some(target_table_id),
604 sink_type: PbSinkType::AppendOnly as i32,
605 ..Default::default()
606 },
607 None,
608 );
609 let creating_sink = mgr
610 .create_job_catalog(
611 &mut sink,
612 &crate::model::StreamContext::default(),
613 &None,
614 1,
615 HashSet::from([source_mv_id.as_object_id()]),
616 risingwave_pb::ddl_service::streaming_job_resource_type::ResourceType::Regular(
617 true,
618 ),
619 &None,
620 None,
621 None,
622 None,
623 None,
624 )
625 .await?;
626 assert_ne!(creating_sink.job_status, JobStatus::Created);
628
629 let sink_id = sink.id();
630 let inner = mgr.inner.read().await;
631 streaming_job::ActiveModel {
632 job_id: Set(sink_id),
633 job_status: Set(JobStatus::Created),
634 ..Default::default()
635 }
636 .update(&inner.db)
637 .await?;
638 let sink_model = risingwave_meta_model::prelude::StreamingJob::find_by_id(sink_id)
639 .one(&inner.db)
640 .await?
641 .expect("sink should exist");
642 assert_eq!(sink_model.job_id, sink_id);
644 assert_eq!(sink_model.job_status, JobStatus::Created);
646 drop(inner);
647
648 let replacement = crate::manager::StreamingJob::MaterializedView(PbTable {
649 id: source_mv_table_id,
650 name: "source_mv".to_owned(),
651 database_id: TEST_DATABASE_ID,
652 schema_id: TEST_SCHEMA_ID,
653 owner: TEST_OWNER_ID as _,
654 ..Default::default()
655 });
656 assert_eq!(replacement.id(), source_mv_id);
658
659 let tmp_model = mgr
660 .create_job_catalog_for_replace(&replacement, None, None, None)
661 .await?;
662
663 assert_ne!(tmp_model.job_id, source_mv_id);
665 assert_eq!(tmp_model.job_status, JobStatus::Initial);
667
668 let inner = mgr.inner.read().await;
669 let db = &inner.db;
670 assert_eq!(
672 ObjectDependency::find()
673 .filter(object_dependency::Column::Oid.eq(source_mv_id.as_object_id()))
674 .filter(object_dependency::Column::UsedBy.eq(sink_id.as_object_id()))
675 .count(db)
676 .await?,
677 1
678 );
679 assert_eq!(
681 ObjectDependency::find()
682 .filter(object_dependency::Column::Oid.eq(source_mv_id.as_object_id()))
683 .filter(object_dependency::Column::UsedBy.eq(tmp_model.job_id.as_object_id()))
684 .count(db)
685 .await?,
686 1
687 );
688
689 Ok(())
690 }
691
692 #[tokio::test]
693 async fn test_table_refill_catalog_snapshot_classifies_table_identity() -> MetaResult<()> {
694 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
695 let inner = mgr.inner.write().await;
696 let txn = inner.db.begin().await?;
697
698 let (mv_job, Some(mv_result), mv_internal) =
699 insert_test_streaming_job(&txn, "mv_both", true, Some(CacheRefillPolicy::Both)).await?
700 else {
701 unreachable!()
702 };
703 let (default_job, Some(_default_result), default_internal) =
704 insert_test_streaming_job(&txn, "mv_default", true, None).await?
705 else {
706 unreachable!()
707 };
708 let (sink_job, None, sink_internal) = insert_test_streaming_job(
709 &txn,
710 "sink_streaming",
711 false,
712 Some(CacheRefillPolicy::Streaming),
713 )
714 .await?
715 else {
716 unreachable!()
717 };
718
719 let result_fragment = FragmentId::new(100);
720 let internal_fragment = FragmentId::new(101);
721 let sink_fragment = FragmentId::new(102);
722 for (fragment_id, job_id, table_ids) in [
723 (result_fragment, mv_job, vec![mv_result, mv_internal]),
724 (internal_fragment, default_job, vec![default_internal]),
725 (sink_fragment, sink_job, vec![sink_internal]),
726 ] {
727 insert_test_fragment(&txn, fragment_id, job_id, TableIdArray(table_ids)).await?;
728 }
729 txn.commit().await?;
730 drop(inner);
731
732 let serving_infos = mgr.fragment_serving_infos().await?;
733 assert_eq!(serving_infos.len(), 3);
734 assert_eq!(
735 serving_infos[&result_fragment].result_table_id,
736 Some(mv_result)
737 );
738 assert_eq!(serving_infos[&internal_fragment].result_table_id, None);
739 assert_eq!(serving_infos[&sink_fragment].result_table_id, None);
740
741 let policies = mgr.table_cache_refill_policies_snapshot().await?;
742 assert_eq!(
743 policies
744 .table_policies
745 .into_iter()
746 .map(|policy| (policy.table_id, policy.policy))
747 .collect::<HashMap<_, _>>(),
748 HashMap::from([(
749 mv_result.as_raw_id(),
750 CacheRefillPolicy::Both.to_protobuf() as i32,
751 )])
752 );
753 assert_eq!(
754 policies
755 .internal_table_policies
756 .into_iter()
757 .map(|policy| (policy.table_id, policy.policy))
758 .collect::<HashMap<_, _>>(),
759 HashMap::from([
760 (
761 mv_internal.as_raw_id(),
762 CacheRefillPolicy::Both.to_protobuf() as i32,
763 ),
764 (
765 sink_internal.as_raw_id(),
766 CacheRefillPolicy::Streaming.to_protobuf() as i32,
767 ),
768 ])
769 );
770
771 Ok(())
772 }
773
774 #[tokio::test]
775 async fn test_foreground_creating_catalog_lifecycle() -> MetaResult<()> {
776 let env = MetaSrvEnv::for_test().await;
777 let (tx, mut notification_rx) = mpsc::unbounded_channel();
778 env.notification_manager().insert_sender(
779 SubscribeType::Frontend,
780 WorkerKey(HostAddress {
781 host: "localhost".to_owned(),
782 port: 1234,
783 }),
784 tx,
785 );
786 let mgr = CatalogController::new(env).await?;
787 let inner = mgr.inner.write().await;
788 let txn = inner.db.begin().await?;
789
790 let (job_id, Some(table_id), internal_table_id) =
792 insert_test_streaming_job(&txn, "creating_table", true, None).await?
793 else {
794 unreachable!()
795 };
796 let associated_source_id = CatalogController::create_object(
797 &txn,
798 ObjectType::Source,
799 TEST_OWNER_ID,
800 Some(TEST_SCHEMA_ID.as_object_id()),
801 )
802 .await?
803 .oid
804 .as_source_id();
805 Source::insert(source::ActiveModel::from(PbSource {
806 id: associated_source_id,
807 schema_id: TEST_SCHEMA_ID,
808 database_id: TEST_DATABASE_ID,
809 name: "creating_table_source".to_owned(),
810 owner: TEST_OWNER_ID as _,
811 optional_associated_table_id: Some(
812 risingwave_pb::catalog::source::OptionalAssociatedTableId::AssociatedTableId(
813 table_id,
814 ),
815 ),
816 ..Default::default()
817 }))
818 .exec(&txn)
819 .await?;
820 table::ActiveModel {
821 table_id: Set(table_id),
822 table_type: Set(TableType::Table),
823 optional_associated_source_id: Set(Some(associated_source_id)),
824 ..Default::default()
825 }
826 .update(&txn)
827 .await?;
828 streaming_job::ActiveModel {
829 job_id: Set(job_id),
830 job_status: Set(JobStatus::Initial),
831 ..Default::default()
832 }
833 .update(&txn)
834 .await?;
835
836 let (_primary_job_id, Some(primary_table_id), _) =
838 insert_test_streaming_job(&txn, "primary_table", true, None).await?
839 else {
840 unreachable!()
841 };
842 let index_job_id = CatalogController::create_object(
843 &txn,
844 ObjectType::Index,
845 TEST_OWNER_ID,
846 Some(TEST_SCHEMA_ID.as_object_id()),
847 )
848 .await?
849 .oid
850 .as_job_id();
851 let index_table_id = index_job_id.as_mv_table_id();
852 insert_test_table(
853 &txn,
854 index_table_id,
855 "creating_index",
856 TableType::Index,
857 None,
858 "",
859 )
860 .await?;
861 index::ActiveModel {
862 index_id: Set(index_job_id.as_index_id()),
863 name: Set("creating_index".to_owned()),
864 index_table_id: Set(index_table_id),
865 primary_table_id: Set(primary_table_id),
866 index_items: Set(vec![].into()),
867 index_column_properties: Set(None),
868 index_columns_len: Set(0),
869 }
870 .insert(&txn)
871 .await?;
872 insert_test_streaming_job_model(&txn, index_job_id, None).await?;
873 streaming_job::ActiveModel {
874 job_id: Set(index_job_id),
875 job_status: Set(JobStatus::Initial),
876 ..Default::default()
877 }
878 .update(&txn)
879 .await?;
880
881 let source_job_id = CatalogController::create_object(
883 &txn,
884 ObjectType::Source,
885 TEST_OWNER_ID,
886 Some(TEST_SCHEMA_ID.as_object_id()),
887 )
888 .await?
889 .oid
890 .as_job_id();
891 Source::insert(source::ActiveModel::from(PbSource {
892 id: source_job_id.as_shared_source_id(),
893 schema_id: TEST_SCHEMA_ID,
894 database_id: TEST_DATABASE_ID,
895 name: "creating_shared_source".to_owned(),
896 owner: TEST_OWNER_ID as _,
897 info: Some(StreamSourceInfo {
898 cdc_source_job: true,
899 ..Default::default()
900 }),
901 ..Default::default()
902 }))
903 .exec(&txn)
904 .await?;
905 insert_test_streaming_job_model(&txn, source_job_id, None).await?;
906 streaming_job::ActiveModel {
907 job_id: Set(source_job_id),
908 job_status: Set(JobStatus::Initial),
909 ..Default::default()
910 }
911 .update(&txn)
912 .await?;
913
914 let sink_job_id = CatalogController::create_object(
915 &txn,
916 ObjectType::Sink,
917 TEST_OWNER_ID,
918 Some(TEST_SCHEMA_ID.as_object_id()),
919 )
920 .await?
921 .oid
922 .as_job_id();
923 Sink::insert(sink::ActiveModel::from(PbSink {
924 id: sink_job_id.as_sink_id(),
925 schema_id: TEST_SCHEMA_ID,
926 database_id: TEST_DATABASE_ID,
927 name: "creating_sink".to_owned(),
928 owner: TEST_OWNER_ID as _,
929 sink_type: PbSinkType::AppendOnly as i32,
930 ..Default::default()
931 }))
932 .exec(&txn)
933 .await?;
934 insert_test_streaming_job_model(&txn, sink_job_id, None).await?;
935 streaming_job::ActiveModel {
936 job_id: Set(sink_job_id),
937 job_status: Set(JobStatus::Initial),
938 ..Default::default()
939 }
940 .update(&txn)
941 .await?;
942
943 txn.commit().await?;
944
945 let (catalog, _) = inner.snapshot().await?;
946 assert!(!catalog.2.iter().any(|table| table.id == table_id));
947 assert!(!catalog.2.iter().any(|table| table.id == internal_table_id));
948 assert!(!catalog.2.iter().any(|table| table.id == index_table_id));
949 assert!(
950 !catalog
951 .3
952 .iter()
953 .any(|source| source.id == associated_source_id)
954 );
955 assert!(
956 !catalog
957 .3
958 .iter()
959 .any(|source| source.id == source_job_id.as_shared_source_id())
960 );
961 assert!(
962 !catalog
963 .6
964 .iter()
965 .any(|index| index.id == index_job_id.as_index_id())
966 );
967 assert!(
968 !catalog
969 .4
970 .iter()
971 .any(|sink| sink.id == sink_job_id.as_sink_id())
972 );
973
974 drop(inner);
975
976 let downstreams = FragmentDownstreamRelation::new();
977 let mut add_notifications = vec![];
978 for creating_job_id in [job_id, index_job_id, source_job_id, sink_job_id] {
979 mgr.post_collect_job_fragments(creating_job_id, &downstreams, None, None, None, true)
980 .await?;
981 let notification = notification_rx
982 .recv()
983 .await
984 .expect("frontend should receive a creating notification")
985 .expect("creating notification should be valid");
986 assert_eq!(notification.operation(), NotificationOperation::Add);
987 let object_group = match notification.info {
988 Some(NotificationInfo::ObjectGroup(object_group)) => object_group,
989 other => panic!("unexpected notification: {other:?}"),
990 };
991 add_notifications.push(object_group);
992 }
993
994 assert!(add_notifications[0].objects.iter().any(|object| matches!(
995 &object.object_info,
996 Some(PbObjectInfo::Table(table)) if table.id == table_id
997 )));
998 assert!(add_notifications[0].objects.iter().any(|object| matches!(
999 &object.object_info,
1000 Some(PbObjectInfo::Table(table)) if table.id == internal_table_id
1001 )));
1002 assert!(add_notifications[0].objects.iter().any(|object| matches!(
1003 &object.object_info,
1004 Some(PbObjectInfo::Source(source)) if source.id == associated_source_id
1005 )));
1006 assert!(add_notifications[1].objects.iter().any(|object| matches!(
1007 &object.object_info,
1008 Some(PbObjectInfo::Table(table)) if table.id == index_table_id
1009 )));
1010 assert!(add_notifications[1].objects.iter().any(|object| matches!(
1011 &object.object_info,
1012 Some(PbObjectInfo::Index(index)) if index.id == index_job_id.as_index_id()
1013 )));
1014 assert!(add_notifications[2].objects.iter().any(|object| matches!(
1015 &object.object_info,
1016 Some(PbObjectInfo::Source(source)) if source.id == source_job_id.as_shared_source_id()
1017 )));
1018 assert!(add_notifications[3].objects.iter().any(|object| matches!(
1019 &object.object_info,
1020 Some(PbObjectInfo::Sink(sink)) if sink.id == sink_job_id.as_sink_id()
1021 )));
1022
1023 let inner = mgr.inner.write().await;
1024 let (catalog, _) = inner.snapshot().await?;
1025 assert!(catalog.2.iter().any(|table| table.id == table_id));
1026 assert!(catalog.2.iter().any(|table| table.id == internal_table_id));
1027 assert!(catalog.2.iter().any(|table| table.id == index_table_id));
1028 assert!(
1029 catalog
1030 .3
1031 .iter()
1032 .any(|source| source.id == source_job_id.as_shared_source_id())
1033 );
1034 assert!(
1035 catalog
1036 .6
1037 .iter()
1038 .any(|index| index.id == index_job_id.as_index_id())
1039 );
1040 assert!(
1041 catalog
1042 .4
1043 .iter()
1044 .any(|sink| sink.id == sink_job_id.as_sink_id())
1045 );
1046
1047 let txn = inner.db.begin().await?;
1048 let (operation, _, _, _) = mgr.finish_streaming_job_inner(&txn, job_id).await?;
1049 assert_eq!(operation, NotificationOperation::Update);
1050 let (operation, _, _, _) = mgr.finish_streaming_job_inner(&txn, index_job_id).await?;
1051 assert_eq!(operation, NotificationOperation::Update);
1052 let (operation, _, _, _) = mgr.finish_streaming_job_inner(&txn, source_job_id).await?;
1053 assert_eq!(operation, NotificationOperation::Update);
1054 let (operation, _, _, _) = mgr.finish_streaming_job_inner(&txn, sink_job_id).await?;
1055 assert_eq!(operation, NotificationOperation::Update);
1056 txn.commit().await?;
1057
1058 Ok(())
1059 }
1060
1061 #[tokio::test]
1062 async fn test_alter_streaming_job_cache_refill_policy_notifies_hummock() -> MetaResult<()> {
1063 let env = MetaSrvEnv::for_test().await;
1064 let (tx, mut rx) = mpsc::unbounded_channel();
1065 env.notification_manager().insert_sender(
1066 SubscribeType::Hummock,
1067 WorkerKey(HostAddress {
1068 host: "localhost".to_owned(),
1069 port: 1234,
1070 }),
1071 tx,
1072 );
1073 let mgr = CatalogController::new(env).await?;
1074
1075 let inner = mgr.inner.write().await;
1076 let txn = inner.db.begin().await?;
1077 let (_job, Some(result_table_id), internal_table_id) =
1078 insert_test_streaming_job(&txn, "mv_cache_refill", true, None).await?
1079 else {
1080 unreachable!()
1081 };
1082 insert_test_fragment(
1083 &txn,
1084 FragmentId::new(200),
1085 result_table_id.as_job_id(),
1086 TableIdArray(vec![result_table_id, internal_table_id]),
1087 )
1088 .await?;
1089 txn.commit().await?;
1090 drop(inner);
1091
1092 mgr.alter_streaming_job_config(
1093 result_table_id.as_job_id(),
1094 HashMap::from([(
1095 STREAMING_CACHE_REFILL_POLICY_CONFIG_PATH.to_owned(),
1096 "\"both\"".to_owned(),
1097 )]),
1098 vec![],
1099 )
1100 .await?;
1101
1102 let response = rx
1103 .recv()
1104 .await
1105 .expect("should receive hummock notification")
1106 .expect("notification should be ok");
1107 assert_eq!(response.operation(), NotificationOperation::Update);
1108 let info = response.info;
1109 let Some(NotificationInfo::TableRefillRuntimeConfig(config)) = info else {
1110 panic!("unexpected notification: {:?}", info);
1111 };
1112 assert!(config.serving_table_vnode_mappings.is_none());
1113 let policies = config
1114 .table_cache_refill_policies
1115 .expect("policy snapshot should be present");
1116 assert_eq!(
1117 policies
1118 .table_policies
1119 .into_iter()
1120 .map(|policy| (policy.table_id, policy.policy))
1121 .collect::<HashMap<_, _>>(),
1122 HashMap::from([(
1123 result_table_id.as_raw_id(),
1124 CacheRefillPolicy::Both.to_protobuf() as i32,
1125 )])
1126 );
1127 assert_eq!(
1128 policies
1129 .internal_table_policies
1130 .into_iter()
1131 .map(|policy| (policy.table_id, policy.policy))
1132 .collect::<HashMap<_, _>>(),
1133 HashMap::from([(
1134 internal_table_id.as_raw_id(),
1135 CacheRefillPolicy::Both.to_protobuf() as i32,
1136 )])
1137 );
1138
1139 Ok(())
1140 }
1141
1142 #[tokio::test]
1143 async fn test_prepare_streaming_job_cache_refill_policy_notifies_hummock() -> MetaResult<()> {
1144 let env = MetaSrvEnv::for_test().await;
1145 let (tx, mut rx) = mpsc::unbounded_channel();
1146 env.notification_manager().insert_sender(
1147 SubscribeType::Hummock,
1148 WorkerKey(HostAddress {
1149 host: "localhost".to_owned(),
1150 port: 1234,
1151 }),
1152 tx,
1153 );
1154 let (local_notification_tx, mut local_notification_rx) = mpsc::unbounded_channel();
1155 env.notification_manager()
1156 .insert_local_sender(local_notification_tx);
1157 let mgr = CatalogController::new(env).await?;
1158
1159 let inner = mgr.inner.write().await;
1160 let txn = inner.db.begin().await?;
1161 let (job_id, Some(result_table_id), internal_table_id) = insert_test_streaming_job(
1162 &txn,
1163 "mv_initial_cache_refill",
1164 true,
1165 Some(CacheRefillPolicy::Both),
1166 )
1167 .await?
1168 else {
1169 unreachable!()
1170 };
1171 let (unprepared_job_id, Some(unprepared_result_table_id), unprepared_internal_table_id) =
1172 insert_test_streaming_job(
1173 &txn,
1174 "mv_unprepared_cache_refill",
1175 true,
1176 Some(CacheRefillPolicy::Serving),
1177 )
1178 .await?
1179 else {
1180 unreachable!()
1181 };
1182 for job_id in [job_id, unprepared_job_id] {
1183 streaming_job::ActiveModel {
1184 job_id: Set(job_id),
1185 job_status: Set(JobStatus::Initial),
1186 ..Default::default()
1187 }
1188 .update(&txn)
1189 .await?;
1190 }
1191 txn.commit().await?;
1192 drop(inner);
1193
1194 let fragments = [Fragment {
1195 fragment_id: FragmentId::new(300),
1196 fragment_type_mask: FragmentTypeMask::default(),
1197 distribution_type: PbFragmentDistributionType::Hash,
1198 state_table_ids: vec![],
1199 maybe_vnode_count: Some(1),
1200 nodes: PbStreamNode::default(),
1201 }];
1202 mgr.prepare_streaming_job(
1203 job_id,
1204 || fragments.iter(),
1205 &FragmentDownstreamRelation::default(),
1206 true,
1207 None,
1208 None,
1209 )
1210 .await?;
1211
1212 let local_notification = local_notification_rx
1213 .try_recv()
1214 .expect("should receive serving fragment mapping notification");
1215 let LocalNotification::ServingFragmentMappingsUpsert(fragment_ids) = local_notification
1216 else {
1217 panic!(
1218 "unexpected local notification before hummock notification: {:?}",
1219 local_notification
1220 );
1221 };
1222 assert_eq!(fragment_ids, vec![FragmentId::new(300).as_raw_id()]);
1223
1224 let response = rx
1225 .recv()
1226 .await
1227 .expect("should receive hummock notification")
1228 .expect("notification should be ok");
1229 assert_eq!(response.operation(), NotificationOperation::Update);
1230 let info = response.info;
1231 let Some(NotificationInfo::TableRefillRuntimeConfig(config)) = info else {
1232 panic!("unexpected notification: {:?}", info);
1233 };
1234 assert!(config.serving_table_vnode_mappings.is_none());
1235 let policies = config
1236 .table_cache_refill_policies
1237 .expect("policy snapshot should be present");
1238 let table_policies = policies
1239 .table_policies
1240 .into_iter()
1241 .map(|policy| (policy.table_id, policy.policy))
1242 .collect::<HashMap<_, _>>();
1243 assert_eq!(
1244 table_policies,
1245 HashMap::from([(
1246 result_table_id.as_raw_id(),
1247 CacheRefillPolicy::Both.to_protobuf() as i32,
1248 )])
1249 );
1250 assert!(!table_policies.contains_key(&unprepared_result_table_id.as_raw_id()));
1251 let internal_table_policies = policies
1252 .internal_table_policies
1253 .into_iter()
1254 .map(|policy| (policy.table_id, policy.policy))
1255 .collect::<HashMap<_, _>>();
1256 assert_eq!(
1257 internal_table_policies,
1258 HashMap::from([(
1259 internal_table_id.as_raw_id(),
1260 CacheRefillPolicy::Both.to_protobuf() as i32,
1261 )])
1262 );
1263 assert!(!internal_table_policies.contains_key(&unprepared_internal_table_id.as_raw_id()));
1264
1265 Ok(())
1266 }
1267
1268 async fn insert_dirty_creating_job_with_fragment(
1269 mgr: &CatalogController,
1270 fragment_id: FragmentId,
1271 vnode_count: i32,
1272 fragment_type_mask: FragmentTypeMask,
1273 ) -> MetaResult<(JobId, TableId)> {
1274 let inner = mgr.inner.write().await;
1275 let txn = inner.db.begin().await?;
1276 let job_obj = CatalogController::create_object(
1277 &txn,
1278 ObjectType::Table,
1279 TEST_OWNER_ID,
1280 Some(TEST_SCHEMA_ID.as_object_id()),
1281 )
1282 .await?;
1283 let job_id = job_obj.oid.as_job_id();
1284 let table_id = job_id.as_mv_table_id();
1285 insert_test_table(
1286 &txn,
1287 table_id,
1288 "mv_dirty_serving_mapping",
1289 TableType::MaterializedView,
1290 None,
1291 "CREATE MATERIALIZED VIEW mv_dirty_serving_mapping AS SELECT 1",
1292 )
1293 .await?;
1294 table::ActiveModel {
1295 table_id: Set(table_id),
1296 engine: Set(Some(table::Engine::Hummock)),
1297 ..Default::default()
1298 }
1299 .update(&txn)
1300 .await?;
1301 streaming_job::ActiveModel {
1302 job_id: Set(job_id),
1303 job_status: Set(JobStatus::Creating),
1304 create_type: Set(CreateType::Foreground),
1305 timezone: Set(None),
1306 config_override: Set(None),
1307 adaptive_parallelism_strategy: Set(None),
1308 parallelism: Set(StreamingParallelism::Adaptive),
1309 backfill_parallelism: Set(None),
1310 backfill_adaptive_parallelism_strategy: Set(None),
1311 backfill_orders: Set(None),
1312 max_parallelism: Set(1),
1313 specific_resource_group: Set(None),
1314 is_serverless_backfill: Set(false),
1315 refresh_interval_sec: Set(None),
1316 }
1317 .insert(&txn)
1318 .await?;
1319 fragment::ActiveModel {
1320 fragment_id: Set(fragment_id),
1321 job_id: Set(job_id),
1322 fragment_type_mask: Set(fragment_type_mask.into()),
1323 distribution_type: Set(DistributionType::Hash),
1324 stream_node: Set(StreamNode::default()),
1325 state_table_ids: Set(Vec::<TableId>::new().into()),
1326 upstream_fragment_id: Set(Vec::<i32>::new().into()),
1327 vnode_count: Set(vnode_count),
1328 parallelism: Set(None),
1329 }
1330 .insert(&txn)
1331 .await?;
1332 txn.commit().await?;
1333 drop(inner);
1334
1335 Ok((job_id, table_id))
1336 }
1337
1338 #[tokio::test]
1339 async fn test_dirty_cleanup_reconcile_removes_stale_serving_vnode_mapping() -> MetaResult<()> {
1340 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1341 let fragment_id = FragmentId::new(42);
1342 insert_dirty_creating_job_with_fragment(
1343 &mgr,
1344 fragment_id,
1345 VirtualNode::COUNT_FOR_TEST as i32,
1346 FragmentTypeMask::from(FragmentTypeFlag::Values as u32),
1347 )
1348 .await?;
1349
1350 let worker = WorkerNode {
1351 id: WorkerId::new(1),
1352 r#type: WorkerType::ComputeNode.into(),
1353 host: Some(HostAddress {
1354 host: "localhost".to_owned(),
1355 port: 1,
1356 }),
1357 state: worker_node::State::Running as i32,
1358 property: Some(worker_node::Property {
1359 is_serving: true,
1360 parallelism: 1,
1361 ..Default::default()
1362 }),
1363 ..Default::default()
1364 };
1365 let serving_vnode_mapping = ServingVnodeMapping::default();
1366 let initial_snapshot = mgr.fragment_serving_infos().await?;
1367 serving_vnode_mapping.upsert(&initial_snapshot, std::slice::from_ref(&worker), None);
1368 assert!(serving_vnode_mapping.all().contains_key(&fragment_id));
1369
1370 mgr.clean_dirty_creating_jobs(Some(TEST_DATABASE_ID))
1371 .await?;
1372 let current_snapshot = mgr.fragment_serving_infos().await?;
1373 assert!(!current_snapshot.contains_key(&fragment_id));
1374
1375 serving_vnode_mapping.reconcile(¤t_snapshot, &[worker], None);
1376 assert!(!serving_vnode_mapping.all().contains_key(&fragment_id));
1377
1378 Ok(())
1379 }
1380
1381 #[tokio::test]
1382 async fn test_clean_dirty_creating_jobs_keeps_job_without_values_fragment() -> MetaResult<()> {
1383 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1384 let fragment_id = FragmentId::new(43);
1385 let (job_id, table_id) = insert_dirty_creating_job_with_fragment(
1386 &mgr,
1387 fragment_id,
1388 1,
1389 FragmentTypeMask::empty(),
1390 )
1391 .await?;
1392
1393 let cleaned = mgr
1394 .clean_dirty_creating_jobs(Some(TEST_DATABASE_ID))
1395 .await?;
1396 assert!(cleaned.streaming_job_ids.is_empty());
1397
1398 let inner = mgr.inner.read().await;
1399 assert!(Object::find_by_id(job_id).one(&inner.db).await?.is_some());
1400 assert!(
1401 StreamingJob::find_by_id(job_id)
1402 .one(&inner.db)
1403 .await?
1404 .is_some()
1405 );
1406 assert!(Table::find_by_id(table_id).one(&inner.db).await?.is_some());
1407
1408 Ok(())
1409 }
1410
1411 #[tokio::test]
1412 async fn test_clean_dirty_creating_jobs_cleans_foreground_job_in_legacy_mode() -> MetaResult<()>
1413 {
1414 let mut opts = MetaOpts::test(false);
1415 opts.clean_all_foreground_jobs_on_recovery = true;
1416 let mgr = CatalogController::new(MetaSrvEnv::for_test_opts(opts, |_| ()).await).await?;
1417 let (job_id, table_id) = insert_dirty_creating_job_with_fragment(
1418 &mgr,
1419 FragmentId::new(44),
1420 1,
1421 FragmentTypeMask::empty(),
1422 )
1423 .await?;
1424
1425 let cleaned = mgr
1426 .clean_dirty_creating_jobs(Some(TEST_DATABASE_ID))
1427 .await?;
1428 assert_eq!(cleaned.streaming_job_ids, vec![job_id]);
1429
1430 let db = &mgr.inner.read().await.db;
1431 assert!(Object::find_by_id(job_id).one(db).await?.is_none());
1432 assert!(StreamingJob::find_by_id(job_id).one(db).await?.is_none());
1433 assert!(Table::find_by_id(table_id).one(db).await?.is_none());
1434
1435 Ok(())
1436 }
1437
1438 #[tokio::test]
1439 async fn test_database_func() -> MetaResult<()> {
1440 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1441 let pb_database = PbDatabase {
1442 name: "db1".to_owned(),
1443 owner: TEST_OWNER_ID as _,
1444 ..Default::default()
1445 };
1446 mgr.create_database(pb_database).await?;
1447
1448 let database_id: DatabaseId = Database::find()
1449 .select_only()
1450 .column(database::Column::DatabaseId)
1451 .filter(database::Column::Name.eq("db1"))
1452 .into_tuple()
1453 .one(&mgr.inner.read().await.db)
1454 .await?
1455 .unwrap();
1456
1457 mgr.alter_name(ObjectType::Database, database_id, "db2")
1458 .await?;
1459 let database = Database::find_by_id(database_id)
1460 .one(&mgr.inner.read().await.db)
1461 .await?
1462 .unwrap();
1463 assert_eq!(database.name, "db2");
1464
1465 let schema_id: SchemaId = Schema::find()
1466 .inner_join(Object)
1467 .select_only()
1468 .column(schema::Column::SchemaId)
1469 .filter(object::Column::DatabaseId.eq(database_id))
1470 .into_tuple()
1471 .one(&mgr.inner.read().await.db)
1472 .await?
1473 .unwrap();
1474 mgr.create_view(
1475 PbView {
1476 schema_id,
1477 database_id,
1478 name: "cross_db_upstream".to_owned(),
1479 owner: TEST_OWNER_ID as _,
1480 sql: "CREATE VIEW cross_db_upstream AS SELECT 1".to_owned(),
1481 ..Default::default()
1482 },
1483 HashSet::new(),
1484 )
1485 .await?;
1486 let upstream_id: ViewId = View::find()
1487 .inner_join(Object)
1488 .select_only()
1489 .column(view::Column::ViewId)
1490 .filter(
1491 object::Column::DatabaseId
1492 .eq(database_id)
1493 .and(view::Column::Name.eq("cross_db_upstream")),
1494 )
1495 .into_tuple()
1496 .one(&mgr.inner.read().await.db)
1497 .await?
1498 .unwrap();
1499
1500 let inner = mgr.inner.write().await;
1501 let txn = inner.db.begin().await?;
1502 let (dependent_job_id, Some(dependent_table_id), _) =
1503 insert_test_streaming_job(&txn, "cross_db_dependent", true, None).await?
1504 else {
1505 unreachable!()
1506 };
1507 ObjectDependency::insert(object_dependency::ActiveModel {
1508 oid: Set(upstream_id.as_object_id()),
1509 used_by: Set(dependent_job_id.as_object_id()),
1510 ..Default::default()
1511 })
1512 .exec(&txn)
1513 .await?;
1514 txn.commit().await?;
1515 drop(inner);
1516
1517 assert!(
1518 mgr.drop_object(ObjectType::Database, database_id, DropMode::Cascade)
1519 .await
1520 .is_err()
1521 );
1522 mgr.drop_object(ObjectType::Table, dependent_table_id, DropMode::Cascade)
1523 .await?;
1524 mgr.drop_object(ObjectType::Database, database_id, DropMode::Cascade)
1525 .await?;
1526
1527 Ok(())
1528 }
1529
1530 #[tokio::test]
1531 async fn test_schema_func() -> MetaResult<()> {
1532 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1533 let pb_schema = PbSchema {
1534 database_id: TEST_DATABASE_ID,
1535 name: "schema1".to_owned(),
1536 owner: TEST_OWNER_ID as _,
1537 ..Default::default()
1538 };
1539 mgr.create_schema(pb_schema.clone()).await?;
1540 assert!(mgr.create_schema(pb_schema).await.is_err());
1541
1542 let schema_id: SchemaId = Schema::find()
1543 .select_only()
1544 .column(schema::Column::SchemaId)
1545 .filter(schema::Column::Name.eq("schema1"))
1546 .into_tuple()
1547 .one(&mgr.inner.read().await.db)
1548 .await?
1549 .unwrap();
1550
1551 mgr.alter_name(ObjectType::Schema, schema_id, "schema2")
1552 .await?;
1553 let schema = Schema::find_by_id(schema_id)
1554 .one(&mgr.inner.read().await.db)
1555 .await?
1556 .unwrap();
1557 assert_eq!(schema.name, "schema2");
1558 mgr.drop_object(ObjectType::Schema, schema_id, DropMode::Restrict)
1559 .await?;
1560
1561 Ok(())
1562 }
1563
1564 #[tokio::test]
1565 async fn test_create_view() -> MetaResult<()> {
1566 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1567 let pb_view = PbView {
1568 schema_id: TEST_SCHEMA_ID,
1569 database_id: TEST_DATABASE_ID,
1570 name: "view".to_owned(),
1571 owner: TEST_OWNER_ID as _,
1572 sql: "CREATE VIEW view AS SELECT 1".to_owned(),
1573 ..Default::default()
1574 };
1575 mgr.create_view(pb_view.clone(), HashSet::new()).await?;
1576 assert!(mgr.create_view(pb_view, HashSet::new()).await.is_err());
1577
1578 let view = View::find().one(&mgr.inner.read().await.db).await?.unwrap();
1579 mgr.drop_object(ObjectType::View, view.view_id, DropMode::Cascade)
1580 .await?;
1581 assert!(
1582 View::find_by_id(view.view_id)
1583 .one(&mgr.inner.read().await.db)
1584 .await?
1585 .is_none()
1586 );
1587
1588 Ok(())
1589 }
1590
1591 #[tokio::test]
1592 async fn test_object_belong_to_cascade() -> MetaResult<()> {
1593 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1594 mgr.create_schema(PbSchema {
1595 database_id: TEST_DATABASE_ID,
1596 name: "belong_to_target".to_owned(),
1597 owner: TEST_OWNER_ID as _,
1598 ..Default::default()
1599 })
1600 .await?;
1601 let target_schema_id: SchemaId = Schema::find()
1602 .select_only()
1603 .column(schema::Column::SchemaId)
1604 .filter(schema::Column::Name.eq("belong_to_target"))
1605 .into_tuple()
1606 .one(&mgr.inner.read().await.db)
1607 .await?
1608 .unwrap();
1609 let txn = mgr.inner.read().await.db.begin().await?;
1610
1611 let mv_obj = CatalogController::create_object(
1612 &txn,
1613 ObjectType::Table,
1614 TEST_OWNER_ID,
1615 Some(TEST_SCHEMA_ID.as_object_id()),
1616 )
1617 .await?;
1618 assert_eq!(mv_obj.belong_to_oid, Some(TEST_SCHEMA_ID.as_object_id()));
1619 assert_eq!(mv_obj.database_id, Some(TEST_DATABASE_ID));
1620 assert_eq!(mv_obj.schema_id, Some(TEST_SCHEMA_ID));
1621 let job_id = mv_obj.oid.as_job_id();
1622 let mv_table_id = job_id.as_mv_table_id();
1623 insert_test_table(
1624 &txn,
1625 mv_table_id,
1626 "mv_belong_to",
1627 TableType::MaterializedView,
1628 None,
1629 "CREATE MATERIALIZED VIEW mv_belong_to AS SELECT 1",
1630 )
1631 .await?;
1632
1633 let internal_obj = CatalogController::create_object(
1634 &txn,
1635 ObjectType::Table,
1636 TEST_OWNER_ID,
1637 Some(job_id.as_object_id()),
1638 )
1639 .await?;
1640 assert_eq!(internal_obj.belong_to_oid, Some(job_id.as_object_id()));
1641 assert_eq!(internal_obj.database_id, Some(TEST_DATABASE_ID));
1642 assert_eq!(internal_obj.schema_id, Some(TEST_SCHEMA_ID));
1643 let internal_table_id = internal_obj.oid.as_table_id();
1644 insert_test_table(
1645 &txn,
1646 internal_table_id,
1647 "__internal_mv_belong_to",
1648 TableType::Internal,
1649 Some(job_id),
1650 "",
1651 )
1652 .await?;
1653 let nested_obj = CatalogController::create_object(
1654 &txn,
1655 ObjectType::Table,
1656 TEST_OWNER_ID,
1657 Some(internal_table_id.as_object_id()),
1658 )
1659 .await?;
1660 txn.commit().await?;
1661
1662 assert!(
1663 mgr.alter_schema(ObjectType::Sink, job_id.as_object_id(), target_schema_id,)
1664 .await
1665 .is_err()
1666 );
1667 mgr.alter_schema(ObjectType::Table, job_id.as_object_id(), target_schema_id)
1668 .await?;
1669
1670 let db = &mgr.inner.read().await.db;
1671 let belonging_object_ids = get_belong_objects(db, job_id.as_object_id())
1672 .await?
1673 .into_iter()
1674 .map(|object| object.oid)
1675 .collect::<HashSet<_>>();
1676 assert_eq!(
1677 belonging_object_ids,
1678 HashSet::from([internal_table_id.as_object_id(), nested_obj.oid])
1679 );
1680 let moved_objects = Object::find()
1681 .filter(object::Column::Oid.is_in([
1682 job_id.as_object_id(),
1683 internal_table_id.as_object_id(),
1684 nested_obj.oid,
1685 ]))
1686 .all(db)
1687 .await?;
1688 assert!(
1689 moved_objects
1690 .iter()
1691 .all(|object| object.schema_id == Some(target_schema_id))
1692 );
1693 assert_eq!(
1694 Object::find_by_id(internal_table_id)
1695 .one(db)
1696 .await?
1697 .unwrap()
1698 .belong_to_oid,
1699 Some(job_id.as_object_id())
1700 );
1701 assert_eq!(
1702 Object::find_by_id(nested_obj.oid)
1703 .one(db)
1704 .await?
1705 .unwrap()
1706 .belong_to_oid,
1707 Some(internal_table_id.as_object_id())
1708 );
1709
1710 Object::delete_by_id(job_id).exec(db).await?;
1711
1712 assert!(Object::find_by_id(job_id).one(db).await?.is_none());
1713 assert!(
1714 Object::find_by_id(internal_table_id)
1715 .one(db)
1716 .await?
1717 .is_none()
1718 );
1719 assert!(Table::find_by_id(mv_table_id).one(db).await?.is_none());
1720 assert!(
1721 Table::find_by_id(internal_table_id)
1722 .one(db)
1723 .await?
1724 .is_none()
1725 );
1726 assert!(Object::find_by_id(nested_obj.oid).one(db).await?.is_none());
1727
1728 Ok(())
1729 }
1730
1731 #[tokio::test]
1732 async fn test_alter_internal_table_schema_rejected() -> MetaResult<()> {
1733 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1734 mgr.create_schema(PbSchema {
1735 database_id: TEST_DATABASE_ID,
1736 name: "internal_table_alter_target".to_owned(),
1737 owner: TEST_OWNER_ID as _,
1738 ..Default::default()
1739 })
1740 .await?;
1741 let target_schema_id: SchemaId = Schema::find()
1742 .select_only()
1743 .column(schema::Column::SchemaId)
1744 .filter(schema::Column::Name.eq("internal_table_alter_target"))
1745 .into_tuple()
1746 .one(&mgr.inner.read().await.db)
1747 .await?
1748 .unwrap();
1749
1750 let txn = mgr.inner.read().await.db.begin().await?;
1751 let parent_obj = CatalogController::create_object(
1752 &txn,
1753 ObjectType::Table,
1754 TEST_OWNER_ID,
1755 Some(TEST_SCHEMA_ID.as_object_id()),
1756 )
1757 .await?;
1758 let parent_job_id = parent_obj.oid.as_job_id();
1759 insert_test_table(
1760 &txn,
1761 parent_job_id.as_mv_table_id(),
1762 "internal_table_parent",
1763 TableType::MaterializedView,
1764 None,
1765 "",
1766 )
1767 .await?;
1768 let internal_obj = CatalogController::create_object(
1769 &txn,
1770 ObjectType::Table,
1771 TEST_OWNER_ID,
1772 Some(parent_job_id.as_object_id()),
1773 )
1774 .await?;
1775 let internal_table_id = internal_obj.oid.as_table_id();
1776 insert_test_table(
1777 &txn,
1778 internal_table_id,
1779 "__internal_table_alter_target",
1780 TableType::Internal,
1781 Some(parent_job_id),
1782 "",
1783 )
1784 .await?;
1785 txn.commit().await?;
1786
1787 for new_schema in [TEST_SCHEMA_ID, target_schema_id] {
1788 assert!(
1789 mgr.alter_schema(
1790 ObjectType::Table,
1791 internal_table_id.as_object_id(),
1792 new_schema,
1793 )
1794 .await
1795 .is_err()
1796 );
1797 }
1798
1799 let internal_obj = Object::find_by_id(internal_table_id)
1800 .one(&mgr.inner.read().await.db)
1801 .await?
1802 .unwrap();
1803 assert_eq!(internal_obj.schema_id, Some(TEST_SCHEMA_ID));
1804 assert_eq!(
1805 internal_obj.belong_to_oid,
1806 Some(parent_job_id.as_object_id())
1807 );
1808
1809 Ok(())
1810 }
1811
1812 #[tokio::test]
1813 async fn test_alter_table_schema_moves_indexes_but_not_subscriptions() -> MetaResult<()> {
1814 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1815 mgr.create_schema(PbSchema {
1816 database_id: TEST_DATABASE_ID,
1817 name: "alter_table_target".to_owned(),
1818 owner: TEST_OWNER_ID as _,
1819 ..Default::default()
1820 })
1821 .await?;
1822 let target_schema_id: SchemaId = Schema::find()
1823 .select_only()
1824 .column(schema::Column::SchemaId)
1825 .filter(schema::Column::Name.eq("alter_table_target"))
1826 .into_tuple()
1827 .one(&mgr.inner.read().await.db)
1828 .await?
1829 .unwrap();
1830
1831 let txn = mgr.inner.read().await.db.begin().await?;
1832 let table_obj = CatalogController::create_object(
1833 &txn,
1834 ObjectType::Table,
1835 TEST_OWNER_ID,
1836 Some(TEST_SCHEMA_ID.as_object_id()),
1837 )
1838 .await?;
1839 let table_id = table_obj.oid.as_table_id();
1840 insert_test_table(
1841 &txn,
1842 table_id,
1843 "mv_with_index_and_subscription",
1844 TableType::MaterializedView,
1845 None,
1846 "CREATE MATERIALIZED VIEW mv_with_index_and_subscription AS SELECT 1",
1847 )
1848 .await?;
1849
1850 let index_obj = CatalogController::create_object(
1851 &txn,
1852 ObjectType::Index,
1853 TEST_OWNER_ID,
1854 Some(TEST_SCHEMA_ID.as_object_id()),
1855 )
1856 .await?;
1857 let index_id = index_obj.oid.as_index_id();
1858 let index_table_id = index_id.as_object_id().as_table_id();
1859 insert_test_table(
1860 &txn,
1861 index_table_id,
1862 "idx_mv_with_index_and_subscription_table",
1863 TableType::Index,
1864 None,
1865 "",
1866 )
1867 .await?;
1868 index::ActiveModel {
1869 index_id: Set(index_id),
1870 name: Set("idx_mv_with_index_and_subscription".to_owned()),
1871 index_table_id: Set(index_table_id),
1872 primary_table_id: Set(table_id),
1873 index_items: Set(Vec::<risingwave_pb::expr::ExprNode>::new().into()),
1874 index_column_properties: Set(None),
1875 index_columns_len: Set(0),
1876 }
1877 .insert(&txn)
1878 .await?;
1879
1880 let index_internal_obj = CatalogController::create_object(
1881 &txn,
1882 ObjectType::Table,
1883 TEST_OWNER_ID,
1884 Some(index_id.as_object_id()),
1885 )
1886 .await?;
1887 let index_internal_table_id = index_internal_obj.oid.as_table_id();
1888 insert_test_table(
1889 &txn,
1890 index_internal_table_id,
1891 "__internal_idx_mv_with_index_and_subscription",
1892 TableType::Internal,
1893 Some(index_id.as_job_id()),
1894 "",
1895 )
1896 .await?;
1897 txn.commit().await?;
1898
1899 let mut subscription = PbSubscription {
1900 name: "subscription_in_original_schema".to_owned(),
1901 definition: "CREATE SUBSCRIPTION subscription_in_original_schema FROM mv_with_index_and_subscription".to_owned(),
1902 retention_seconds: 86400,
1903 database_id: TEST_DATABASE_ID,
1904 schema_id: TEST_SCHEMA_ID,
1905 dependent_table_id: table_id,
1906 owner: TEST_OWNER_ID as _,
1907 subscription_state: SubscriptionState::Created as _,
1908 ..Default::default()
1909 };
1910 mgr.create_subscription_catalog(&mut subscription).await?;
1911
1912 {
1913 let inner = mgr.inner.read().await;
1914 assert_eq!(
1915 Object::find_by_id(index_id)
1916 .one(&inner.db)
1917 .await?
1918 .unwrap()
1919 .belong_to_oid,
1920 Some(TEST_SCHEMA_ID.as_object_id())
1921 );
1922 assert_eq!(
1923 Object::find_by_id(subscription.id)
1924 .one(&inner.db)
1925 .await?
1926 .unwrap()
1927 .belong_to_oid,
1928 Some(TEST_SCHEMA_ID.as_object_id())
1929 );
1930 }
1931
1932 mgr.alter_schema(ObjectType::Table, table_id.as_object_id(), target_schema_id)
1933 .await?;
1934
1935 let db = &mgr.inner.read().await.db;
1936 for object_id in [table_id.as_object_id(), index_id.as_object_id()] {
1937 let object = Object::find_by_id(object_id).one(db).await?.unwrap();
1938 assert_eq!(object.schema_id, Some(target_schema_id));
1939 assert_eq!(object.belong_to_oid, Some(target_schema_id.as_object_id()));
1940 }
1941 let index_internal_object = Object::find_by_id(index_internal_table_id)
1942 .one(db)
1943 .await?
1944 .unwrap();
1945 assert_eq!(index_internal_object.schema_id, Some(target_schema_id));
1946 assert_eq!(
1947 index_internal_object.belong_to_oid,
1948 Some(index_id.as_object_id())
1949 );
1950
1951 let subscription_object = Object::find_by_id(subscription.id).one(db).await?.unwrap();
1952 assert_eq!(subscription_object.schema_id, Some(TEST_SCHEMA_ID));
1953 assert_eq!(
1954 subscription_object.belong_to_oid,
1955 Some(TEST_SCHEMA_ID.as_object_id())
1956 );
1957
1958 Ok(())
1959 }
1960
1961 #[tokio::test]
1962 async fn test_create_function() -> MetaResult<()> {
1963 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
1964 let test_data_type = risingwave_pb::data::DataType {
1965 type_name: risingwave_pb::data::data_type::TypeName::Int32 as _,
1966 ..Default::default()
1967 };
1968 let arg_types = vec![test_data_type.clone()];
1969 let pb_function = PbFunction {
1970 schema_id: TEST_SCHEMA_ID,
1971 database_id: TEST_DATABASE_ID,
1972 name: "test_function".to_owned(),
1973 owner: TEST_OWNER_ID as _,
1974 arg_types,
1975 return_type: Some(test_data_type.clone()),
1976 language: "python".to_owned(),
1977 kind: Some(risingwave_pb::catalog::function::Kind::Scalar(
1978 Default::default(),
1979 )),
1980 ..Default::default()
1981 };
1982 mgr.create_function(pb_function.clone()).await?;
1983 assert!(mgr.create_function(pb_function).await.is_err());
1984
1985 let function = Function::find()
1986 .inner_join(Object)
1987 .filter(
1988 object::Column::DatabaseId
1989 .eq(TEST_DATABASE_ID)
1990 .and(object::Column::SchemaId.eq(TEST_SCHEMA_ID))
1991 .add(function::Column::Name.eq("test_function")),
1992 )
1993 .one(&mgr.inner.read().await.db)
1994 .await?
1995 .unwrap();
1996 assert_eq!(function.return_type.to_protobuf(), test_data_type);
1997 assert_eq!(function.arg_types.to_protobuf().len(), 1);
1998 assert_eq!(function.language, "python");
1999
2000 mgr.create_schema(PbSchema {
2001 database_id: TEST_DATABASE_ID,
2002 name: "function_target".to_owned(),
2003 owner: TEST_OWNER_ID as _,
2004 ..Default::default()
2005 })
2006 .await?;
2007 let target_schema_id: SchemaId = Schema::find()
2008 .select_only()
2009 .column(schema::Column::SchemaId)
2010 .filter(schema::Column::Name.eq("function_target"))
2011 .into_tuple()
2012 .one(&mgr.inner.read().await.db)
2013 .await?
2014 .unwrap();
2015 mgr.alter_schema(
2016 ObjectType::Function,
2017 function.function_id.as_object_id(),
2018 target_schema_id,
2019 )
2020 .await?;
2021 assert_eq!(
2022 Object::find_by_id(function.function_id)
2023 .one(&mgr.inner.read().await.db)
2024 .await?
2025 .unwrap()
2026 .schema_id,
2027 Some(target_schema_id)
2028 );
2029
2030 mgr.drop_object(
2031 ObjectType::Function,
2032 function.function_id,
2033 DropMode::Restrict,
2034 )
2035 .await?;
2036 assert!(
2037 Object::find_by_id(function.function_id)
2038 .one(&mgr.inner.read().await.db)
2039 .await?
2040 .is_none()
2041 );
2042
2043 Ok(())
2044 }
2045
2046 #[tokio::test]
2047 async fn test_alter_relation_rename() -> MetaResult<()> {
2048 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2049 let pb_source = PbSource {
2050 schema_id: TEST_SCHEMA_ID,
2051 database_id: TEST_DATABASE_ID,
2052 name: "s1".to_owned(),
2053 owner: TEST_OWNER_ID as _,
2054 definition: r#"CREATE SOURCE s1 (v1 int) with (
2055 connector = 'kafka',
2056 topic = 'kafka_alter',
2057 properties.bootstrap.server = 'message_queue:29092',
2058 scan.startup.mode = 'earliest'
2059) FORMAT PLAIN ENCODE JSON"#
2060 .to_owned(),
2061 info: Some(StreamSourceInfo {
2062 ..Default::default()
2063 }),
2064 ..Default::default()
2065 };
2066 mgr.create_source(pb_source, None).await?;
2067 let source_id: SourceId = Source::find()
2068 .select_only()
2069 .column(source::Column::SourceId)
2070 .filter(source::Column::Name.eq("s1"))
2071 .into_tuple()
2072 .one(&mgr.inner.read().await.db)
2073 .await?
2074 .unwrap();
2075
2076 let pb_view = PbView {
2077 schema_id: TEST_SCHEMA_ID,
2078 database_id: TEST_DATABASE_ID,
2079 name: "view_1".to_owned(),
2080 owner: TEST_OWNER_ID as _,
2081 sql: "CREATE VIEW view_1 AS SELECT v1 FROM s1".to_owned(),
2082 ..Default::default()
2083 };
2084 mgr.create_view(pb_view, HashSet::from([source_id.as_object_id()]))
2085 .await?;
2086 let view_id: ViewId = View::find()
2087 .select_only()
2088 .column(view::Column::ViewId)
2089 .filter(view::Column::Name.eq("view_1"))
2090 .into_tuple()
2091 .one(&mgr.inner.read().await.db)
2092 .await?
2093 .unwrap();
2094
2095 mgr.alter_name(ObjectType::Source, source_id, "s2").await?;
2096 let source = Source::find_by_id(source_id)
2097 .one(&mgr.inner.read().await.db)
2098 .await?
2099 .unwrap();
2100 assert_eq!(source.name, "s2");
2101 assert_eq!(
2102 source.definition,
2103 "CREATE SOURCE s2 (v1 INT) WITH (\
2104 connector = 'kafka', \
2105 topic = 'kafka_alter', \
2106 properties.bootstrap.server = 'message_queue:29092', \
2107 scan.startup.mode = 'earliest'\
2108) FORMAT PLAIN ENCODE JSON"
2109 );
2110
2111 let view = View::find_by_id(view_id)
2112 .one(&mgr.inner.read().await.db)
2113 .await?
2114 .unwrap();
2115 assert_eq!(
2116 view.definition,
2117 "CREATE VIEW view_1 AS SELECT v1 FROM s2 AS s1"
2118 );
2119
2120 mgr.drop_object(ObjectType::Source, source_id, DropMode::Cascade)
2121 .await?;
2122 assert!(
2123 View::find_by_id(view_id)
2124 .one(&mgr.inner.read().await.db)
2125 .await?
2126 .is_none()
2127 );
2128
2129 Ok(())
2130 }
2131
2132 #[tokio::test]
2133 async fn test_cancel_creating_table_deletes_associated_source() -> MetaResult<()> {
2134 let env = MetaSrvEnv::for_test().await;
2135 let (tx, mut notification_rx) = mpsc::unbounded_channel();
2136 env.notification_manager().insert_sender(
2137 SubscribeType::Frontend,
2138 WorkerKey(HostAddress {
2139 host: "localhost".to_owned(),
2140 port: 1234,
2141 }),
2142 tx,
2143 );
2144 let mgr = CatalogController::new(env).await?;
2145
2146 let mut inner = mgr.inner.write().await;
2147 let txn = inner.db.begin().await?;
2148 let obj = CatalogController::create_object(
2149 &txn,
2150 ObjectType::Table,
2151 TEST_OWNER_ID,
2152 Some(TEST_SCHEMA_ID.as_object_id()),
2153 )
2154 .await?;
2155 let job_id = obj.oid.as_job_id();
2156 let source_obj = CatalogController::create_object(
2157 &txn,
2158 ObjectType::Source,
2159 TEST_OWNER_ID,
2160 Some(job_id.as_object_id()),
2161 )
2162 .await?;
2163 Source::insert(source::ActiveModel::from(PbSource {
2164 id: source_obj.oid.as_source_id(),
2165 schema_id: TEST_SCHEMA_ID,
2166 database_id: TEST_DATABASE_ID,
2167 name: "source_abort_initial".to_owned(),
2168 owner: TEST_OWNER_ID as _,
2169 ..Default::default()
2170 }))
2171 .exec(&txn)
2172 .await?;
2173
2174 table::ActiveModel {
2175 table_id: Set(obj.oid.as_table_id()),
2176 name: Set("table_abort_initial".to_owned()),
2177 optional_associated_source_id: Set(Some(source_obj.oid.as_source_id())),
2178 table_type: Set(TableType::Table),
2179 belongs_to_job_id: Set(None),
2180 columns: Set(vec![].into()),
2181 pk: Set(vec![].into()),
2182 distribution_key: Set(Vec::<i32>::new().into()),
2183 stream_key: Set(Vec::<i32>::new().into()),
2184 append_only: Set(false),
2185 fragment_id: Set(None),
2186 vnode_col_index: Set(None),
2187 row_id_index: Set(None),
2188 value_indices: Set(Vec::<i32>::new().into()),
2189 definition: Set("CREATE TABLE table_abort_initial (v1 INT)".to_owned()),
2190 handle_pk_conflict_behavior: Set(HandleConflictBehavior::NoCheck),
2191 version_column_indices: Set(None),
2192 read_prefix_len_hint: Set(0),
2193 watermark_indices: Set(Vec::<i32>::new().into()),
2194 dist_key_in_pk: Set(Vec::<i32>::new().into()),
2195 dml_fragment_id: Set(None),
2196 cardinality: Set(None),
2197 cleaned_by_watermark: Set(false),
2198 description: Set(None),
2199 version: Set(None),
2200 retention_seconds: Set(None),
2201 cdc_table_id: Set(None),
2202 vnode_count: Set(1),
2203 webhook_info: Set(None),
2204 engine: Set(None),
2205 clean_watermark_index_in_pk: Set(None),
2206 clean_watermark_indices: Set(None),
2207 refreshable: Set(false),
2208 vector_index_info: Set(None),
2209 cdc_table_type: Set(None),
2210 }
2211 .insert(&txn)
2212 .await?;
2213
2214 let internal_obj = CatalogController::create_object(
2215 &txn,
2216 ObjectType::Table,
2217 TEST_OWNER_ID,
2218 Some(job_id.as_object_id()),
2219 )
2220 .await?;
2221 let internal_table_id = internal_obj.oid.as_table_id();
2222 insert_test_table(
2223 &txn,
2224 internal_table_id,
2225 "__internal_mv_abort_initial",
2226 TableType::Internal,
2227 Some(job_id),
2228 "",
2229 )
2230 .await?;
2231
2232 streaming_job::ActiveModel {
2233 job_id: Set(job_id),
2234 job_status: Set(JobStatus::Creating),
2235 create_type: Set(CreateType::Foreground),
2236 timezone: Set(None),
2237 config_override: Set(None),
2238 adaptive_parallelism_strategy: Set(None),
2239 parallelism: Set(StreamingParallelism::Adaptive),
2240 backfill_parallelism: Set(None),
2241 backfill_adaptive_parallelism_strategy: Set(None),
2242 backfill_orders: Set(None),
2243 max_parallelism: Set(1),
2244 specific_resource_group: Set(None),
2245 is_serverless_backfill: Set(false),
2246 refresh_interval_sec: Set(None),
2247 }
2248 .insert(&txn)
2249 .await?;
2250
2251 let (tx, rx) = oneshot::channel();
2252 inner.register_finish_notifier(TEST_DATABASE_ID, job_id, tx);
2253 txn.commit().await?;
2254 drop(inner);
2255
2256 let abort_result = mgr.try_abort_creating_streaming_job(job_id, true).await?;
2257 assert!(abort_result.aborted);
2258 assert_eq!(abort_result.database_id, Some(TEST_DATABASE_ID));
2259
2260 let err = rx
2261 .await
2262 .expect("finish notifier should be notified")
2263 .expect_err("creating job cancellation should fail the create wait");
2264 assert!(err.contains("cancelled"));
2265
2266 let db = &mgr.inner.read().await.db;
2267 assert!(Object::find_by_id(job_id).one(db).await?.is_none());
2268 assert!(StreamingJob::find_by_id(job_id).one(db).await?.is_none());
2269 assert!(
2270 Table::find_by_id(job_id.as_mv_table_id())
2271 .one(db)
2272 .await?
2273 .is_none()
2274 );
2275 assert!(
2276 Object::find_by_id(internal_table_id)
2277 .one(db)
2278 .await?
2279 .is_none()
2280 );
2281 assert!(
2282 Table::find_by_id(internal_table_id)
2283 .one(db)
2284 .await?
2285 .is_none()
2286 );
2287 assert!(
2288 mgr.inner
2289 .read()
2290 .await
2291 .dropped_tables
2292 .contains_key(&internal_table_id)
2293 );
2294 assert!(
2295 Source::find_by_id(source_obj.oid.as_source_id())
2296 .one(db)
2297 .await?
2298 .is_none()
2299 );
2300
2301 let notification = notification_rx
2302 .recv()
2303 .await
2304 .expect("frontend should receive an abort notification")
2305 .expect("abort notification should be valid");
2306 assert_eq!(notification.operation(), NotificationOperation::Delete);
2307 let object_group = match notification.info {
2308 Some(NotificationInfo::ObjectGroup(object_group)) => object_group,
2309 other => panic!("unexpected notification: {other:?}"),
2310 };
2311 assert!(object_group.objects.iter().any(|object| matches!(
2312 &object.object_info,
2313 Some(PbObjectInfo::Table(table)) if table.id == job_id.as_mv_table_id()
2314 )));
2315 assert!(object_group.objects.iter().any(|object| matches!(
2316 &object.object_info,
2317 Some(PbObjectInfo::Source(source)) if source.id == source_obj.oid.as_source_id()
2318 )));
2319
2320 Ok(())
2321 }
2322
2323 #[tokio::test]
2324 async fn test_failed_foreground_creating_job_is_preserved() -> MetaResult<()> {
2325 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2326 let (job_id, table_id) = insert_dirty_creating_job_with_fragment(
2327 &mgr,
2328 FragmentId::new(45),
2329 1,
2330 FragmentTypeMask::empty(),
2331 )
2332 .await?;
2333
2334 let abort_result = mgr.try_abort_creating_streaming_job(job_id, false).await?;
2335 assert!(!abort_result.aborted);
2336 assert_eq!(abort_result.database_id, Some(TEST_DATABASE_ID));
2337
2338 let db = &mgr.inner.read().await.db;
2339 assert!(Object::find_by_id(job_id).one(db).await?.is_some());
2340 assert!(StreamingJob::find_by_id(job_id).one(db).await?.is_some());
2341 assert!(Table::find_by_id(table_id).one(db).await?.is_some());
2342
2343 Ok(())
2344 }
2345
2346 #[tokio::test]
2347 async fn test_failed_created_job_is_preserved() -> MetaResult<()> {
2348 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2349 let (job_id, table_id) = insert_dirty_creating_job_with_fragment(
2350 &mgr,
2351 FragmentId::new(46),
2352 1,
2353 FragmentTypeMask::empty(),
2354 )
2355 .await?;
2356
2357 {
2358 let inner = mgr.inner.read().await;
2359 streaming_job::ActiveModel {
2360 job_id: Set(job_id),
2361 job_status: Set(JobStatus::Created),
2362 ..Default::default()
2363 }
2364 .update(&inner.db)
2365 .await?;
2366 }
2367
2368 let abort_result = mgr.try_abort_creating_streaming_job(job_id, false).await?;
2369 assert!(!abort_result.aborted);
2370 assert_eq!(abort_result.database_id, Some(TEST_DATABASE_ID));
2371 let db = &mgr.inner.read().await.db;
2372 assert!(Object::find_by_id(job_id).one(db).await?.is_some());
2373 assert!(StreamingJob::find_by_id(job_id).one(db).await?.is_some());
2374 assert!(Table::find_by_id(table_id).one(db).await?.is_some());
2375
2376 Ok(())
2377 }
2378
2379 #[tokio::test]
2380 async fn test_clean_dirty_creating_jobs_records_dropped_tables_for_per_db_recovery()
2381 -> MetaResult<()> {
2382 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2383
2384 let inner = mgr.inner.write().await;
2385 let txn = inner.db.begin().await?;
2386 let mv_obj = CatalogController::create_object(
2387 &txn,
2388 ObjectType::Table,
2389 TEST_OWNER_ID,
2390 Some(TEST_SCHEMA_ID.as_object_id()),
2391 )
2392 .await?;
2393 let job_id = mv_obj.oid.as_job_id();
2394 let mv_table_id = job_id.as_mv_table_id();
2395 insert_test_table(
2396 &txn,
2397 mv_table_id,
2398 "mv_dirty",
2399 TableType::MaterializedView,
2400 None,
2401 "CREATE MATERIALIZED VIEW mv_dirty AS SELECT 1",
2402 )
2403 .await?;
2404
2405 let internal_obj = CatalogController::create_object(
2406 &txn,
2407 ObjectType::Table,
2408 TEST_OWNER_ID,
2409 Some(job_id.as_object_id()),
2410 )
2411 .await?;
2412 let internal_table_id = internal_obj.oid.as_table_id();
2413 insert_test_table(
2414 &txn,
2415 internal_table_id,
2416 "__internal_mv_dirty",
2417 TableType::Internal,
2418 Some(job_id),
2419 "",
2420 )
2421 .await?;
2422
2423 streaming_job::ActiveModel {
2424 job_id: Set(job_id),
2425 job_status: Set(JobStatus::Initial),
2426 create_type: Set(CreateType::Foreground),
2427 timezone: Set(None),
2428 config_override: Set(None),
2429 adaptive_parallelism_strategy: Set(None),
2430 parallelism: Set(StreamingParallelism::Adaptive),
2431 backfill_parallelism: Set(None),
2432 backfill_adaptive_parallelism_strategy: Set(None),
2433 backfill_orders: Set(None),
2434 max_parallelism: Set(1),
2435 specific_resource_group: Set(None),
2436 is_serverless_backfill: Set(false),
2437 refresh_interval_sec: Set(None),
2438 }
2439 .insert(&txn)
2440 .await?;
2441 txn.commit().await?;
2442 drop(inner);
2443
2444 let cleaned = mgr
2445 .clean_dirty_creating_jobs(Some(TEST_DATABASE_ID))
2446 .await?;
2447 assert_eq!(cleaned.streaming_job_ids, vec![job_id]);
2448 assert!(cleaned.source_ids.is_empty());
2449 let mut dropped_table_ids = cleaned.dropped_table_ids;
2450 dropped_table_ids.sort_unstable();
2451 assert_eq!(dropped_table_ids, vec![mv_table_id, internal_table_id]);
2452
2453 let inner = mgr.inner.read().await;
2454 assert!(inner.dropped_tables.contains_key(&mv_table_id));
2455 assert!(inner.dropped_tables.contains_key(&internal_table_id));
2456 assert!(Object::find_by_id(job_id).one(&inner.db).await?.is_none());
2457 assert!(
2458 Object::find_by_id(internal_table_id)
2459 .one(&inner.db)
2460 .await?
2461 .is_none()
2462 );
2463 assert!(
2464 StreamingJob::find_by_id(job_id)
2465 .one(&inner.db)
2466 .await?
2467 .is_none()
2468 );
2469 assert!(
2470 Table::find_by_id(mv_table_id)
2471 .one(&inner.db)
2472 .await?
2473 .is_none()
2474 );
2475 assert!(
2476 Table::find_by_id(internal_table_id)
2477 .one(&inner.db)
2478 .await?
2479 .is_none()
2480 );
2481
2482 Ok(())
2483 }
2484
2485 #[tokio::test]
2486 async fn test_clean_dirty_creating_jobs_notifies_serving_mapping_fragment_delete()
2487 -> MetaResult<()> {
2488 let env = MetaSrvEnv::for_test().await;
2489 let (local_notification_tx, mut local_notification_rx) = mpsc::unbounded_channel();
2490 env.notification_manager()
2491 .insert_local_sender(local_notification_tx);
2492 let mgr = CatalogController::new(env).await?;
2493 let fragment_id = FragmentId::new(3);
2494 let (job_id, mv_table_id) = insert_dirty_creating_job_with_fragment(
2495 &mgr,
2496 fragment_id,
2497 1,
2498 FragmentTypeMask::from(FragmentTypeFlag::Values as u32),
2499 )
2500 .await?;
2501
2502 assert!(
2503 mgr.fragment_serving_infos()
2504 .await?
2505 .contains_key(&fragment_id)
2506 );
2507
2508 let cleaned = mgr
2509 .clean_dirty_creating_jobs(Some(TEST_DATABASE_ID))
2510 .await?;
2511 assert_eq!(cleaned.streaming_job_ids, vec![job_id]);
2512
2513 let inner = mgr.inner.read().await;
2514 assert!(Object::find_by_id(job_id).one(&inner.db).await?.is_none());
2515 assert!(
2516 StreamingJob::find_by_id(job_id)
2517 .one(&inner.db)
2518 .await?
2519 .is_none()
2520 );
2521 assert!(
2522 Table::find_by_id(mv_table_id)
2523 .one(&inner.db)
2524 .await?
2525 .is_none()
2526 );
2527 drop(inner);
2528 assert!(
2529 !mgr.fragment_serving_infos()
2530 .await?
2531 .contains_key(&fragment_id)
2532 );
2533
2534 let notification = local_notification_rx.try_recv().expect(
2535 "dirty-job cleanup must notify the serving mapping worker about deleted fragments",
2536 );
2537 match notification {
2538 LocalNotification::ServingFragmentMappingsDelete(fragment_ids) => {
2539 assert_eq!(fragment_ids, vec![fragment_id]);
2540 }
2541 notification => panic!("unexpected local notification: {notification:?}"),
2542 }
2543
2544 Ok(())
2545 }
2546
2547 #[tokio::test]
2548 async fn test_abort_creating_subscription_commits_delete() -> MetaResult<()> {
2549 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2550 let pb_view = PbView {
2551 schema_id: TEST_SCHEMA_ID,
2552 database_id: TEST_DATABASE_ID,
2553 name: "subscription_dep_view".to_owned(),
2554 owner: TEST_OWNER_ID as _,
2555 sql: "CREATE VIEW subscription_dep_view AS SELECT 1".to_owned(),
2556 ..Default::default()
2557 };
2558 mgr.create_view(pb_view, HashSet::new()).await?;
2559
2560 let view_id: ViewId = View::find()
2561 .select_only()
2562 .column(view::Column::ViewId)
2563 .filter(view::Column::Name.eq("subscription_dep_view"))
2564 .into_tuple()
2565 .one(&mgr.inner.read().await.db)
2566 .await?
2567 .unwrap();
2568
2569 let mut pb_subscription = PbSubscription {
2570 name: "subscription_to_abort".to_owned(),
2571 definition: "CREATE SUBSCRIPTION subscription_to_abort FROM subscription_dep_view"
2572 .to_owned(),
2573 retention_seconds: 86400,
2574 database_id: TEST_DATABASE_ID,
2575 schema_id: TEST_SCHEMA_ID,
2576 dependent_table_id: view_id.as_object_id().as_table_id(),
2577 owner: TEST_OWNER_ID as _,
2578 subscription_state: SubscriptionState::Init as _,
2579 ..Default::default()
2580 };
2581 mgr.create_subscription_catalog(&mut pb_subscription)
2582 .await?;
2583
2584 mgr.try_abort_creating_subscription(pb_subscription.id)
2585 .await?;
2586
2587 assert!(
2588 Subscription::find_by_id(pb_subscription.id)
2589 .one(&mgr.inner.read().await.db)
2590 .await?
2591 .is_none()
2592 );
2593
2594 Ok(())
2595 }
2596
2597 #[tokio::test]
2598 async fn test_drop_table_cascade_drops_dependent_subscription() -> MetaResult<()> {
2599 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2600
2601 let inner = mgr.inner.write().await;
2602 let txn = inner.db.begin().await?;
2603 let table_obj = CatalogController::create_object(
2604 &txn,
2605 ObjectType::Table,
2606 TEST_OWNER_ID,
2607 Some(TEST_SCHEMA_ID.as_object_id()),
2608 )
2609 .await?;
2610 let table_id = table_obj.oid.as_table_id();
2611 insert_test_table(
2612 &txn,
2613 table_id,
2614 "subscription_dep_table",
2615 TableType::Table,
2616 None,
2617 "CREATE TABLE subscription_dep_table (v1 INT)",
2618 )
2619 .await?;
2620 txn.commit().await?;
2621 drop(inner);
2622
2623 let mut pb_subscription = PbSubscription {
2624 name: "subscription_to_drop_with_table".to_owned(),
2625 definition:
2626 "CREATE SUBSCRIPTION subscription_to_drop_with_table FROM subscription_dep_table"
2627 .to_owned(),
2628 retention_seconds: 86400,
2629 database_id: TEST_DATABASE_ID,
2630 schema_id: TEST_SCHEMA_ID,
2631 dependent_table_id: table_id,
2632 owner: TEST_OWNER_ID as _,
2633 subscription_state: SubscriptionState::Created as _,
2634 ..Default::default()
2635 };
2636 mgr.create_subscription_catalog(&mut pb_subscription)
2637 .await?;
2638
2639 mgr.drop_object(ObjectType::Table, table_id, DropMode::Cascade)
2640 .await?;
2641
2642 let db = &mgr.inner.read().await.db;
2643 assert!(Table::find_by_id(table_id).one(db).await?.is_none());
2644 assert!(
2645 Object::find_by_id(table_id.as_object_id())
2646 .one(db)
2647 .await?
2648 .is_none()
2649 );
2650 assert!(
2651 Subscription::find_by_id(pb_subscription.id)
2652 .one(db)
2653 .await?
2654 .is_none()
2655 );
2656 assert!(
2657 Object::find_by_id(pb_subscription.id.as_object_id())
2658 .one(db)
2659 .await?
2660 .is_none()
2661 );
2662
2663 Ok(())
2664 }
2665
2666 #[tokio::test]
2667 async fn test_get_table_change_log_truncate_info() -> MetaResult<()> {
2668 let mgr = CatalogController::new(MetaSrvEnv::for_test().await).await?;
2669 let pb_view = PbView {
2670 schema_id: TEST_SCHEMA_ID,
2671 database_id: TEST_DATABASE_ID,
2672 name: "change_log_upstream".to_owned(),
2673 owner: TEST_OWNER_ID as _,
2674 sql: "CREATE VIEW change_log_upstream AS SELECT 1".to_owned(),
2675 ..Default::default()
2676 };
2677 mgr.create_view(pb_view, HashSet::new()).await?;
2678 let upstream_table_id: TableId = View::find()
2679 .select_only()
2680 .column(view::Column::ViewId)
2681 .filter(view::Column::Name.eq("change_log_upstream"))
2682 .into_tuple::<ViewId>()
2683 .one(&mgr.inner.read().await.db)
2684 .await?
2685 .unwrap()
2686 .as_object_id()
2687 .as_table_id();
2688 let mut subscription = PbSubscription {
2689 name: "change_log_subscription".to_owned(),
2690 definition: "CREATE SUBSCRIPTION change_log_subscription FROM change_log_upstream"
2691 .to_owned(),
2692 retention_seconds: 123,
2693 database_id: TEST_DATABASE_ID,
2694 schema_id: TEST_SCHEMA_ID,
2695 dependent_table_id: upstream_table_id,
2696 owner: TEST_OWNER_ID as _,
2697 subscription_state: SubscriptionState::Created as _,
2698 ..Default::default()
2699 };
2700 mgr.create_subscription_catalog(&mut subscription).await?;
2701
2702 let inner = mgr.inner.write().await;
2703 let txn = inner.db.begin().await?;
2704 let (job_id, _, state_table_id) =
2705 insert_test_streaming_job(&txn, "snapshot_job", true, None).await?;
2706 let mut job = streaming_job::Entity::find_by_id(job_id)
2707 .one(&txn)
2708 .await?
2709 .unwrap()
2710 .into_active_model();
2711 job.job_status = Set(JobStatus::Creating);
2712 job.update(&txn).await?;
2713 fragment::ActiveModel {
2714 fragment_id: Set(FragmentId::new(100)),
2715 job_id: Set(job_id),
2716 fragment_type_mask: Set(FragmentTypeFlag::SnapshotBackfillStreamScan as i32),
2717 distribution_type: Set(fragment::DistributionType::Hash),
2718 stream_node: Set(StreamNode::from(&PbStreamNode {
2719 node_body: Some(PbNodeBody::StreamScan(Box::new(StreamScanNode {
2720 table_id: upstream_table_id,
2721 stream_scan_type: StreamScanType::SnapshotBackfill as i32,
2722 snapshot_backfill_epoch: None,
2723 ..Default::default()
2724 }))),
2725 ..Default::default()
2726 })),
2727 state_table_ids: Set(vec![state_table_id].into()),
2728 upstream_fragment_id: Set(I32Array::default()),
2729 vnode_count: Set(1),
2730 parallelism: Set(None),
2731 }
2732 .insert(&txn)
2733 .await?;
2734 txn.commit().await?;
2735 drop(inner);
2736
2737 let truncate_info = mgr.get_table_change_log_truncate_info().await?;
2738 assert_eq!(
2739 truncate_info.subscription_retention_seconds,
2740 HashMap::from([(upstream_table_id, 123)])
2741 );
2742 assert_eq!(truncate_info.independent_jobs.len(), 1);
2743 let independent_job = &truncate_info.independent_jobs[0];
2744 assert_eq!(independent_job.job_id, job_id);
2745 assert_eq!(
2746 independent_job.state_table_ids,
2747 HashSet::from([state_table_id])
2748 );
2749 assert_eq!(
2750 independent_job.upstream_table_snapshot_epochs,
2751 HashMap::from([(upstream_table_id, None)])
2752 );
2753
2754 Ok(())
2755 }
2756}