risingwave_meta/controller/catalog/
drop_op.rs1use risingwave_common::catalog::ICEBERG_SINK_PREFIX;
16use risingwave_pb::catalog::subscription::PbSubscriptionState;
17use risingwave_pb::telemetry::PbTelemetryDatabaseObject;
18use sea_orm::{ColumnTrait, DatabaseTransaction, EntityTrait, ModelTrait, QueryFilter};
19
20use super::*;
21impl CatalogController {
22 pub async fn drop_object(
25 &self,
26 object_type: ObjectType,
27 object_id: impl Into<ObjectId>,
28 drop_mode: DropMode,
29 ) -> MetaResult<(ReleaseContext, NotificationVersion)> {
30 let object_id = object_id.into();
31 let mut inner = self.inner.write().await;
32 let txn = inner.db.begin().await?;
33
34 let obj: PartialObject = Object::find_by_id(object_id)
35 .into_partial_model()
36 .one(&txn)
37 .await?
38 .ok_or_else(|| MetaError::catalog_id_not_found(object_type.as_str(), object_id))?;
39 assert_eq!(obj.obj_type, object_type);
40 let database_id = if object_type == ObjectType::Database {
41 object_id.as_database_id()
42 } else {
43 obj.database_id
44 .ok_or_else(|| anyhow!("dropped object should have database_id"))?
45 };
46
47 if obj.obj_type == ObjectType::Subscription {
49 validate_subscription_deletion(&txn, object_id.as_subscription_id()).await?;
50 }
51
52 let mut removed_objects = match drop_mode {
53 DropMode::Cascade => {
54 get_referring_objects_cascade(object_id, object_type, &txn).await?
55 }
56 DropMode::Restrict => match object_type {
57 ObjectType::Database => unreachable!("database always be dropped in cascade mode"),
58 ObjectType::Schema => {
59 ensure_schema_empty(object_id.as_schema_id(), &txn).await?;
60 Default::default()
61 }
62 ObjectType::Table => {
63 let objects = validate_restrict_drop_and_collect_owned_objects(
64 object_type,
65 object_id,
66 &txn,
67 )
68 .await?;
69 for obj in objects.iter().filter(|object| {
70 object.obj_type == ObjectType::Source || object.obj_type == ObjectType::Sink
71 }) {
72 report_drop_object(obj.obj_type, obj.oid, &txn).await;
73 }
74 assert!(
75 objects.iter().all(|obj| obj.obj_type == ObjectType::Index
76 || obj.obj_type == ObjectType::Sink),
77 "only index and iceberg sink could be dropped in restrict mode"
78 );
79 for obj in &objects {
80 validate_restrict_drop_and_collect_owned_objects(
81 obj.obj_type,
82 obj.oid,
83 &txn,
84 )
85 .await?;
86 }
87 objects
88 }
89 object_type @ (ObjectType::Source | ObjectType::Sink) => {
90 validate_restrict_drop_and_collect_owned_objects(object_type, object_id, &txn)
91 .await?;
92 report_drop_object(object_type, object_id, &txn).await;
93 vec![]
94 }
95
96 ObjectType::View
97 | ObjectType::Index
98 | ObjectType::Function
99 | ObjectType::Connection
100 | ObjectType::Subscription
101 | ObjectType::Secret => {
102 validate_restrict_drop_and_collect_owned_objects(object_type, object_id, &txn)
103 .await?;
104 vec![]
105 }
106 },
107 };
108
109 removed_objects.push(obj);
110 let mut removed_object_ids: HashSet<_> =
111 removed_objects.iter().map(|obj| obj.oid).collect();
112
113 for obj in &removed_objects {
114 if obj.obj_type == ObjectType::Sink {
115 let sink = Sink::find_by_id(obj.oid.as_sink_id())
116 .one(&txn)
117 .await?
118 .ok_or_else(|| MetaError::catalog_id_not_found("sink", obj.oid))?;
119
120 if let Some(target_table) = sink.target_table
121 && !removed_object_ids.contains(&target_table.as_object_id())
122 && !has_table_been_migrated(&txn, target_table).await?
123 {
124 return Err(anyhow::anyhow!(
125 "Dropping sink into table is not allowed for unmigrated table {}. Please migrate it first.",
126 target_table
127 ).into());
128 }
129 }
130 }
131
132 let root_objects = Object::find()
135 .filter(object::Column::Oid.is_in(removed_object_ids.iter().copied()))
136 .all(&txn)
137 .await?;
138 let belonging_objects =
139 get_belong_objects_by_ids(&txn, removed_objects.iter().map(|obj| obj.oid)).await?;
140 removed_object_ids.extend(belonging_objects.iter().map(|obj| obj.oid));
141 let mut objects_to_remove = root_objects.clone();
142 objects_to_remove.extend(belonging_objects.iter().cloned());
143 let removed_catalog_models = load_object_models(&txn, &objects_to_remove).await?;
144 removed_objects.extend(belonging_objects.into_iter().map(|obj| PartialObject {
145 oid: obj.oid,
146 obj_type: obj.obj_type,
147 schema_id: obj.schema_id,
148 database_id: obj.database_id,
149 }));
150
151 let removed_table_ids = removed_objects
152 .iter()
153 .filter(|obj| obj.obj_type == ObjectType::Table || obj.obj_type == ObjectType::Index)
154 .map(|obj| obj.oid.as_table_id());
155
156 let removed_iceberg_table_sinks: Vec<PbSink> = removed_catalog_models
157 .iter()
158 .filter_map(|object_info| match object_info {
159 PbObjectInfo::Sink(sink) if sink.name.starts_with(ICEBERG_SINK_PREFIX) => {
160 Some(sink.clone())
161 }
162 _ => None,
163 })
164 .collect();
165
166 let mut removed_iceberg_sink_ids: Vec<SinkId> = Vec::new();
170 let mut removed_iceberg_pk_index_sink_ids: Vec<SinkId> = Vec::new();
171 for object_info in &removed_catalog_models {
172 if let PbObjectInfo::Sink(sink) = object_info {
173 if crate::manager::iceberg_compaction::is_iceberg_sink(&sink.properties) {
174 removed_iceberg_sink_ids.push(sink.id);
175 }
176 if crate::manager::iceberg_pk_index_sink::is_iceberg_pk_index_sink(&sink.properties)
177 {
178 removed_iceberg_pk_index_sink_ids.push(sink.id);
179 }
180 }
181 }
182
183 let removed_streaming_job_ids: Vec<JobId> = StreamingJob::find()
184 .select_only()
185 .column(streaming_job::Column::JobId)
186 .filter(streaming_job::Column::JobId.is_in(removed_object_ids))
187 .into_tuple()
188 .all(&txn)
189 .await?;
190
191 if !removed_streaming_job_ids.is_empty() {
193 let creating = StreamingJob::find()
194 .filter(
195 streaming_job::Column::JobStatus
196 .ne(JobStatus::Created)
197 .and(streaming_job::Column::JobId.is_in(removed_streaming_job_ids.clone())),
198 )
199 .count(&txn)
200 .await?;
201 if creating != 0 {
202 if creating == 1 && object_type == ObjectType::Sink {
203 info!("dropping creating sink job, it will be cancelled");
204 } else {
205 return Err(MetaError::permission_denied(format!(
206 "cannot drop {creating} streaming job(s) that are still being created; please cancel them first"
207 )));
208 }
209 }
210 }
211
212 let removed_state_table_ids: HashSet<_> = removed_table_ids.clone().collect();
213
214 let removed_source_ids: HashSet<_> = removed_objects
215 .iter()
216 .filter(|obj| obj.obj_type == ObjectType::Source)
217 .map(|obj| obj.oid.as_source_id())
218 .collect();
219
220 let removed_secret_ids = removed_objects
221 .iter()
222 .filter(|obj| obj.obj_type == ObjectType::Secret)
223 .map(|obj| obj.oid.as_secret_id())
224 .collect();
225
226 let removed_objects: HashMap<_, _> = removed_objects
227 .into_iter()
228 .map(|obj| (obj.oid, obj))
229 .collect();
230
231 for obj in removed_objects.values() {
233 if let Some(obj_database_id) = obj.database_id
234 && obj_database_id != database_id
235 {
236 return Err(MetaError::permission_denied(format!(
237 "Referenced by other objects in database {obj_database_id}, please drop them manually"
238 )));
239 }
240 }
241
242 let (removed_source_fragments, removed_sink_fragments, removed_fragments) =
243 get_fragments_for_jobs(&txn, removed_streaming_job_ids.clone()).await?;
244
245 let sink_target_fragments = fetch_target_fragments(&txn, removed_sink_fragments).await?;
246 let mut removed_sink_fragment_by_targets = HashMap::new();
247 for (sink_fragment, target_fragments) in sink_target_fragments {
248 assert!(
249 target_fragments.len() <= 1,
250 "sink should have at most one downstream fragment"
251 );
252 if let Some(target_fragment) = target_fragments.first()
253 && !removed_fragments.contains(target_fragment)
254 {
255 removed_sink_fragment_by_targets
256 .entry(*target_fragment)
257 .or_insert_with(Vec::new)
258 .push(sink_fragment);
259 }
260 }
261
262 let updated_user_ids: Vec<UserId> = UserPrivilege::find()
264 .select_only()
265 .distinct()
266 .column(user_privilege::Column::UserId)
267 .filter(user_privilege::Column::Oid.is_in(removed_objects.keys().cloned()))
268 .into_tuple()
269 .all(&txn)
270 .await?;
271 let dropped_tables = removed_catalog_models
272 .into_iter()
273 .filter_map(|object_info| match object_info {
274 PbObjectInfo::Table(table) if removed_state_table_ids.contains(&table.id) => {
275 Some(table)
276 }
277 _ => None,
278 });
279 let res = Object::delete_many()
282 .filter(object::Column::Oid.is_in(root_objects.iter().map(|obj| obj.oid)))
283 .exec(&txn)
284 .await?;
285 if res.rows_affected == 0 {
286 return Err(MetaError::catalog_id_not_found(
287 object_type.as_str(),
288 object_id,
289 ));
290 }
291 let user_infos = list_user_info_by_ids(updated_user_ids, &txn).await?;
292
293 txn.commit().await?;
294
295 self.notify_users_update(user_infos).await;
297 inner
298 .dropped_tables
299 .extend(dropped_tables.map(|t| (t.id, t)));
300
301 let version = match object_type {
302 ObjectType::Database => {
303 self.notify_frontend(
305 NotificationOperation::Delete,
306 NotificationInfo::Database(PbDatabase {
307 id: database_id,
308 ..Default::default()
309 }),
310 )
311 .await
312 }
313 ObjectType::Schema => {
314 let (schema_obj, mut to_notify_objs): (Vec<_>, Vec<_>) = removed_objects
315 .into_values()
316 .partition(|obj| obj.obj_type == ObjectType::Schema && obj.oid == object_id);
317 let schema_obj = Itertools::exactly_one(schema_obj.into_iter())
318 .expect("schema object not found");
319 to_notify_objs.push(schema_obj);
320
321 let relation_group = build_object_group_for_delete(to_notify_objs);
322 self.notify_frontend(NotificationOperation::Delete, relation_group)
323 .await
324 }
325 _ => {
326 let relation_group =
329 build_object_group_for_delete(removed_objects.into_values().collect());
330 self.notify_frontend(NotificationOperation::Delete, relation_group)
331 .await
332 }
333 };
334
335 Ok((
336 ReleaseContext {
337 database_id,
338 removed_streaming_job_ids,
339 removed_state_table_ids: removed_state_table_ids.into_iter().collect(),
340 removed_source_ids: removed_source_ids.into_iter().collect(),
341 removed_secret_ids,
342 removed_source_fragments,
343 removed_fragments,
344 removed_sink_fragment_by_targets,
345 removed_iceberg_table_sinks,
346 removed_iceberg_sink_ids,
347 removed_iceberg_pk_index_sink_ids,
348 },
349 version,
350 ))
351 }
352
353 pub async fn try_abort_creating_subscription(
354 &self,
355 subscription_id: SubscriptionId,
356 ) -> MetaResult<()> {
357 let inner = self.inner.write().await;
358 let txn = inner.db.begin().await?;
359
360 let subscription = Subscription::find_by_id(subscription_id).one(&txn).await?;
361 let Some(subscription) = subscription else {
362 tracing::warn!(
363 %subscription_id,
364 "subscription not found when aborting creation, might be cleaned by recovery"
365 );
366 return Ok(());
367 };
368
369 if subscription.subscription_state == PbSubscriptionState::Created as i32 {
370 tracing::warn!(
371 %subscription_id,
372 "subscription is already created when aborting creation"
373 );
374 return Ok(());
375 }
376
377 subscription.delete(&txn).await?;
378 txn.commit().await?;
379 Ok(())
380 }
381}
382
383async fn report_drop_object(
384 object_type: ObjectType,
385 object_id: ObjectId,
386 txn: &DatabaseTransaction,
387) {
388 let connector_name = {
389 match object_type {
390 ObjectType::Sink => Sink::find_by_id(object_id.as_sink_id())
391 .select_only()
392 .column(sink::Column::Properties)
393 .into_tuple::<Property>()
394 .one(txn)
395 .await
396 .ok()
397 .flatten()
398 .and_then(|properties| properties.inner_ref().get("connector").cloned()),
399 ObjectType::Source => Source::find_by_id(object_id.as_source_id())
400 .select_only()
401 .column(source::Column::WithProperties)
402 .into_tuple::<Property>()
403 .one(txn)
404 .await
405 .ok()
406 .flatten()
407 .and_then(|properties| properties.inner_ref().get("connector").cloned()),
408 _ => unreachable!(),
409 }
410 };
411 if let Some(connector_name) = connector_name {
412 report_event(
413 PbTelemetryEventStage::DropStreamJob,
414 "source",
415 object_id.as_raw_id() as _,
416 Some(connector_name),
417 Some(match object_type {
418 ObjectType::Source => PbTelemetryDatabaseObject::Source,
419 ObjectType::Sink => PbTelemetryDatabaseObject::Sink,
420 _ => unreachable!(),
421 }),
422 None,
423 );
424 }
425}