Skip to main content

risingwave_meta/controller/catalog/
drop_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::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    // Drop all kinds of objects including databases,
23    // schemas, relations, connections, functions, etc.
24    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        // Check the cross-db dependency info to see if the subscription can be dropped.
48        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        // Load all objects that belong to the dropped objects before deletion. Cascaded rows are
133        // still needed for notifications and resource cleanup.
134        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        // Iceberg sinks (and the pk-index subset) can be user-created with arbitrary
167        // names, so unlike the iceberg-table cleanup above, identify them by
168        // inspecting properties rather than by name prefix.
169        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        // Check if there are any streaming jobs that are creating.
192        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        // TODO: Support drop cascade for cross-database query.
232        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        // Find affect users with privileges on all this objects.
263        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        // Delete the explicitly selected objects. The self foreign key cascades to every object
280        // that belongs to them.
281        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        // notify about them.
296        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                // TODO: Notify objects in other databases when the cross-database query is supported.
304                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                // Hummock observers and compactor observers are notified once the corresponding barrier is completed.
327                // They only need RelationInfo::Table.
328                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}