1use 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#[derive(Debug)]
54pub enum Control<T> {
55 Continue(T),
56 Stop,
57}
58
59impl<T> Control<T> {
60 pub fn stop(&self) -> bool {
62 matches!(self, Control::Stop)
63 }
64}
65
66impl DropTableExecutor {
67 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
77pub struct DropTableExecutor {
82 table: TableName,
83 table_id: TableId,
84 drop_if_exists: bool,
85}
86
87impl DropTableExecutor {
88 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 #[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 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 #[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 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 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 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 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 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 #[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 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 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 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 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 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 #[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 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 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 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}