Skip to main content

common_meta/ddl/
drop_table.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
15pub mod executor;
16mod metadata;
17
18use std::collections::HashMap;
19
20use async_trait::async_trait;
21use common_error::ext::BoxedError;
22use common_event_recorder::Event;
23use common_procedure::error::{ExternalSnafu, FromJsonSnafu, ToJsonSnafu};
24use common_procedure::{
25    Context as ProcedureContext, Error as ProcedureError, EventContext, EventTrigger, LockKey,
26    Procedure, Result as ProcedureResult, Status,
27};
28use common_telemetry::info;
29use common_telemetry::tracing::warn;
30#[cfg(feature = "enterprise")]
31use common_time::util::current_time_millis;
32use common_wal::options::WalOptions;
33use serde::{Deserialize, Serialize};
34use snafu::{OptionExt, ResultExt};
35use store_api::storage::RegionNumber;
36use strum::AsRefStr;
37use table::metadata::TableId;
38use table::table_reference::TableReference;
39#[cfg(feature = "enterprise")]
40use uuid::Uuid;
41
42use self::executor::DropTableExecutor;
43use crate::ddl::DdlContext;
44use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator};
45use crate::ddl::utils::{convert_region_routes_to_detecting_regions, map_to_procedure_error};
46use crate::error::{self, Result};
47use crate::key::table_route::TableRouteValue;
48use crate::lock_key::{CatalogLock, SchemaLock, TableLock, TableNameLock};
49use crate::metrics;
50use crate::region_keeper::OperatingRegionGuard;
51use crate::rpc::ddl::DropTableTask;
52use crate::rpc::router::{RegionRoute, operating_leader_region_roles};
53
54#[cfg(feature = "enterprise")]
55fn ensure_retry_later(err: error::Error) -> error::Error {
56    if err.is_retry_later() {
57        err
58    } else {
59        error::Error::retry_later(err)
60    }
61}
62
63pub struct DropTableProcedure {
64    /// The context of procedure runtime.
65    pub context: DdlContext,
66    /// The serializable data.
67    pub data: DropTableData,
68    /// The guards of opening regions.
69    pub(crate) dropping_regions: Vec<OperatingRegionGuard>,
70    /// The drop table executor.
71    executor: DropTableExecutor,
72}
73
74impl DropTableProcedure {
75    pub const TYPE_NAME: &'static str = "metasrv-procedure::DropTable";
76
77    pub fn new(task: DropTableTask, context: DdlContext) -> Self {
78        let data = DropTableData::new(
79            task,
80            cfg!(feature = "enterprise") && context.soft_drop_enabled,
81            context.soft_drop_retention,
82        );
83        let executor = data.build_executor();
84        Self {
85            context,
86            data,
87            dropping_regions: vec![],
88            executor,
89        }
90    }
91
92    pub fn from_json(json: &str, context: DdlContext) -> ProcedureResult<Self> {
93        let data: DropTableData = serde_json::from_str(json).context(FromJsonSnafu)?;
94        #[cfg(feature = "enterprise")]
95        let mut data = data;
96        #[cfg(feature = "enterprise")]
97        if data.state == DropTableState::Prepare
98            && data.soft_drop_enabled
99            && data.dropped_at.is_none()
100            && data.soft_drop_retention_millis.is_none()
101        {
102            data.soft_drop_retention_millis = context
103                .soft_drop_retention
104                .and_then(|retention| i64::try_from(retention.as_millis()).ok());
105        }
106        let executor = data.build_executor();
107
108        Ok(Self {
109            context,
110            data,
111            dropping_regions: vec![],
112            executor,
113        })
114    }
115
116    #[cfg(feature = "enterprise")]
117    async fn prepare_soft_drop(&mut self) -> Result<()> {
118        self.executor
119            .check_tombstone_conflict(&self.context, self.data.soft_drop_enabled)
120            .await?;
121        if self.data.soft_drop_enabled && self.data.dropped_at.is_none() {
122            let dropped_at = current_time_millis();
123            let retention_millis =
124                self.data
125                    .soft_drop_retention_millis
126                    .context(error::UnexpectedSnafu {
127                        err_msg: "Soft-drop retention is missing from the procedure".to_string(),
128                    })?;
129            let retention_expires_at = dropped_at.checked_add(retention_millis).context(
130                error::UnexpectedSnafu {
131                    err_msg: format!(
132                        "Soft-drop retention deadline overflows i64: dropped_at={dropped_at}, retention_millis={retention_millis}"
133                    ),
134                },
135            )?;
136            self.data.dropped_at = Some(dropped_at);
137            self.data.retention_expires_at = Some(retention_expires_at);
138        }
139        if self.data.soft_drop_enabled && self.data.drop_generation.is_none() {
140            self.data.drop_generation = Some(Uuid::new_v4().to_string());
141        }
142        Ok(())
143    }
144
145    pub(crate) async fn on_prepare(&mut self) -> Result<Status> {
146        if self.executor.on_prepare(&self.context).await?.stop() {
147            return Ok(Status::done());
148        }
149        self.fill_table_metadata().await?;
150        #[cfg(feature = "enterprise")]
151        self.prepare_soft_drop().await?;
152        self.data.state = DropTableState::DeleteMetadata;
153
154        Ok(Status::executing(true))
155    }
156
157    /// Register dropping regions if doesn't exist.
158    fn register_dropping_regions(&mut self) -> Result<()> {
159        let dropping_regions = operating_leader_region_roles(&self.data.physical_region_routes);
160
161        if !self.dropping_regions.is_empty() {
162            return Ok(());
163        }
164
165        let mut dropping_region_guards = Vec::with_capacity(dropping_regions.len());
166
167        for (region_id, datanode_id, role) in dropping_regions {
168            let guard = self
169                .context
170                .memory_region_keeper
171                .register_with_role(datanode_id, region_id, role)
172                .context(error::RegionOperatingRaceSnafu {
173                    region_id,
174                    peer_id: datanode_id,
175                })?;
176            dropping_region_guards.push(guard);
177        }
178
179        self.dropping_regions = dropping_region_guards;
180        Ok(())
181    }
182
183    #[cfg(not(feature = "enterprise"))]
184    async fn delete_metadata(&mut self) -> Result<()> {
185        let table_route_value = &TableRouteValue::new(
186            self.data.task.table_id,
187            // Safety: checked
188            self.data.physical_table_id.unwrap(),
189            self.data.physical_region_routes.clone(),
190        );
191        self.executor
192            .on_delete_metadata(
193                &self.context,
194                table_route_value,
195                &self.data.region_wal_options,
196            )
197            .await
198    }
199
200    #[cfg(feature = "enterprise")]
201    async fn delete_metadata(&mut self) -> Result<()> {
202        if !self.data.soft_drop_enabled {
203            let table_route_value = &TableRouteValue::new(
204                self.data.task.table_id,
205                // Safety: checked
206                self.data.physical_table_id.unwrap(),
207                self.data.physical_region_routes.clone(),
208            );
209            return self
210                .executor
211                .on_delete_metadata(
212                    &self.context,
213                    table_route_value,
214                    &self.data.region_wal_options,
215                )
216                .await;
217        }
218
219        let storage = self
220            .context
221            .table_metadata_manager
222            .table_route_manager()
223            .table_route_storage();
224        storage
225            .remap_region_routes(&mut self.data.physical_region_routes)
226            .await
227            .map_err(ensure_retry_later)?;
228        let table_route_value = &TableRouteValue::new(
229            self.data.task.table_id,
230            // Safety: checked
231            self.data.physical_table_id.unwrap(),
232            self.data.physical_region_routes.clone(),
233        );
234        self.executor
235            .on_close_regions(
236                &self.context.node_manager,
237                &self.context.leader_region_registry,
238                &self.data.physical_region_routes,
239                true,
240            )
241            .await
242            .map_err(ensure_retry_later)?;
243        self.executor
244            .on_soft_delete_metadata(
245                &self.context,
246                table_route_value,
247                &self.data.region_wal_options,
248                self.data.dropped_at,
249                self.data.retention_expires_at,
250                self.data.drop_generation.as_deref(),
251            )
252            .await
253            .map_err(ensure_retry_later)?;
254        self.data.allow_rollback = false;
255        Ok(())
256    }
257
258    pub(crate) async fn on_delete_metadata(&mut self) -> Result<Status> {
259        self.register_dropping_regions()?;
260        // TODO(weny): Considers introducing a RegionStatus to indicate the region is dropping.
261        let table_id = self.data.table_id();
262        self.delete_metadata().await?;
263        info!("Deleted table metadata for table {table_id}");
264        self.data.state = DropTableState::InvalidateTableCache;
265        Ok(Status::executing(true))
266    }
267
268    /// Broadcasts invalidate table cache instruction.
269    async fn on_broadcast(&mut self) -> Result<Status> {
270        let result = self.executor.invalidate_table_cache(&self.context).await;
271        #[cfg(feature = "enterprise")]
272        if self.data.soft_drop_enabled {
273            result.map_err(ensure_retry_later)?;
274        } else {
275            result?;
276        }
277        #[cfg(not(feature = "enterprise"))]
278        result?;
279
280        self.data.state = DropTableState::DatanodeDropRegions;
281
282        Ok(Status::executing(true))
283    }
284
285    pub async fn on_datanode_drop_regions(&mut self, retrying: bool) -> Result<Status> {
286        if retrying {
287            info!(
288                "Remapping region routes addresses for retrying drop regions for table_id: {}",
289                self.data.table_id()
290            );
291            let storage = self
292                .context
293                .table_metadata_manager
294                .table_route_manager()
295                .table_route_storage();
296            // The peer addresses may change during retries,
297            // so we always remap the region routes.
298            storage
299                .remap_region_routes(&mut self.data.physical_region_routes)
300                .await?;
301        }
302
303        #[cfg(feature = "enterprise")]
304        if self.data.soft_drop_enabled {
305            self.context
306                .deregister_failure_detectors(convert_region_routes_to_detecting_regions(
307                    &self.data.physical_region_routes,
308                ))
309                .await;
310            self.dropping_regions.clear();
311            return Ok(Status::done());
312        }
313
314        self.executor
315            .on_drop_regions(
316                &self.context.node_manager,
317                &self.context.leader_region_registry,
318                &self.data.physical_region_routes,
319                false,
320                false,
321                false,
322            )
323            .await?;
324        self.context
325            .deregister_failure_detectors(convert_region_routes_to_detecting_regions(
326                &self.data.physical_region_routes,
327            ))
328            .await;
329
330        self.data.state = DropTableState::DeleteTombstone;
331        Ok(Status::executing(true))
332    }
333
334    /// Deletes metadata tombstone.
335    async fn on_delete_metadata_tombstone(&mut self) -> Result<Status> {
336        let table_route_value = &TableRouteValue::new(
337            self.data.task.table_id,
338            // Safety: checked
339            self.data.physical_table_id.unwrap(),
340            self.data.physical_region_routes.clone(),
341        );
342        self.executor
343            .on_delete_metadata_tombstone(
344                &self.context,
345                table_route_value,
346                &self.data.region_wal_options,
347            )
348            .await?;
349
350        self.dropping_regions.clear();
351        Ok(Status::done())
352    }
353}
354
355#[async_trait]
356impl Procedure for DropTableProcedure {
357    fn type_name(&self) -> &str {
358        Self::TYPE_NAME
359    }
360
361    fn recover(&mut self) -> ProcedureResult<()> {
362        // Only registers regions if the metadata is deleted.
363        let register_operating_regions = matches!(
364            self.data.state,
365            DropTableState::DeleteMetadata
366                | DropTableState::InvalidateTableCache
367                | DropTableState::DatanodeDropRegions
368        );
369        if register_operating_regions {
370            self.register_dropping_regions()
371                .map_err(BoxedError::new)
372                .context(ExternalSnafu {
373                    clean_poisons: false,
374                })?;
375        }
376
377        Ok(())
378    }
379
380    async fn execute(&mut self, ctx: &ProcedureContext) -> ProcedureResult<Status> {
381        let state = &self.data.state;
382        let _timer = metrics::METRIC_META_PROCEDURE_DROP_TABLE
383            .with_label_values(&[state.as_ref()])
384            .start_timer();
385
386        match self.data.state {
387            DropTableState::Prepare => self.on_prepare().await,
388            DropTableState::DeleteMetadata => self.on_delete_metadata().await,
389            DropTableState::InvalidateTableCache => self.on_broadcast().await,
390            DropTableState::DatanodeDropRegions => {
391                let retrying = ctx.is_retrying().await.unwrap_or(false);
392                self.on_datanode_drop_regions(retrying).await
393            }
394            DropTableState::DeleteTombstone => self.on_delete_metadata_tombstone().await,
395        }
396        .map_err(map_to_procedure_error)
397    }
398
399    fn dump(&self) -> ProcedureResult<String> {
400        serde_json::to_string(&self.data).context(ToJsonSnafu)
401    }
402
403    fn lock_key(&self) -> LockKey {
404        let table_ref = &self.data.table_ref();
405        let table_id = self.data.table_id();
406        let lock_key = vec![
407            CatalogLock::Read(table_ref.catalog).into(),
408            SchemaLock::read(table_ref.catalog, table_ref.schema).into(),
409            TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(),
410            TableLock::Write(table_id).into(),
411        ];
412
413        LockKey::new(lock_key)
414    }
415
416    fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn Event>> {
417        if !ctx
418            .event_type_filter
419            .allows(TableDdlEventType::DropTable.as_str())
420        {
421            return None;
422        }
423        let task = &self.data.task;
424        let locator = TableDdlLocator::new(&task.catalog, &task.schema, &task.table)
425            .with_table_id(task.table_id);
426        let event = match &ctx.trigger {
427            EventTrigger::Submitted => {
428                TableDdlEvent::drop_table_submitted(locator, task.drop_if_exists)
429            }
430            _ => TableDdlEvent::lifecycle(TableDdlEventType::DropTable, [locator]),
431        };
432
433        Some(Box::new(event))
434    }
435
436    fn rollback_supported(&self) -> bool {
437        !matches!(self.data.state, DropTableState::Prepare) && self.data.allow_rollback
438    }
439
440    async fn rollback(&mut self, _: &ProcedureContext) -> ProcedureResult<()> {
441        warn!(
442            "Rolling back the drop table procedure, table: {}",
443            self.data.table_id()
444        );
445
446        let table_route_value = &TableRouteValue::new(
447            self.data.task.table_id,
448            // Safety: checked
449            self.data.physical_table_id.unwrap(),
450            self.data.physical_region_routes.clone(),
451        );
452        self.executor
453            .on_restore_metadata(
454                &self.context,
455                table_route_value,
456                &self.data.region_wal_options,
457            )
458            .await
459            .map_err(ProcedureError::external)?;
460
461        self.dropping_regions.clear();
462        Ok(())
463    }
464}
465
466#[derive(Debug, Serialize, Deserialize)]
467pub struct DropTableData {
468    pub state: DropTableState,
469    pub task: DropTableTask,
470    pub physical_region_routes: Vec<RegionRoute>,
471    pub physical_table_id: Option<TableId>,
472    #[serde(default)]
473    pub region_wal_options: HashMap<RegionNumber, WalOptions>,
474    #[serde(default)]
475    pub allow_rollback: bool,
476    #[serde(default)]
477    pub soft_drop_enabled: bool,
478    #[serde(default)]
479    pub dropped_at: Option<i64>,
480    #[serde(default)]
481    pub retention_expires_at: Option<i64>,
482    #[serde(default)]
483    pub soft_drop_retention_millis: Option<i64>,
484    #[serde(default)]
485    pub drop_generation: Option<String>,
486}
487
488impl DropTableData {
489    pub fn new(
490        task: DropTableTask,
491        soft_drop_enabled: bool,
492        soft_drop_retention: Option<std::time::Duration>,
493    ) -> Self {
494        Self {
495            state: DropTableState::Prepare,
496            task,
497            physical_region_routes: vec![],
498            physical_table_id: None,
499            region_wal_options: HashMap::new(),
500            allow_rollback: false,
501            soft_drop_enabled,
502            dropped_at: None,
503            retention_expires_at: None,
504            soft_drop_retention_millis: soft_drop_retention_millis(soft_drop_retention),
505            drop_generation: None,
506        }
507    }
508
509    fn table_ref(&self) -> TableReference<'_> {
510        self.task.table_ref()
511    }
512
513    fn table_id(&self) -> TableId {
514        self.task.table_id
515    }
516
517    fn build_executor(&self) -> DropTableExecutor {
518        DropTableExecutor::new(
519            self.task.table_name(),
520            self.task.table_id,
521            self.task.drop_if_exists,
522        )
523    }
524}
525
526#[cfg(feature = "enterprise")]
527fn soft_drop_retention_millis(soft_drop_retention: Option<std::time::Duration>) -> Option<i64> {
528    soft_drop_retention.and_then(|retention| i64::try_from(retention.as_millis()).ok())
529}
530
531#[cfg(not(feature = "enterprise"))]
532fn soft_drop_retention_millis(_: Option<std::time::Duration>) -> Option<i64> {
533    None
534}
535
536/// The state of drop table.
537#[derive(Debug, Serialize, Deserialize, AsRefStr, PartialEq)]
538pub enum DropTableState {
539    /// Prepares to drop the table
540    Prepare,
541    /// Deletes metadata logically
542    DeleteMetadata,
543    /// Invalidates Table Cache
544    InvalidateTableCache,
545    /// Drops regions on Datanode
546    DatanodeDropRegions,
547    /// Deletes metadata tombstone permanently
548    DeleteTombstone,
549}