Skip to main content

risingwave_meta/controller/catalog/
list_op.rs

1// Copyright 2024 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use risingwave_common::catalog::FragmentTypeMask;
16use risingwave_common::id::JobId;
17use risingwave_meta_model::refresh_job::{self, RefreshState};
18use risingwave_pb::meta::ObjectDependency as PbObjectDependency;
19use sea_orm::prelude::DateTime;
20
21use super::*;
22use crate::controller::fragment::FragmentTypeMaskExt;
23use crate::controller::utils::load_streaming_jobs_by_ids;
24
25impl CatalogController {
26    pub async fn list_time_travel_table_ids(&self) -> MetaResult<Vec<TableId>> {
27        self.inner.read().await.list_time_travel_table_ids().await
28    }
29
30    pub async fn list_refresh_jobs(&self) -> MetaResult<Vec<refresh_job::Model>> {
31        let inner = self.inner.read().await;
32        Ok(RefreshJob::find().all(&inner.db).await?)
33    }
34
35    pub async fn get_refresh_job_state_by_table_id(
36        &self,
37        table_id: TableId,
38    ) -> MetaResult<RefreshState> {
39        let inner = self.inner.read().await;
40        let (refresh_job_state,): (RefreshState,) = RefreshJob::find_by_id(table_id)
41            .select_only()
42            .select_column(refresh_job::Column::CurrentStatus)
43            .into_tuple()
44            .one(&inner.db)
45            .await?
46            .ok_or_else(|| MetaError::catalog_id_not_found("refresh_job", table_id))?;
47        Ok(refresh_job_state)
48    }
49
50    pub async fn list_refreshable_table_ids(&self) -> MetaResult<Vec<TableId>> {
51        let inner = self.inner.read().await;
52        Ok(Table::find()
53            .select_only()
54            .column(table::Column::TableId)
55            .filter(table::Column::Refreshable.eq(true))
56            .into_tuple()
57            .all(&inner.db)
58            .await?)
59    }
60
61    pub async fn list_stream_job_desc_for_telemetry(
62        &self,
63    ) -> MetaResult<Vec<MetaTelemetryJobDesc>> {
64        let inner = self.inner.read().await;
65        let info: Vec<(TableId, Option<Property>)> = Table::find()
66            .select_only()
67            .column(table::Column::TableId)
68            .column(source::Column::WithProperties)
69            .join(JoinType::LeftJoin, table::Relation::Source.def())
70            .filter(
71                table::Column::TableType
72                    .eq(TableType::Table)
73                    .or(table::Column::TableType.eq(TableType::MaterializedView)),
74            )
75            .into_tuple()
76            .all(&inner.db)
77            .await?;
78
79        Ok(info
80            .into_iter()
81            .map(|(table_id, properties)| {
82                let connector_info = if let Some(inner_props) = properties {
83                    inner_props
84                        .inner_ref()
85                        .get(UPSTREAM_SOURCE_KEY)
86                        .map(|v| v.to_lowercase())
87                } else {
88                    None
89                };
90                MetaTelemetryJobDesc {
91                    table_id,
92                    connector: connector_info,
93                    optimization: vec![],
94                }
95            })
96            .collect())
97    }
98
99    pub async fn list_creating_jobs(
100        &self,
101        include_initial: bool,
102        database_id: Option<DatabaseId>,
103    ) -> MetaResult<Vec<(JobId, String, DateTime, CreateType, bool)>> {
104        let inner = self.inner.read().await;
105        let status_cond = if include_initial {
106            streaming_job::Column::JobStatus.is_in([JobStatus::Initial, JobStatus::Creating])
107        } else {
108            streaming_job::Column::JobStatus.eq(JobStatus::Creating)
109        };
110        let database_cond = database_id
111            .map(|database_id| object::Column::DatabaseId.eq(database_id))
112            .unwrap_or_else(|| SimpleExpr::from(true));
113        let filter_cond = status_cond.and(database_cond);
114        let object_columns = [object::Column::InitializedAt];
115        let streaming_job_columns = [
116            streaming_job::Column::CreateType,
117            streaming_job::Column::IsServerlessBackfill,
118        ];
119        let mut table_info: Vec<(JobId, String, DateTime, CreateType, bool)> = Table::find()
120            .select_only()
121            .columns([table::Column::TableId, table::Column::Definition])
122            .columns(object_columns)
123            .columns(streaming_job_columns)
124            .join(JoinType::LeftJoin, table::Relation::Object1.def())
125            .join(JoinType::LeftJoin, object::Relation::StreamingJob.def())
126            .filter(filter_cond.clone())
127            .into_tuple()
128            .all(&inner.db)
129            .await?;
130        let sink_info: Vec<(JobId, String, DateTime, CreateType, bool)> = Sink::find()
131            .select_only()
132            .columns([sink::Column::SinkId, sink::Column::Definition])
133            .columns(object_columns)
134            .columns(streaming_job_columns)
135            .join(JoinType::LeftJoin, sink::Relation::Object.def())
136            .join(JoinType::LeftJoin, object::Relation::StreamingJob.def())
137            .filter(filter_cond)
138            .into_tuple()
139            .all(&inner.db)
140            .await?;
141
142        table_info.extend(sink_info);
143
144        Ok(table_info)
145    }
146
147    pub async fn list_databases(&self) -> MetaResult<Vec<PbDatabase>> {
148        let inner = self.inner.read().await;
149        inner.list_databases().await
150    }
151
152    pub async fn list_all_object_dependencies(&self) -> MetaResult<Vec<PbObjectDependency>> {
153        let inner = self.inner.read().await;
154        let txn = inner.db.begin().await?;
155        let dependencies = list_object_dependencies(&txn, true).await?;
156        txn.commit().await?;
157        Ok(dependencies)
158    }
159
160    pub async fn list_created_object_dependencies(&self) -> MetaResult<Vec<PbObjectDependency>> {
161        let inner = self.inner.read().await;
162        let txn = inner.db.begin().await?;
163        let dependencies = list_object_dependencies(&txn, false).await?;
164        txn.commit().await?;
165        Ok(dependencies)
166    }
167
168    pub async fn list_schemas(&self) -> MetaResult<Vec<PbSchema>> {
169        let inner = self.inner.read().await;
170        inner.list_schemas().await
171    }
172
173    /// [`Self::list_tables_by_type`] with all types.
174    pub async fn list_all_state_tables(&self) -> MetaResult<Vec<PbTable>> {
175        let inner = self.inner.read().await;
176        inner.list_all_state_tables().await
177    }
178
179    pub async fn list_readonly_table_ids(&self, schema_id: SchemaId) -> MetaResult<Vec<TableId>> {
180        let inner = self.inner.read().await;
181        let table_ids: Vec<TableId> = Table::find()
182            .select_only()
183            .column(table::Column::TableId)
184            .join(JoinType::InnerJoin, table::Relation::Object1.def())
185            .filter(
186                object::Column::SchemaId
187                    .eq(schema_id)
188                    .and(table::Column::TableType.ne(TableType::Table)),
189            )
190            .into_tuple()
191            .all(&inner.db)
192            .await?;
193        Ok(table_ids)
194    }
195
196    pub async fn list_dml_table_ids(&self, schema_id: SchemaId) -> MetaResult<Vec<TableId>> {
197        let inner = self.inner.read().await;
198        let table_ids: Vec<TableId> = Table::find()
199            .select_only()
200            .column(table::Column::TableId)
201            .join(JoinType::InnerJoin, table::Relation::Object1.def())
202            .filter(
203                object::Column::SchemaId
204                    .eq(schema_id)
205                    .and(table::Column::TableType.eq(TableType::Table)),
206            )
207            .into_tuple()
208            .all(&inner.db)
209            .await?;
210        Ok(table_ids)
211    }
212
213    pub async fn list_view_ids(&self, schema_id: SchemaId) -> MetaResult<Vec<ViewId>> {
214        let inner = self.inner.read().await;
215        let view_ids: Vec<ViewId> = View::find()
216            .select_only()
217            .column(view::Column::ViewId)
218            .join(JoinType::InnerJoin, view::Relation::Object.def())
219            .filter(object::Column::SchemaId.eq(schema_id))
220            .into_tuple()
221            .all(&inner.db)
222            .await?;
223        Ok(view_ids)
224    }
225
226    /// Use [`Self::list_all_state_tables`] to get all types.
227    pub async fn list_tables_by_type(&self, table_type: TableType) -> MetaResult<Vec<PbTable>> {
228        let inner = self.inner.read().await;
229        let table_objs = Table::find()
230            .find_also_related(Object)
231            .filter(table::Column::TableType.eq(table_type))
232            .all(&inner.db)
233            .await?;
234        let streaming_jobs = load_streaming_jobs_by_ids(
235            &inner.db,
236            table_objs.iter().map(|(table, _)| table.job_id()),
237        )
238        .await?;
239        Ok(table_objs
240            .into_iter()
241            .map(|(table, obj)| {
242                let job_id = table.job_id();
243                let streaming_job = streaming_jobs.get(&job_id).cloned();
244                ObjectModel(table, obj.unwrap(), streaming_job).into()
245            })
246            .collect())
247    }
248
249    pub async fn list_sources(&self) -> MetaResult<Vec<PbSource>> {
250        let inner = self.inner.read().await;
251        inner.list_sources().await
252    }
253
254    // Return a hashmap to distinguish whether each source is shared or not.
255    pub async fn list_source_id_with_shared_types(&self) -> MetaResult<HashMap<SourceId, bool>> {
256        let inner = self.inner.read().await;
257        let source_ids: Vec<(SourceId, Option<StreamSourceInfo>)> = Source::find()
258            .select_only()
259            .columns([source::Column::SourceId, source::Column::SourceInfo])
260            .into_tuple()
261            .all(&inner.db)
262            .await?;
263
264        Ok(source_ids
265            .into_iter()
266            .map(|(source_id, info)| {
267                (
268                    source_id,
269                    info.map(|info| info.to_protobuf().cdc_source_job)
270                        .unwrap_or(false),
271                )
272            })
273            .collect())
274    }
275
276    pub async fn list_connections(&self) -> MetaResult<Vec<PbConnection>> {
277        let inner = self.inner.read().await;
278        let conn_objs = Connection::find()
279            .find_also_related(Object)
280            .all(&inner.db)
281            .await?;
282        Ok(conn_objs
283            .into_iter()
284            .map(|(conn, obj)| ObjectModel(conn, obj.unwrap(), None).into())
285            .collect())
286    }
287
288    pub async fn list_source_ids(&self, schema_id: SchemaId) -> MetaResult<Vec<SourceId>> {
289        let inner = self.inner.read().await;
290        let source_ids: Vec<SourceId> = Source::find()
291            .select_only()
292            .column(source::Column::SourceId)
293            .join(JoinType::InnerJoin, source::Relation::Object.def())
294            .filter(object::Column::SchemaId.eq(schema_id))
295            .into_tuple()
296            .all(&inner.db)
297            .await?;
298        Ok(source_ids)
299    }
300
301    pub async fn list_indexes(&self) -> MetaResult<Vec<PbIndex>> {
302        let inner = self.inner.read().await;
303        inner.list_indexes().await
304    }
305
306    pub async fn list_sinks(&self) -> MetaResult<Vec<PbSink>> {
307        let inner = self.inner.read().await;
308        inner.list_sinks().await
309    }
310
311    pub async fn list_subscriptions(&self) -> MetaResult<Vec<PbSubscription>> {
312        let inner = self.inner.read().await;
313        inner.list_subscriptions().await
314    }
315
316    pub async fn list_views(&self) -> MetaResult<Vec<PbView>> {
317        let inner = self.inner.read().await;
318        inner.list_views().await
319    }
320
321    pub async fn list_users(&self) -> MetaResult<Vec<PbUserInfo>> {
322        let inner = self.inner.read().await;
323        inner.list_users().await
324    }
325
326    pub async fn list_functions(&self) -> MetaResult<Vec<PbFunction>> {
327        let inner = self.inner.read().await;
328        inner.list_functions().await
329    }
330
331    /// `Unmigrated` refers to table-fragments that have not yet been migrated to the new plan (for now, this means
332    /// table-fragments that do not use `UpstreamSinkUnion` operator to receive multiple upstream sinks)
333    pub async fn list_unmigrated_tables(&self) -> MetaResult<Vec<PbTable>> {
334        let inner = self.inner.read().await;
335
336        let table_objs = Table::find()
337            .filter(table::Column::TableType.eq(TableType::Table))
338            .find_also_related(Object)
339            .join(JoinType::InnerJoin, object::Relation::Fragment.def())
340            .filter(FragmentTypeMask::intersects(FragmentTypeFlag::Mview).and(
341                FragmentTypeMask::disjoint(FragmentTypeFlag::UpstreamSinkUnion),
342            ))
343            .all(&inner.db)
344            .await?;
345        let streaming_jobs = load_streaming_jobs_by_ids(
346            &inner.db,
347            table_objs.iter().map(|(table, _)| table.job_id()),
348        )
349        .await?;
350
351        Ok(table_objs
352            .into_iter()
353            .map(|(table, obj)| {
354                let job_id = table.job_id();
355                let streaming_job = streaming_jobs.get(&job_id).cloned();
356                ObjectModel(table, obj.unwrap(), streaming_job).into()
357            })
358            .collect())
359    }
360
361    pub async fn list_sink_ids(&self, database_id: Option<DatabaseId>) -> MetaResult<Vec<SinkId>> {
362        let inner = self.inner.read().await;
363
364        let mut query = Sink::find().select_only().column(sink::Column::SinkId);
365
366        if let Some(database_id) = database_id {
367            query = query
368                .join(JoinType::InnerJoin, sink::Relation::Object.def())
369                .filter(object::Column::DatabaseId.eq(database_id));
370        }
371
372        let sink_ids: Vec<SinkId> = query.into_tuple().all(&inner.db).await?;
373
374        Ok(sink_ids)
375    }
376}