risingwave_meta/controller/catalog/
list_op.rs1use 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 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 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 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 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}