Skip to main content

common_meta/ddl/drop_table/
executor.rs

1// Copyright 2023 Greptime Team
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 std::collections::HashMap;
16#[cfg(feature = "enterprise")]
17use std::collections::HashSet;
18
19use api::v1::region::{
20    CleanUpRequest as PbCleanUpRequest, CloseRequest as PbCloseRegionRequest,
21    DropRequest as PbDropRegionRequest, RegionRequest, RegionRequestHeader, region_request,
22};
23use common_error::ext::ErrorExt;
24use common_error::status_code::StatusCode;
25use common_telemetry::tracing_context::TracingContext;
26use common_telemetry::{debug, error};
27use common_wal::options::WalOptions;
28use futures::future::join_all;
29use snafu::ensure;
30use store_api::storage::{RegionId, RegionNumber};
31use table::metadata::{TableId, TableInfo};
32use table::table_name::TableName;
33
34use crate::cache_invalidator::Context;
35use crate::ddl::utils::{
36    add_peer_context_if_needed, convert_region_routes_to_detecting_regions, region_storage_path,
37};
38use crate::ddl::{CreateRequestBuilder, DdlContext, build_template_from_raw_table_info};
39use crate::error::{self, Result};
40use crate::instruction::CacheIdent;
41#[cfg(feature = "enterprise")]
42use crate::key::DroppedTableLifecycle;
43use crate::key::table_name::TableNameKey;
44use crate::key::table_route::TableRouteValue;
45use crate::node_manager::NodeManagerRef;
46use crate::region_registry::LeaderRegionRegistryRef;
47use crate::rpc::router::{
48    RegionRoute, find_follower_regions, find_followers, find_leader_regions, find_leaders,
49    operating_leader_regions,
50};
51
52/// [Control] indicated to the caller whether to go to the next step.
53#[derive(Debug)]
54pub enum Control<T> {
55    Continue(T),
56    Stop,
57}
58
59impl<T> Control<T> {
60    /// Returns true if it's [Control::Stop].
61    pub fn stop(&self) -> bool {
62        matches!(self, Control::Stop)
63    }
64}
65
66impl DropTableExecutor {
67    /// Returns the [DropTableExecutor].
68    pub fn new(table: TableName, table_id: TableId, drop_if_exists: bool) -> Self {
69        Self {
70            table,
71            table_id,
72            drop_if_exists,
73        }
74    }
75}
76
77/// [DropTableExecutor] performs:
78/// - Drops the metadata of the table.
79/// - Invalidates the cache on the Frontend nodes.
80/// - Drops the regions on the Datanode nodes.
81pub struct DropTableExecutor {
82    table: TableName,
83    table_id: TableId,
84    drop_if_exists: bool,
85}
86
87impl DropTableExecutor {
88    /// Checks whether table exists.
89    /// - Early returns if table not exists and `drop_if_exists` is `true`.
90    /// - Throws an error if table not exists and `drop_if_exists` is `false`.
91    pub async fn on_prepare(&self, ctx: &DdlContext) -> Result<Control<()>> {
92        let table_ref = self.table.table_ref();
93
94        let exist = ctx
95            .table_metadata_manager
96            .table_name_manager()
97            .exists(TableNameKey::new(
98                table_ref.catalog,
99                table_ref.schema,
100                table_ref.table,
101            ))
102            .await?;
103
104        if !exist && self.drop_if_exists {
105            return Ok(Control::Stop);
106        }
107
108        ensure!(
109            exist,
110            error::TableNotFoundSnafu {
111                table_name: table_ref.to_string()
112            }
113        );
114
115        Ok(Control::Continue(()))
116    }
117
118    /// Rejects dropping a recreated live table while an older tombstone still owns the same
119    /// fully qualified name.
120    #[cfg(feature = "enterprise")]
121    pub async fn check_tombstone_conflict(
122        &self,
123        ctx: &DdlContext,
124        soft_drop_enabled: bool,
125    ) -> Result<()> {
126        let table_ref = self.table.table_ref();
127        if let Some(dropped_table) = ctx
128            .table_metadata_manager
129            .get_dropped_table(&self.table)
130            .await?
131            && dropped_table.table_id != self.table_id
132            && (soft_drop_enabled || dropped_table.dropped_at.is_some())
133        {
134            return error::TableNameTombstoneConflictSnafu {
135                table_name: table_ref.to_string(),
136                existing_table_id: dropped_table.table_id,
137                dropping_table_id: self.table_id,
138            }
139            .fail();
140        }
141
142        Ok(())
143    }
144
145    /// Deletes the table metadata **logically**.
146    pub async fn on_delete_metadata(
147        &self,
148        ctx: &DdlContext,
149        table_route_value: &TableRouteValue,
150        region_wal_options: &HashMap<RegionNumber, WalOptions>,
151    ) -> Result<()> {
152        ctx.table_metadata_manager
153            .delete_table_metadata(
154                self.table_id,
155                &self.table,
156                table_route_value,
157                region_wal_options,
158                None,
159            )
160            .await
161    }
162
163    /// Soft-deletes table metadata and retains its tombstone lifecycle markers.
164    #[cfg(feature = "enterprise")]
165    pub async fn on_soft_delete_metadata(
166        &self,
167        ctx: &DdlContext,
168        table_route_value: &TableRouteValue,
169        region_wal_options: &HashMap<RegionNumber, WalOptions>,
170        dropped_at: Option<i64>,
171        retention_expires_at: Option<i64>,
172        drop_generation: Option<&str>,
173    ) -> Result<()> {
174        ctx.table_metadata_manager
175            .delete_table_metadata_with_retention_and_generation(
176                self.table_id,
177                &self.table,
178                table_route_value,
179                region_wal_options,
180                DroppedTableLifecycle {
181                    dropped_at,
182                    retention_expires_at,
183                    drop_generation,
184                },
185            )
186            .await
187    }
188
189    /// Deletes the table metadata tombstone **permanently**.
190    pub async fn on_delete_metadata_tombstone(
191        &self,
192        ctx: &DdlContext,
193        table_route_value: &TableRouteValue,
194        region_wal_options: &HashMap<u32, WalOptions>,
195    ) -> Result<()> {
196        ctx.table_metadata_manager
197            .delete_table_metadata_tombstone(
198                self.table_id,
199                &self.table,
200                table_route_value,
201                region_wal_options,
202            )
203            .await
204    }
205
206    /// Deletes metadata for table **permanently**.
207    pub async fn on_destroy_metadata(
208        &self,
209        ctx: &DdlContext,
210        table_route_value: &TableRouteValue,
211        region_wal_options: &HashMap<u32, WalOptions>,
212    ) -> Result<()> {
213        ctx.table_metadata_manager
214            .destroy_table_metadata(
215                self.table_id,
216                &self.table,
217                table_route_value,
218                region_wal_options,
219            )
220            .await?;
221
222        let detecting_regions = if table_route_value.is_physical() {
223            // Safety: checked.
224            let regions = table_route_value.region_routes().unwrap();
225            convert_region_routes_to_detecting_regions(regions)
226        } else {
227            vec![]
228        };
229        ctx.deregister_failure_detectors(detecting_regions).await;
230        Ok(())
231    }
232
233    /// Restores the table metadata.
234    pub async fn on_restore_metadata(
235        &self,
236        ctx: &DdlContext,
237        table_route_value: &TableRouteValue,
238        region_wal_options: &HashMap<u32, WalOptions>,
239    ) -> Result<()> {
240        ctx.table_metadata_manager
241            .restore_table_metadata(
242                self.table_id,
243                &self.table,
244                table_route_value,
245                region_wal_options,
246            )
247            .await
248    }
249
250    /// Invalidates caches for the table.
251    pub async fn invalidate_table_cache(&self, ctx: &DdlContext) -> Result<()> {
252        let cache_invalidator = &ctx.cache_invalidator;
253        let ctx = Context {
254            subject: Some(format!(
255                "Invalidate table cache by dropping table {}, table_id: {}",
256                self.table.table_ref(),
257                self.table_id,
258            )),
259        };
260
261        cache_invalidator
262            .invalidate(
263                &ctx,
264                &[
265                    CacheIdent::TableName(self.table.table_ref().into()),
266                    CacheIdent::TableId(self.table_id),
267                ],
268            )
269            .await?;
270
271        Ok(())
272    }
273
274    /// Drops regions on datanodes.
275    ///
276    /// Arguments:
277    /// - `node_manager`: resolves datanode clients from peers in `region_routes`.
278    /// - `leader_region_registry`: tracks in-flight leader region operations.
279    /// - `region_routes`: table region placement; leaders receive drop requests and followers
280    ///   receive close requests.
281    /// - `fast_path`: forwards to datanode drop requests to skip extra cleanup when safe.
282    /// - `force`: forwards to datanode drop requests to allow forced region removal.
283    /// - `partial_drop`: forwards to datanode drop requests for partial table/region drops.
284    #[allow(clippy::too_many_arguments)]
285    pub async fn on_drop_regions(
286        &self,
287        node_manager: &NodeManagerRef,
288        leader_region_registry: &LeaderRegionRegistryRef,
289        region_routes: &[RegionRoute],
290        fast_path: bool,
291        force: bool,
292        partial_drop: bool,
293    ) -> Result<()> {
294        // Drops leader regions on datanodes.
295        let leaders = find_leaders(region_routes);
296        let mut drop_region_tasks = Vec::with_capacity(leaders.len());
297        let table_id = self.table_id;
298        for datanode in leaders {
299            let requester = node_manager.datanode(&datanode).await;
300            let regions = find_leader_regions(region_routes, &datanode);
301            let region_ids = regions
302                .iter()
303                .map(|region_number| RegionId::new(table_id, *region_number))
304                .collect::<Vec<_>>();
305
306            for region_id in region_ids {
307                debug!("Dropping region {region_id} on Datanode {datanode:?}");
308                let request = RegionRequest {
309                    header: Some(RegionRequestHeader {
310                        tracing_context: TracingContext::from_current_span().to_w3c(),
311                        ..Default::default()
312                    }),
313                    body: Some(region_request::Body::Drop(PbDropRegionRequest {
314                        region_id: region_id.as_u64(),
315                        fast_path,
316                        force,
317                        partial_drop,
318                        soft_drop: false,
319                    })),
320                };
321                let datanode = datanode.clone();
322                let requester = requester.clone();
323                drop_region_tasks.push(async move {
324                    if let Err(err) = requester.handle(request).await
325                        && err.status_code() != StatusCode::RegionNotFound
326                    {
327                        return Err(add_peer_context_if_needed(datanode)(err));
328                    }
329                    Ok(())
330                });
331            }
332        }
333
334        join_all(drop_region_tasks)
335            .await
336            .into_iter()
337            .collect::<Result<Vec<_>>>()?;
338
339        // Drops follower regions on datanodes.
340        let followers = find_followers(region_routes);
341        let mut close_region_tasks = Vec::with_capacity(followers.len());
342        for datanode in followers {
343            let requester = node_manager.datanode(&datanode).await;
344            let regions = find_follower_regions(region_routes, &datanode);
345            let region_ids = regions
346                .iter()
347                .map(|region_number| RegionId::new(table_id, *region_number))
348                .collect::<Vec<_>>();
349
350            for region_id in region_ids {
351                debug!("Closing region {region_id} on Datanode {datanode:?}");
352                let request = RegionRequest {
353                    header: Some(RegionRequestHeader {
354                        tracing_context: TracingContext::from_current_span().to_w3c(),
355                        ..Default::default()
356                    }),
357                    body: Some(region_request::Body::Close(PbCloseRegionRequest {
358                        region_id: region_id.as_u64(),
359                        flush_on_close: false,
360                    })),
361                };
362
363                let datanode = datanode.clone();
364                let requester = requester.clone();
365                close_region_tasks.push(async move {
366                    if let Err(err) = requester.handle(request).await
367                        && err.status_code() != StatusCode::RegionNotFound
368                    {
369                        return Err(add_peer_context_if_needed(datanode)(err));
370                    }
371                    Ok(())
372                });
373            }
374        }
375
376        // Failure to close follower regions is not critical.
377        // When a leader region is dropped, follower regions will be unable to renew their leases via metasrv.
378        // Eventually, these follower regions will be automatically closed by the region livekeeper.
379        if let Err(err) = join_all(close_region_tasks)
380            .await
381            .into_iter()
382            .collect::<Result<Vec<_>>>()
383        {
384            error!(err; "Failed to close follower regions on datanodes, table_id: {}", table_id);
385        }
386
387        // Deletes the leader region from registry.
388        let region_ids = operating_leader_regions(region_routes);
389        leader_region_registry.batch_delete(region_ids.into_iter().map(|(region_id, _)| region_id));
390
391        Ok(())
392    }
393
394    /// Cleans leader regions on datanodes without reopening them as live regions.
395    pub async fn on_cleanup_regions_offline(
396        &self,
397        node_manager: &NodeManagerRef,
398        leader_region_registry: &LeaderRegionRegistryRef,
399        table_info: &TableInfo,
400        region_routes: &[RegionRoute],
401        region_wal_options: &HashMap<RegionNumber, WalOptions>,
402    ) -> Result<()> {
403        let template = build_template_from_raw_table_info(table_info)?;
404        let builder = CreateRequestBuilder::new(template, None);
405        let storage_path = region_storage_path(&self.table.catalog_name, &self.table.schema_name);
406
407        let leaders = find_leaders(region_routes);
408        let mut cleanup_region_tasks = Vec::with_capacity(leaders.len());
409        let table_id = self.table_id;
410        for datanode in leaders {
411            let requester = node_manager.datanode(&datanode).await;
412            let regions = find_leader_regions(region_routes, &datanode);
413            let region_ids = regions
414                .iter()
415                .map(|region_number| RegionId::new(table_id, *region_number))
416                .collect::<Vec<_>>();
417
418            for region_id in region_ids {
419                debug!("Cleaning region {region_id} offline on Datanode {datanode:?}");
420                let create_request = builder.build_one(
421                    region_id,
422                    storage_path.clone(),
423                    region_wal_options,
424                    &HashMap::new(),
425                )?;
426                let request = RegionRequest {
427                    header: Some(RegionRequestHeader {
428                        tracing_context: TracingContext::from_current_span().to_w3c(),
429                        ..Default::default()
430                    }),
431                    body: Some(region_request::Body::CleanUp(PbCleanUpRequest {
432                        region_id: create_request.region_id,
433                        engine: create_request.engine,
434                        path: create_request.path,
435                        options: create_request.options,
436                    })),
437                };
438                let datanode = datanode.clone();
439                let requester = requester.clone();
440                cleanup_region_tasks.push(async move {
441                    if let Err(err) = requester.handle(request).await
442                        && err.status_code() != StatusCode::RegionNotFound
443                    {
444                        return Err(add_peer_context_if_needed(datanode)(err));
445                    }
446                    Ok(())
447                });
448            }
449        }
450
451        join_all(cleanup_region_tasks)
452            .await
453            .into_iter()
454            .collect::<Result<Vec<_>>>()?;
455
456        let region_ids = operating_leader_regions(region_routes);
457        leader_region_registry.batch_delete(region_ids.into_iter().map(|(region_id, _)| region_id));
458
459        Ok(())
460    }
461
462    /// Closes all table regions on datanodes without deleting region files or metadata tombstones.
463    /// When `flush_leaders_on_close` is set, only leader regions are flushed before close.
464    #[cfg(feature = "enterprise")]
465    pub async fn on_close_regions(
466        &self,
467        node_manager: &NodeManagerRef,
468        leader_region_registry: &LeaderRegionRegistryRef,
469        region_routes: &[RegionRoute],
470        flush_leaders_on_close: bool,
471    ) -> Result<()> {
472        let table_id = self.table_id;
473        let mut seen_peer_ids = HashSet::new();
474        let peers = find_leaders(region_routes)
475            .into_iter()
476            .chain(find_followers(region_routes))
477            .filter(|peer| seen_peer_ids.insert(peer.id));
478        let close_region_tasks = peers.map(|datanode| {
479            let region_ids = find_leader_regions(region_routes, &datanode)
480                .into_iter()
481                .map(|region_number| {
482                    (
483                        RegionId::new(table_id, region_number),
484                        flush_leaders_on_close,
485                    )
486                })
487                .chain(
488                    find_follower_regions(region_routes, &datanode)
489                        .into_iter()
490                        .map(|region_number| (RegionId::new(table_id, region_number), false)),
491                )
492                .collect::<Vec<_>>();
493
494            async move {
495                let requester = node_manager.datanode(&datanode).await;
496                let close_region_tasks =
497                    region_ids.into_iter().map(|(region_id, flush_on_close)| {
498                        debug!("Closing region {region_id} on Datanode {datanode:?}");
499                        let request = RegionRequest {
500                            header: Some(RegionRequestHeader {
501                                tracing_context: TracingContext::from_current_span().to_w3c(),
502                                ..Default::default()
503                            }),
504                            body: Some(region_request::Body::Close(PbCloseRegionRequest {
505                                region_id: region_id.as_u64(),
506                                flush_on_close,
507                            })),
508                        };
509
510                        let datanode = datanode.clone();
511                        let requester = requester.clone();
512                        async move {
513                            if let Err(err) = requester.handle(request).await
514                                && err.status_code() != StatusCode::RegionNotFound
515                            {
516                                return Err(add_peer_context_if_needed(datanode)(err));
517                            }
518                            Ok(())
519                        }
520                    });
521
522                join_all(close_region_tasks)
523                    .await
524                    .into_iter()
525                    .collect::<Result<Vec<_>>>()?;
526                Ok(())
527            }
528        });
529
530        join_all(close_region_tasks)
531            .await
532            .into_iter()
533            .collect::<Result<Vec<_>>>()?;
534
535        let region_ids = operating_leader_regions(region_routes);
536        leader_region_registry.batch_delete(region_ids.into_iter().map(|(region_id, _)| region_id));
537
538        Ok(())
539    }
540}
541
542#[cfg(test)]
543mod tests {
544    use std::assert_matches;
545    use std::collections::HashMap;
546    use std::sync::Arc;
547
548    use api::v1::{ColumnDataType, SemanticType};
549    use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
550    use table::metadata::TableInfo;
551    use table::table_name::TableName;
552
553    use super::*;
554    use crate::ddl::test_util::columns::TestColumnDefBuilder;
555    use crate::ddl::test_util::create_table::{
556        TestCreateTableExprBuilder, build_raw_table_info_from_expr,
557    };
558    use crate::key::table_route::TableRouteValue;
559    use crate::test_util::{MockDatanodeManager, new_ddl_context};
560
561    fn test_create_raw_table_info(name: &str) -> TableInfo {
562        let create_table = TestCreateTableExprBuilder::default()
563            .column_defs([
564                TestColumnDefBuilder::default()
565                    .name("ts")
566                    .data_type(ColumnDataType::TimestampMillisecond)
567                    .semantic_type(SemanticType::Timestamp)
568                    .build()
569                    .unwrap()
570                    .into(),
571                TestColumnDefBuilder::default()
572                    .name("host")
573                    .data_type(ColumnDataType::String)
574                    .semantic_type(SemanticType::Tag)
575                    .build()
576                    .unwrap()
577                    .into(),
578                TestColumnDefBuilder::default()
579                    .name("cpu")
580                    .data_type(ColumnDataType::Float64)
581                    .semantic_type(SemanticType::Field)
582                    .build()
583                    .unwrap()
584                    .into(),
585            ])
586            .time_index("ts")
587            .primary_keys(["host".into()])
588            .table_name(name)
589            .build()
590            .unwrap()
591            .into();
592        build_raw_table_info_from_expr(&create_table)
593    }
594
595    #[tokio::test]
596    async fn test_on_prepare() {
597        // Drops if exists
598        let node_manager = Arc::new(MockDatanodeManager::new(()));
599        let ctx = new_ddl_context(node_manager);
600        let executor = DropTableExecutor::new(
601            TableName::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "my_table"),
602            1024,
603            true,
604        );
605        let ctrl = executor.on_prepare(&ctx).await.unwrap();
606        assert!(ctrl.stop());
607
608        // Drops a non-exists table
609        let executor = DropTableExecutor::new(
610            TableName::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "my_table"),
611            1024,
612            false,
613        );
614        let err = executor.on_prepare(&ctx).await.unwrap_err();
615        assert_matches!(err, error::Error::TableNotFound { .. });
616
617        // Drops a exists table
618        let executor = DropTableExecutor::new(
619            TableName::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "my_table"),
620            1024,
621            false,
622        );
623        let raw_table_info = test_create_raw_table_info("my_table");
624        ctx.table_metadata_manager
625            .create_table_metadata(
626                raw_table_info,
627                TableRouteValue::physical(vec![]),
628                HashMap::new(),
629            )
630            .await
631            .unwrap();
632        let ctrl = executor.on_prepare(&ctx).await.unwrap();
633        assert!(!ctrl.stop());
634    }
635}