Skip to main content

risingwave_meta/controller/catalog/
test.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15#[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            // Use create_job_catalog to trigger construct_sink_cycle_check_query for regression
382            // testing purpose to ensure no circular issue causing infinite recursion when
383            // cte_referencing
384            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        // Ensure the test sinks were created
437        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        // Ensure no object_dependency created for (target_table, sink) which could cause circular
450        // issue
451        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        // Check that the cascade drop successfully dropped
483        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        // Sanity checks that sources were not dropped
492        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            // Ensures the replacement failed because a referring streaming job is still creating.
503            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        // Ensures the incoming sink is still creating, which should block replacement.
555        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        // Ensures the replacement targets the upstream MV that the creating sink depends on.
566        assert_eq!(replacement.id(), target_mv_id);
567
568        // Ensures replacement rejects the upstream MV while its referring sink is creating.
569        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        // Ensures the rejection error reports the expected concurrency reason.
574        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        // Ensures the incoming sink starts as a creating job before the test marks it created.
627        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        // Ensures the persisted sink row is the sink created by this test.
643        assert_eq!(sink_model.job_id, sink_id);
644        // Ensures a fully created incoming sink does not block upstream MV replacement.
645        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        // Ensures the replacement targets the upstream MV that the created sink depends on.
657        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        // Ensures replacement creates a distinct temporary job instead of reusing the original MV id.
664        assert_ne!(tmp_model.job_id, source_mv_id);
665        // Ensures the temporary replacement job is created but not finished yet.
666        assert_eq!(tmp_model.job_status, JobStatus::Initial);
667
668        let inner = mgr.inner.read().await;
669        let db = &inner.db;
670        // Ensures the created sink still depends on the upstream MV being replaced.
671        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        // Ensures replacement records the temporary job as a dependent of the original MV.
680        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        // A foreground table and its internal table are both visible while creating.
791        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        // A foreground index and its index table are both visible while creating.
837        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        // A creating shared source is included in restart snapshots as well.
882        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(&current_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}