1pub 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 pub context: DdlContext,
66 pub data: DropTableData,
68 pub(crate) dropping_regions: Vec<OperatingRegionGuard>,
70 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 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 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 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 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 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 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 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 async fn on_delete_metadata_tombstone(&mut self) -> Result<Status> {
336 let table_route_value = &TableRouteValue::new(
337 self.data.task.table_id,
338 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 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 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#[derive(Debug, Serialize, Deserialize, AsRefStr, PartialEq)]
538pub enum DropTableState {
539 Prepare,
541 DeleteMetadata,
543 InvalidateTableCache,
545 DatanodeDropRegions,
547 DeleteTombstone,
549}