1use std::sync::Arc;
16use std::time::Duration;
17
18use api::v1::Repartition;
19use api::v1::alter_table_expr::Kind;
20use api::v1::repartition::Source as PbRepartitionSource;
21use common_error::ext::BoxedError;
22use common_event_recorder::PersistentEventContext;
23#[cfg(feature = "enterprise")]
24use common_event_recorder::TriggerReason;
25use common_procedure::{
26 BoxedProcedure, BoxedProcedureLoader, Output, ProcedureContext, ProcedureId,
27 ProcedureManagerRef, ProcedureWithId, watcher,
28};
29use common_telemetry::tracing_context::{FutureExt, TracingContext};
30use common_telemetry::{debug, info, tracing};
31use derive_builder::Builder;
32use snafu::{OptionExt, ResultExt, ensure};
33use store_api::storage::{RegionId, TableId};
34use table::table_name::TableName;
35
36use crate::ddl::alter_database::AlterDatabaseProcedure;
37use crate::ddl::alter_logical_tables::AlterLogicalTablesProcedure;
38use crate::ddl::alter_table::{AlterTableProcedure, RegionRouteChanged, only_enables_skip_wal};
39use crate::ddl::comment_on::CommentOnProcedure;
40use crate::ddl::create_database::{CreateDatabaseMetadataCommitterRef, CreateDatabaseProcedure};
41use crate::ddl::create_flow::CreateFlowProcedure;
42use crate::ddl::create_logical_tables::CreateLogicalTablesProcedure;
43use crate::ddl::create_table::CreateTableProcedure;
44use crate::ddl::create_view::CreateViewProcedure;
45use crate::ddl::drop_database::DropDatabaseProcedure;
46use crate::ddl::drop_flow::DropFlowProcedure;
47use crate::ddl::drop_table::DropTableProcedure;
48use crate::ddl::drop_view::DropViewProcedure;
49#[cfg(feature = "enterprise")]
50use crate::ddl::purge_dropped_table::PurgeDroppedTableProcedure;
51use crate::ddl::truncate_table::TruncateTableProcedure;
52#[cfg(feature = "enterprise")]
53use crate::ddl::undrop_table::UndropTableProcedure;
54use crate::ddl::{DdlContext, utils};
55use crate::error::{
56 self, CreateRepartitionProcedureSnafu, EmptyDdlTasksSnafu,
57 PersistRepartitionGcRequirementSnafu, ProcedureOutputSnafu, RegisterProcedureLoaderSnafu,
58 RegisterRepartitionProcedureLoaderSnafu, Result, SubmitProcedureSnafu, TableInfoNotFoundSnafu,
59 TableNotFoundSnafu, TableRouteNotFoundSnafu, UnexpectedLogicalRouteTableSnafu,
60 UnsupportedSnafu, WaitProcedureSnafu,
61};
62use crate::key::table_info::TableInfoValue;
63use crate::key::table_name::TableNameKey;
64use crate::key::{DeserializedValueWithBytes, TableMetadataManagerRef};
65use crate::procedure_executor::ExecutorContext;
66#[cfg(feature = "enterprise")]
67use crate::rpc::ddl::DdlTask::CreateTrigger;
68#[cfg(feature = "enterprise")]
69use crate::rpc::ddl::DdlTask::DropTrigger;
70use crate::rpc::ddl::DdlTask::{
71 AlterDatabase, AlterLogicalTables, AlterTable, CommentOn, CreateDatabase, CreateFlow,
72 CreateLogicalTables, CreateTable, CreateView, DropDatabase, DropFlow, DropLogicalTables,
73 DropTable, DropView, PurgeDroppedTable, TruncateTable, UndropTable,
74};
75#[cfg(feature = "enterprise")]
76use crate::rpc::ddl::trigger::CreateTriggerTask;
77#[cfg(feature = "enterprise")]
78use crate::rpc::ddl::trigger::DropTriggerTask;
79use crate::rpc::ddl::{
80 AlterDatabaseTask, AlterTableTask, CommentOnTask, CreateDatabaseTask, CreateFlowTask,
81 CreateTableTask, CreateViewTask, DropDatabaseTask, DropFlowTask, DropTableTask, DropViewTask,
82 PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest, SubmitDdlTaskResponse,
83 TruncateTableTask, UndropTableTask,
84};
85
86const MAX_REGION_ROUTE_CHANGE_RETRIES: usize = 3;
87
88#[async_trait::async_trait]
90pub trait DdlManagerConfigurator<C>: Send + Sync {
91 async fn configure(
93 &self,
94 ddl_manager: DdlManager,
95 ctx: C,
96 ) -> std::result::Result<DdlManager, BoxedError>;
97}
98
99pub type DdlManagerConfiguratorRef<C> = Arc<dyn DdlManagerConfigurator<C>>;
100
101pub type DdlManagerRef = Arc<DdlManager>;
102
103pub type BoxedProcedureLoaderFactory = dyn Fn(DdlContext) -> BoxedProcedureLoader;
104
105#[derive(Builder)]
107pub struct DdlManager {
108 ddl_context: DdlContext,
109 procedure_manager: ProcedureManagerRef,
110 repartition_procedure_factory: RepartitionProcedureFactoryRef,
111 #[cfg(feature = "enterprise")]
112 trigger_ddl_manager: Option<TriggerDdlManagerRef>,
113}
114
115#[cfg(feature = "enterprise")]
118#[async_trait::async_trait]
119pub trait TriggerDdlManager: Send + Sync {
120 async fn create_trigger(
121 &self,
122 create_trigger_task: CreateTriggerTask,
123 procedure_manager: ProcedureManagerRef,
124 ddl_context: DdlContext,
125 query_context: QueryContext,
126 procedure_context: ProcedureContext,
127 ) -> Result<SubmitDdlTaskResponse>;
128
129 async fn drop_trigger(
130 &self,
131 drop_trigger_task: DropTriggerTask,
132 procedure_manager: ProcedureManagerRef,
133 ddl_context: DdlContext,
134 query_context: QueryContext,
135 procedure_context: ProcedureContext,
136 ) -> Result<SubmitDdlTaskResponse>;
137
138 fn as_any(&self) -> &dyn std::any::Any;
139}
140
141#[cfg(feature = "enterprise")]
142pub type TriggerDdlManagerRef = Arc<dyn TriggerDdlManager>;
143
144macro_rules! procedure_loader_entry {
145 ($procedure:ident) => {
146 (
147 $procedure::TYPE_NAME,
148 &|context: DdlContext| -> common_procedure::BoxedProcedureLoader {
149 Box::new(move |json: &str| {
150 let context = context.clone();
151 $procedure::from_json(json, context).map(|p| Box::new(p) as _)
152 })
153 },
154 )
155 };
156}
157
158macro_rules! procedure_loader {
159 ($($procedure:ident),*) => {
160 vec![
161 $(procedure_loader_entry!($procedure)),*
162 ]
163 };
164}
165
166pub type RepartitionProcedureFactoryRef = Arc<dyn RepartitionProcedureFactory>;
167
168pub enum RepartitionSource {
169 Partitioned {
170 exprs: Vec<String>,
171 target_partition_columns: Option<Vec<String>>,
177 },
178 Unpartitioned {
179 partition_columns: Vec<String>,
180 },
181}
182
183#[async_trait::async_trait]
184pub trait RepartitionProcedureFactory: Send + Sync {
185 #[allow(clippy::too_many_arguments)]
186 fn create(
187 &self,
188 ddl_ctx: &DdlContext,
189 table_name: TableName,
190 table_id: TableId,
191 source: RepartitionSource,
192 to_exprs: Vec<String>,
193 timeout: Option<Duration>,
194 ) -> std::result::Result<BoxedProcedure, BoxedError>;
195
196 fn register_loaders(
197 &self,
198 ddl_ctx: &DdlContext,
199 procedure_manager: &ProcedureManagerRef,
200 ) -> std::result::Result<(), BoxedError>;
201
202 async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError>;
204}
205
206#[derive(Debug, Clone, Copy)]
211pub struct DdlOptions {
212 pub timeout: Duration,
216 pub wait: bool,
223}
224
225impl DdlManager {
226 pub fn new(
228 ddl_context: DdlContext,
229 procedure_manager: ProcedureManagerRef,
230 repartition_procedure_factory: RepartitionProcedureFactoryRef,
231 ) -> Self {
232 Self {
233 ddl_context,
234 procedure_manager,
235 repartition_procedure_factory,
236 #[cfg(feature = "enterprise")]
237 trigger_ddl_manager: None,
238 }
239 }
240
241 #[cfg(feature = "enterprise")]
242 pub fn with_trigger_ddl_manager(mut self, trigger_ddl_manager: TriggerDdlManagerRef) -> Self {
243 self.trigger_ddl_manager = Some(trigger_ddl_manager);
244 self
245 }
246
247 pub fn with_create_database_metadata_committer(
248 mut self,
249 committer: CreateDatabaseMetadataCommitterRef,
250 ) -> Self {
251 self.ddl_context.create_database_metadata_committer = Some(committer);
252 self
253 }
254
255 pub fn table_metadata_manager(&self) -> &TableMetadataManagerRef {
257 &self.ddl_context.table_metadata_manager
258 }
259
260 pub fn create_context(&self) -> DdlContext {
262 self.ddl_context.clone()
263 }
264
265 pub fn register_loaders(&self) -> Result<()> {
267 let loaders: Vec<(&str, &BoxedProcedureLoaderFactory)> = procedure_loader!(
268 CreateTableProcedure,
269 CreateLogicalTablesProcedure,
270 CreateViewProcedure,
271 CreateFlowProcedure,
272 CreateDatabaseProcedure,
273 AlterTableProcedure,
274 AlterLogicalTablesProcedure,
275 AlterDatabaseProcedure,
276 DropTableProcedure,
277 DropFlowProcedure,
278 TruncateTableProcedure,
279 DropDatabaseProcedure,
280 DropViewProcedure,
281 CommentOnProcedure
282 );
283 #[cfg(feature = "enterprise")]
284 let loaders = {
285 let soft_drop_loaders: Vec<(&str, &BoxedProcedureLoaderFactory)> =
286 procedure_loader!(UndropTableProcedure, PurgeDroppedTableProcedure);
287 loaders
288 .into_iter()
289 .chain(soft_drop_loaders)
290 .collect::<Vec<_>>()
291 };
292
293 for (type_name, loader_factory) in loaders {
294 let context = self.create_context();
295 self.procedure_manager
296 .register_loader(type_name, loader_factory(context))
297 .context(RegisterProcedureLoaderSnafu { type_name })?;
298 }
299
300 #[cfg(feature = "enterprise")]
301 {
302 let type_name = PurgeDroppedTableProcedure::EXPIRED_TYPE_NAME;
303 let context = self.create_context();
304 self.procedure_manager
305 .register_loader(
306 type_name,
307 Box::new(move |json: &str| {
308 PurgeDroppedTableProcedure::from_json(json, context.clone())
309 .map(|procedure| Box::new(procedure) as _)
310 }),
311 )
312 .context(RegisterProcedureLoaderSnafu { type_name })?;
313 }
314
315 self.repartition_procedure_factory
316 .register_loaders(&self.ddl_context, &self.procedure_manager)
317 .context(RegisterRepartitionProcedureLoaderSnafu)?;
318
319 Ok(())
320 }
321
322 async fn submit_repartition_task(
341 &self,
342 table_id: TableId,
343 table_name: TableName,
344 repartition: Repartition,
345 wait: bool,
346 timeout: Duration,
347 procedure_context: ProcedureContext,
348 ) -> Result<(ProcedureId, Option<Output>)> {
349 let context = self.create_context();
350
351 let into_partition_exprs = repartition.into_partition_exprs;
352 let source = repartition.source;
353
354 let source = match source {
355 Some(PbRepartitionSource::PartitionExprs(source)) => RepartitionSource::Partitioned {
356 exprs: source.exprs,
357 target_partition_columns: source
358 .target_partition_columns
359 .map(|columns| columns.columns),
360 },
361 Some(PbRepartitionSource::Unpartitioned(source)) => RepartitionSource::Unpartitioned {
362 partition_columns: source.partition_columns,
363 },
364 None => {
365 #[allow(deprecated)]
367 RepartitionSource::Partitioned {
368 exprs: repartition.from_partition_exprs,
369 target_partition_columns: None,
370 }
371 }
372 };
373
374 let procedure = self
375 .repartition_procedure_factory
376 .create(
377 &context,
378 table_name,
379 table_id,
380 source,
381 into_partition_exprs,
382 Some(timeout),
383 )
384 .context(CreateRepartitionProcedureSnafu)?;
385 self.repartition_procedure_factory
386 .ensure_gc_requirement()
387 .await
388 .context(PersistRepartitionGcRequirementSnafu)?;
389 let procedure_with_id =
390 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
391 if wait {
392 self.execute_procedure_and_wait(procedure_with_id).await
393 } else {
394 self.submit_procedure(procedure_with_id)
395 .await
396 .map(|p| (p, None))
397 }
398 }
399
400 #[tracing::instrument(skip_all)]
402 pub async fn submit_alter_table_task(
403 &self,
404 table_id: TableId,
405 alter_table_task: AlterTableTask,
406 procedure_context: ProcedureContext,
407 ddl_options: DdlOptions,
408 ) -> Result<(ProcedureId, Option<Output>)> {
409 let mut alter_table_task = alter_table_task;
411 if let Some(Kind::Repartition(_)) = alter_table_task.alter_table.kind.as_ref()
412 && let Kind::Repartition(repartition) =
413 alter_table_task.alter_table.kind.take().unwrap()
414 {
415 let table_name = TableName::new(
416 alter_table_task.alter_table.catalog_name,
417 alter_table_task.alter_table.schema_name,
418 alter_table_task.alter_table.table_name,
419 );
420 return self
421 .submit_repartition_task(
422 table_id,
423 table_name,
424 repartition,
425 ddl_options.wait,
426 ddl_options.timeout,
427 procedure_context,
428 )
429 .await;
430 }
431
432 let lock_regions = alter_table_task
433 .alter_table
434 .kind
435 .as_ref()
436 .is_some_and(only_enables_skip_wal);
437
438 let mut route_change_retries = 0;
439 loop {
440 let region_ids_to_lock = if lock_regions {
443 let (_, route) = self
444 .table_metadata_manager()
445 .table_route_manager()
446 .get_physical_table_route(table_id)
447 .await?;
448 route
449 .region_routes
450 .iter()
451 .map(|route| route.region.id)
452 .collect::<Vec<RegionId>>()
453 } else {
454 vec![]
455 };
456
457 let context = self.create_context();
458 let procedure = AlterTableProcedure::new_with_region_locks(
459 table_id,
460 alter_table_task.clone(),
461 region_ids_to_lock,
462 context,
463 )?;
464
465 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
466 .with_context(procedure_context.clone());
467 let result = self.execute_procedure_and_wait(procedure_with_id).await?;
468 if result
469 .1
470 .as_ref()
471 .is_some_and(|output| output.is::<RegionRouteChanged>())
472 {
473 if route_change_retries == MAX_REGION_ROUTE_CHANGE_RETRIES {
474 let source = error::UnexpectedSnafu {
475 err_msg: format!(
476 "Region route kept changing while altering table {table_id}"
477 ),
478 }
479 .build();
480 return Err(error::Error::retry_later(source));
481 }
482 route_change_retries += 1;
483 continue;
484 }
485 return Ok(result);
486 }
487 }
488
489 #[tracing::instrument(skip_all)]
491 pub async fn submit_create_table_task(
492 &self,
493 create_table_task: CreateTableTask,
494 query_context: QueryContext,
495 procedure_context: ProcedureContext,
496 ) -> Result<(ProcedureId, Option<Output>)> {
497 let context = self.create_context();
498
499 let procedure = CreateTableProcedure::new_with_query_context(
500 create_table_task,
501 query_context,
502 context,
503 )?;
504
505 let procedure_with_id =
506 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
507
508 self.execute_procedure_and_wait(procedure_with_id).await
509 }
510
511 #[tracing::instrument(skip_all)]
513 pub async fn submit_create_view_task(
514 &self,
515 create_view_task: CreateViewTask,
516 procedure_context: ProcedureContext,
517 ) -> Result<(ProcedureId, Option<Output>)> {
518 let context = self.create_context();
519
520 let procedure = CreateViewProcedure::new(create_view_task, context);
521
522 let procedure_with_id =
523 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
524
525 self.execute_procedure_and_wait(procedure_with_id).await
526 }
527
528 #[tracing::instrument(skip_all)]
530 pub async fn submit_create_logical_table_tasks(
531 &self,
532 create_table_tasks: Vec<CreateTableTask>,
533 physical_table_id: TableId,
534 procedure_context: ProcedureContext,
535 ) -> Result<(ProcedureId, Option<Output>)> {
536 let context = self.create_context();
537
538 let procedure =
539 CreateLogicalTablesProcedure::new(create_table_tasks, physical_table_id, context);
540
541 let procedure_with_id =
542 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
543
544 self.execute_procedure_and_wait(procedure_with_id).await
545 }
546
547 #[tracing::instrument(skip_all)]
549 pub async fn submit_alter_logical_table_tasks(
550 &self,
551 alter_table_tasks: Vec<AlterTableTask>,
552 physical_table_id: TableId,
553 procedure_context: ProcedureContext,
554 ) -> Result<(ProcedureId, Option<Output>)> {
555 let context = self.create_context();
556
557 let procedure =
558 AlterLogicalTablesProcedure::new(alter_table_tasks, physical_table_id, context);
559
560 let procedure_with_id =
561 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
562
563 self.execute_procedure_and_wait(procedure_with_id).await
564 }
565
566 #[tracing::instrument(skip_all)]
568 pub async fn submit_drop_table_task(
569 &self,
570 drop_table_task: DropTableTask,
571 procedure_context: ProcedureContext,
572 ) -> Result<(ProcedureId, Option<Output>)> {
573 let context = self.create_context();
574
575 let procedure = DropTableProcedure::new(drop_table_task, context);
576
577 let procedure_with_id =
578 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
579
580 self.execute_procedure_and_wait(procedure_with_id).await
581 }
582
583 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
585 #[tracing::instrument(skip_all)]
586 pub async fn submit_undrop_table_task(
587 &self,
588 undrop_table_task: UndropTableTask,
589 procedure_context: ProcedureContext,
590 ) -> Result<(ProcedureId, Option<Output>)> {
591 #[cfg(not(feature = "enterprise"))]
592 {
593 use crate::error::UnsupportedSnafu;
594
595 return UnsupportedSnafu {
596 operation: "undrop table is only available in GreptimeDB Enterprise Edition",
597 }
598 .fail();
599 }
600 #[cfg(feature = "enterprise")]
601 {
602 let context = self.create_context();
603 let original_table_name = context
604 .table_metadata_manager
605 .get_dropped_table_by_id(undrop_table_task.table_id)
606 .await?
607 .with_context(|| TableNotFoundSnafu {
608 table_name: undrop_table_task.table_id.to_string(),
609 })?
610 .table_name;
611 let procedure = UndropTableProcedure::new_with_original_table_name(
612 undrop_table_task,
613 context,
614 Some(original_table_name),
615 );
616 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
617 .with_context(procedure_context);
618
619 self.execute_procedure_and_wait(procedure_with_id).await
620 }
621 }
622
623 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
625 #[tracing::instrument(skip_all)]
626 pub async fn submit_purge_dropped_table_task(
627 &self,
628 purge_dropped_table_task: PurgeDroppedTableTask,
629 procedure_context: ProcedureContext,
630 ) -> Result<(ProcedureId, Option<Output>)> {
631 #[cfg(not(feature = "enterprise"))]
632 {
633 use crate::error::UnsupportedSnafu;
634
635 return UnsupportedSnafu {
636 operation: "purge dropped table is only available in GreptimeDB Enterprise Edition",
637 }
638 .fail();
639 }
640 #[cfg(feature = "enterprise")]
641 {
642 let context = self.create_context();
643 let procedure = PurgeDroppedTableProcedure::new(purge_dropped_table_task, context);
644 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
645 .with_context(procedure_context);
646
647 self.execute_procedure_and_wait(procedure_with_id).await
648 }
649 }
650
651 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
653 #[tracing::instrument(skip_all)]
654 pub async fn submit_expired_purge_dropped_table_task(
655 &self,
656 purge_dropped_table_task: PurgeDroppedTableTask,
657 ) -> Result<(ProcedureId, Option<Output>)> {
658 #[cfg(not(feature = "enterprise"))]
659 {
660 use crate::error::UnsupportedSnafu;
661
662 return UnsupportedSnafu {
663 operation:
664 "purge expired dropped table is only available in GreptimeDB Enterprise Edition",
665 }
666 .fail();
667 }
668 #[cfg(feature = "enterprise")]
669 {
670 let context = self.create_context();
671 let procedure_context = ProcedureContext::from_event_context(
672 PersistentEventContext::new(TriggerReason::ScheduledGc),
673 );
674 let procedure =
675 PurgeDroppedTableProcedure::new_if_expired(purge_dropped_table_task, context);
676 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
677 .with_context(procedure_context);
678
679 self.execute_procedure_and_wait(procedure_with_id).await
680 }
681 }
682
683 #[tracing::instrument(skip_all)]
685 pub async fn submit_create_database(
686 &self,
687 CreateDatabaseTask {
688 catalog,
689 schema,
690 create_if_not_exists,
691 options,
692 creator,
693 }: CreateDatabaseTask,
694 procedure_context: ProcedureContext,
695 ) -> Result<(ProcedureId, Option<Output>)> {
696 let context = self.create_context();
697 let procedure = CreateDatabaseProcedure::new(
698 catalog,
699 schema,
700 create_if_not_exists,
701 options,
702 creator,
703 context,
704 );
705 let procedure_with_id =
706 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
707
708 self.execute_procedure_and_wait(procedure_with_id).await
709 }
710
711 #[tracing::instrument(skip_all)]
713 pub async fn submit_drop_database(
714 &self,
715 DropDatabaseTask {
716 catalog,
717 schema,
718 drop_if_exists,
719 }: DropDatabaseTask,
720 procedure_context: ProcedureContext,
721 ) -> Result<(ProcedureId, Option<Output>)> {
722 let context = self.create_context();
723 let procedure = DropDatabaseProcedure::new(catalog, schema, drop_if_exists, context);
724 let procedure_with_id =
725 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
726
727 self.execute_procedure_and_wait(procedure_with_id).await
728 }
729
730 pub async fn submit_alter_database(
731 &self,
732 alter_database_task: AlterDatabaseTask,
733 procedure_context: ProcedureContext,
734 ) -> Result<(ProcedureId, Option<Output>)> {
735 let context = self.create_context();
736 let procedure = AlterDatabaseProcedure::new(alter_database_task, context)?;
737 let procedure_with_id =
738 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
739
740 self.execute_procedure_and_wait(procedure_with_id).await
741 }
742
743 #[tracing::instrument(skip_all)]
745 pub async fn submit_create_flow_task(
746 &self,
747 create_flow: CreateFlowTask,
748 query_context: QueryContext,
749 procedure_context: ProcedureContext,
750 ) -> Result<(ProcedureId, Option<Output>)> {
751 let context = self.create_context();
752 let procedure = CreateFlowProcedure::new(create_flow, query_context, context);
753 let procedure_with_id =
754 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
755
756 self.execute_procedure_and_wait(procedure_with_id).await
757 }
758
759 #[tracing::instrument(skip_all)]
761 pub async fn submit_drop_flow_task(
762 &self,
763 drop_flow: DropFlowTask,
764 procedure_context: ProcedureContext,
765 ) -> Result<(ProcedureId, Option<Output>)> {
766 let context = self.create_context();
767 let procedure = DropFlowProcedure::new(drop_flow, context);
768 let procedure_with_id =
769 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
770
771 self.execute_procedure_and_wait(procedure_with_id).await
772 }
773
774 #[tracing::instrument(skip_all)]
776 pub async fn submit_drop_view_task(
777 &self,
778 drop_view: DropViewTask,
779 procedure_context: ProcedureContext,
780 ) -> Result<(ProcedureId, Option<Output>)> {
781 let context = self.create_context();
782 let procedure = DropViewProcedure::new(drop_view, context);
783 let procedure_with_id =
784 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
785
786 self.execute_procedure_and_wait(procedure_with_id).await
787 }
788
789 #[tracing::instrument(skip_all)]
791 pub async fn submit_truncate_table_task(
792 &self,
793 truncate_table_task: TruncateTableTask,
794 table_info_value: DeserializedValueWithBytes<TableInfoValue>,
795 procedure_context: ProcedureContext,
796 ) -> Result<(ProcedureId, Option<Output>)> {
797 let context = self.create_context();
798 let procedure = TruncateTableProcedure::new(truncate_table_task, table_info_value, context);
799
800 let procedure_with_id =
801 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
802
803 self.execute_procedure_and_wait(procedure_with_id).await
804 }
805
806 #[tracing::instrument(skip_all)]
808 pub async fn submit_comment_on_task(
809 &self,
810 mut comment_on_task: CommentOnTask,
811 procedure_context: ProcedureContext,
812 ) -> Result<(ProcedureId, Option<Output>)> {
813 let context = self.create_context();
814 comment_on_task
815 .enrich_object_id(
816 context.table_metadata_manager.table_name_manager(),
817 context.flow_metadata_manager.flow_name_manager(),
818 )
819 .await?;
820 let procedure = CommentOnProcedure::new(comment_on_task, context);
821 let procedure_with_id =
822 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
823
824 self.execute_procedure_and_wait(procedure_with_id).await
825 }
826
827 async fn execute_procedure_and_wait(
829 &self,
830 procedure_with_id: ProcedureWithId,
831 ) -> Result<(ProcedureId, Option<Output>)> {
832 let procedure_id = procedure_with_id.id;
833
834 let mut watcher = self
835 .procedure_manager
836 .submit(procedure_with_id)
837 .await
838 .context(SubmitProcedureSnafu)?;
839
840 let output = watcher::wait(&mut watcher)
841 .await
842 .context(WaitProcedureSnafu)?;
843
844 Ok((procedure_id, output))
845 }
846
847 async fn submit_procedure(&self, procedure_with_id: ProcedureWithId) -> Result<ProcedureId> {
849 let procedure_id = procedure_with_id.id;
850 let _ = self
851 .procedure_manager
852 .submit(procedure_with_id)
853 .await
854 .context(SubmitProcedureSnafu)?;
855
856 Ok(procedure_id)
857 }
858
859 pub async fn submit_ddl_task(
860 &self,
861 context: ExecutorContext,
862 request: SubmitDdlTaskRequest,
863 ) -> Result<SubmitDdlTaskResponse> {
864 let ExecutorContext {
865 tracing_context,
866 query_context,
867 actor,
868 event_input,
869 } = context;
870 let query_context = query_context.context(UnsupportedSnafu {
871 operation: "submit_ddl_task without query context",
872 })?;
873 let procedure_context = ProcedureContext {
874 actor,
875 event_context: event_input
876 .map(|input| PersistentEventContext::from((input, query_context.protocol()))),
877 };
878 let span = tracing_context
879 .as_ref()
880 .map(TracingContext::from_w3c)
881 .unwrap_or_else(TracingContext::from_current_span)
882 .attach(tracing::info_span!("DdlManager::submit_ddl_task"));
883 let SubmitDdlTaskRequest {
884 wait,
885 timeout,
886 task,
887 } = request;
888 let ddl_options = DdlOptions { wait, timeout };
889 async move {
890 debug!("Submitting Ddl task: {:?}", task);
891 match task {
892 CreateTable(create_table_task) => {
893 handle_create_table_task(
894 self,
895 create_table_task,
896 query_context,
897 procedure_context,
898 )
899 .await
900 }
901 DropTable(drop_table_task) => {
902 handle_drop_table_task(self, drop_table_task, procedure_context).await
903 }
904 UndropTable(undrop_table_task) => {
905 handle_undrop_table_task(self, undrop_table_task, procedure_context).await
906 }
907 PurgeDroppedTable(purge_dropped_table_task) => {
908 handle_purge_dropped_table_task(
909 self,
910 purge_dropped_table_task,
911 procedure_context,
912 )
913 .await
914 }
915 AlterTable(alter_table_task) => {
916 handle_alter_table_task(self, alter_table_task, ddl_options, procedure_context)
917 .await
918 }
919 TruncateTable(truncate_table_task) => {
920 handle_truncate_table_task(self, truncate_table_task, procedure_context).await
921 }
922 CreateLogicalTables(create_table_tasks) => {
923 handle_create_logical_table_tasks(self, create_table_tasks, procedure_context)
924 .await
925 }
926 AlterLogicalTables(alter_table_tasks) => {
927 handle_alter_logical_table_tasks(self, alter_table_tasks, procedure_context)
928 .await
929 }
930 DropLogicalTables(_) => todo!(),
931 CreateDatabase(create_database_task) => {
932 handle_create_database_task(self, create_database_task, procedure_context).await
933 }
934 DropDatabase(drop_database_task) => {
935 handle_drop_database_task(self, drop_database_task, procedure_context).await
936 }
937 AlterDatabase(alter_database_task) => {
938 handle_alter_database_task(self, alter_database_task, procedure_context).await
939 }
940 CreateFlow(create_flow_task) => {
941 handle_create_flow_task(
942 self,
943 create_flow_task,
944 query_context,
945 procedure_context,
946 )
947 .await
948 }
949 DropFlow(drop_flow_task) => {
950 handle_drop_flow_task(self, drop_flow_task, procedure_context).await
951 }
952 CreateView(create_view_task) => {
953 handle_create_view_task(self, create_view_task, procedure_context).await
954 }
955 DropView(drop_view_task) => {
956 handle_drop_view_task(self, drop_view_task, procedure_context).await
957 }
958 CommentOn(comment_on_task) => {
959 handle_comment_on_task(self, comment_on_task, procedure_context).await
960 }
961 #[cfg(feature = "enterprise")]
962 CreateTrigger(create_trigger_task) => {
963 handle_create_trigger_task(
964 self,
965 create_trigger_task,
966 query_context,
967 procedure_context,
968 )
969 .await
970 }
971 #[cfg(feature = "enterprise")]
972 DropTrigger(drop_trigger_task) => {
973 handle_drop_trigger_task(
974 self,
975 drop_trigger_task,
976 query_context,
977 procedure_context,
978 )
979 .await
980 }
981 }
982 }
983 .trace(span)
984 .await
985 }
986}
987
988async fn handle_truncate_table_task(
989 ddl_manager: &DdlManager,
990 truncate_table_task: TruncateTableTask,
991 procedure_context: ProcedureContext,
992) -> Result<SubmitDdlTaskResponse> {
993 let table_id = truncate_table_task.table_id;
994 let table_metadata_manager = &ddl_manager.table_metadata_manager();
995 let table_ref = truncate_table_task.table_ref();
996
997 let table_info_value = table_metadata_manager
998 .table_info_manager()
999 .get(table_id)
1000 .await?
1001 .with_context(|| TableInfoNotFoundSnafu {
1002 table: table_ref.to_string(),
1003 })?;
1004 let physical_table_id = table_metadata_manager
1005 .table_route_manager()
1006 .get_physical_table_id(table_id)
1007 .await?;
1008 ensure!(
1009 physical_table_id == table_id,
1010 error::UnexpectedSnafu {
1011 err_msg: "Truncate table is only supported for physical tables."
1012 }
1013 );
1014
1015 let (id, _) = ddl_manager
1016 .submit_truncate_table_task(truncate_table_task, table_info_value, procedure_context)
1017 .await?;
1018
1019 info!("Table: {table_id} is truncated via procedure_id {id:?}");
1020
1021 Ok(SubmitDdlTaskResponse {
1022 key: id.to_string().into(),
1023 ..Default::default()
1024 })
1025}
1026
1027async fn handle_alter_table_task(
1028 ddl_manager: &DdlManager,
1029 alter_table_task: AlterTableTask,
1030 ddl_options: DdlOptions,
1031 procedure_context: ProcedureContext,
1032) -> Result<SubmitDdlTaskResponse> {
1033 let table_ref = alter_table_task.table_ref();
1034
1035 let table_id = ddl_manager
1036 .table_metadata_manager()
1037 .table_name_manager()
1038 .get(TableNameKey::new(
1039 table_ref.catalog,
1040 table_ref.schema,
1041 table_ref.table,
1042 ))
1043 .await?
1044 .with_context(|| TableNotFoundSnafu {
1045 table_name: table_ref.to_string(),
1046 })?
1047 .table_id();
1048
1049 let table_route_value = ddl_manager
1050 .table_metadata_manager()
1051 .table_route_manager()
1052 .table_route_storage()
1053 .get(table_id)
1054 .await?
1055 .context(TableRouteNotFoundSnafu { table_id })?;
1056 ensure!(
1057 table_route_value.is_physical(),
1058 UnexpectedLogicalRouteTableSnafu {
1059 err_msg: format!("{:?} is a non-physical TableRouteValue.", table_ref),
1060 }
1061 );
1062
1063 let (id, _) = ddl_manager
1064 .submit_alter_table_task(table_id, alter_table_task, procedure_context, ddl_options)
1065 .await?;
1066
1067 info!("Table: {table_id} is altered via procedure_id {id:?}");
1068
1069 Ok(SubmitDdlTaskResponse {
1070 key: id.to_string().into(),
1071 ..Default::default()
1072 })
1073}
1074
1075async fn handle_drop_table_task(
1076 ddl_manager: &DdlManager,
1077 drop_table_task: DropTableTask,
1078 procedure_context: ProcedureContext,
1079) -> Result<SubmitDdlTaskResponse> {
1080 let table_id = drop_table_task.table_id;
1081 let (id, _) = ddl_manager
1082 .submit_drop_table_task(drop_table_task, procedure_context)
1083 .await?;
1084
1085 info!("Table: {table_id} is dropped via procedure_id {id:?}");
1086
1087 Ok(SubmitDdlTaskResponse {
1088 key: id.to_string().into(),
1089 ..Default::default()
1090 })
1091}
1092
1093async fn handle_undrop_table_task(
1094 ddl_manager: &DdlManager,
1095 undrop_table_task: UndropTableTask,
1096 procedure_context: ProcedureContext,
1097) -> Result<SubmitDdlTaskResponse> {
1098 let table_id = undrop_table_task.table_id;
1099 let (id, _) = ddl_manager
1100 .submit_undrop_table_task(undrop_table_task, procedure_context)
1101 .await?;
1102
1103 info!("Table: {table_id} is undropped via procedure_id {id:?}");
1104
1105 Ok(SubmitDdlTaskResponse {
1106 key: id.to_string().into(),
1107 ..Default::default()
1108 })
1109}
1110
1111async fn handle_purge_dropped_table_task(
1112 ddl_manager: &DdlManager,
1113 purge_dropped_table_task: PurgeDroppedTableTask,
1114 procedure_context: ProcedureContext,
1115) -> Result<SubmitDdlTaskResponse> {
1116 let (id, _) = ddl_manager
1117 .submit_purge_dropped_table_task(purge_dropped_table_task, procedure_context)
1118 .await?;
1119
1120 info!("Dropped table is purged via procedure_id {id:?}");
1121
1122 Ok(SubmitDdlTaskResponse {
1123 key: id.to_string().into(),
1124 ..Default::default()
1125 })
1126}
1127
1128async fn handle_create_table_task(
1129 ddl_manager: &DdlManager,
1130 create_table_task: CreateTableTask,
1131 query_context: QueryContext,
1132 procedure_context: ProcedureContext,
1133) -> Result<SubmitDdlTaskResponse> {
1134 let (id, output) = ddl_manager
1135 .submit_create_table_task(create_table_task, query_context, procedure_context)
1136 .await?;
1137
1138 let procedure_id = id.to_string();
1139 let output = output.context(ProcedureOutputSnafu {
1140 procedure_id: &procedure_id,
1141 err_msg: "empty output",
1142 })?;
1143 let table_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1144 procedure_id: &procedure_id,
1145 err_msg: "downcast to `u32`",
1146 })?);
1147 info!("Table: {table_id} is created via procedure_id {id:?}");
1148
1149 Ok(SubmitDdlTaskResponse {
1150 key: procedure_id.into(),
1151 table_ids: vec![table_id],
1152 })
1153}
1154
1155async fn handle_create_logical_table_tasks(
1156 ddl_manager: &DdlManager,
1157 create_table_tasks: Vec<CreateTableTask>,
1158 procedure_context: ProcedureContext,
1159) -> Result<SubmitDdlTaskResponse> {
1160 ensure!(
1161 !create_table_tasks.is_empty(),
1162 EmptyDdlTasksSnafu {
1163 name: "create logical tables"
1164 }
1165 );
1166 let physical_table_id = utils::check_and_get_physical_table_id(
1167 ddl_manager.table_metadata_manager(),
1168 &create_table_tasks,
1169 )
1170 .await?;
1171 let num_logical_tables = create_table_tasks.len();
1172
1173 let (id, output) = ddl_manager
1174 .submit_create_logical_table_tasks(create_table_tasks, physical_table_id, procedure_context)
1175 .await?;
1176
1177 info!(
1178 "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is created via procedure_id {id:?}"
1179 );
1180
1181 let procedure_id = id.to_string();
1182 let output = output.context(ProcedureOutputSnafu {
1183 procedure_id: &procedure_id,
1184 err_msg: "empty output",
1185 })?;
1186 let table_ids = output
1187 .downcast_ref::<Vec<TableId>>()
1188 .context(ProcedureOutputSnafu {
1189 procedure_id: &procedure_id,
1190 err_msg: "downcast to `Vec<TableId>`",
1191 })?
1192 .clone();
1193
1194 Ok(SubmitDdlTaskResponse {
1195 key: procedure_id.into(),
1196 table_ids,
1197 })
1198}
1199
1200async fn handle_create_database_task(
1201 ddl_manager: &DdlManager,
1202 create_database_task: CreateDatabaseTask,
1203 procedure_context: ProcedureContext,
1204) -> Result<SubmitDdlTaskResponse> {
1205 let catalog = create_database_task.catalog.clone();
1206 let schema = create_database_task.schema.clone();
1207 let (id, _) = ddl_manager
1208 .submit_create_database(create_database_task, procedure_context)
1209 .await?;
1210
1211 let procedure_id = id.to_string();
1212 info!(
1213 "Database {}.{} is created via procedure_id {id:?}",
1214 catalog, schema
1215 );
1216
1217 Ok(SubmitDdlTaskResponse {
1218 key: procedure_id.into(),
1219 ..Default::default()
1220 })
1221}
1222
1223async fn handle_drop_database_task(
1224 ddl_manager: &DdlManager,
1225 drop_database_task: DropDatabaseTask,
1226 procedure_context: ProcedureContext,
1227) -> Result<SubmitDdlTaskResponse> {
1228 let (id, _) = ddl_manager
1229 .submit_drop_database(drop_database_task.clone(), procedure_context)
1230 .await?;
1231
1232 let procedure_id = id.to_string();
1233 info!(
1234 "Database {}.{} is dropped via procedure_id {id:?}",
1235 drop_database_task.catalog, drop_database_task.schema
1236 );
1237
1238 Ok(SubmitDdlTaskResponse {
1239 key: procedure_id.into(),
1240 ..Default::default()
1241 })
1242}
1243
1244async fn handle_alter_database_task(
1245 ddl_manager: &DdlManager,
1246 alter_database_task: AlterDatabaseTask,
1247 procedure_context: ProcedureContext,
1248) -> Result<SubmitDdlTaskResponse> {
1249 let (id, _) = ddl_manager
1250 .submit_alter_database(alter_database_task.clone(), procedure_context)
1251 .await?;
1252
1253 let procedure_id = id.to_string();
1254 info!(
1255 "Database {}.{} is altered via procedure_id {id:?}",
1256 alter_database_task.catalog(),
1257 alter_database_task.schema()
1258 );
1259
1260 Ok(SubmitDdlTaskResponse {
1261 key: procedure_id.into(),
1262 ..Default::default()
1263 })
1264}
1265
1266async fn handle_drop_flow_task(
1267 ddl_manager: &DdlManager,
1268 drop_flow_task: DropFlowTask,
1269 procedure_context: ProcedureContext,
1270) -> Result<SubmitDdlTaskResponse> {
1271 let (id, _) = ddl_manager
1272 .submit_drop_flow_task(drop_flow_task.clone(), procedure_context)
1273 .await?;
1274
1275 let procedure_id = id.to_string();
1276 info!(
1277 "Flow {}.{}({}) is dropped via procedure_id {id:?}",
1278 drop_flow_task.catalog_name, drop_flow_task.flow_name, drop_flow_task.flow_id,
1279 );
1280
1281 Ok(SubmitDdlTaskResponse {
1282 key: procedure_id.into(),
1283 ..Default::default()
1284 })
1285}
1286
1287#[cfg(feature = "enterprise")]
1288async fn handle_drop_trigger_task(
1289 ddl_manager: &DdlManager,
1290 drop_trigger_task: DropTriggerTask,
1291 query_context: QueryContext,
1292 procedure_context: ProcedureContext,
1293) -> Result<SubmitDdlTaskResponse> {
1294 let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1295 use crate::error::UnsupportedSnafu;
1296
1297 return UnsupportedSnafu {
1298 operation: "drop trigger",
1299 }
1300 .fail();
1301 };
1302
1303 m.drop_trigger(
1304 drop_trigger_task,
1305 ddl_manager.procedure_manager.clone(),
1306 ddl_manager.ddl_context.clone(),
1307 query_context,
1308 procedure_context,
1309 )
1310 .await
1311}
1312
1313async fn handle_drop_view_task(
1314 ddl_manager: &DdlManager,
1315 drop_view_task: DropViewTask,
1316 procedure_context: ProcedureContext,
1317) -> Result<SubmitDdlTaskResponse> {
1318 let (id, _) = ddl_manager
1319 .submit_drop_view_task(drop_view_task.clone(), procedure_context)
1320 .await?;
1321
1322 let procedure_id = id.to_string();
1323 info!(
1324 "View {}({}) is dropped via procedure_id {id:?}",
1325 drop_view_task.table_ref(),
1326 drop_view_task.view_id,
1327 );
1328
1329 Ok(SubmitDdlTaskResponse {
1330 key: procedure_id.into(),
1331 ..Default::default()
1332 })
1333}
1334
1335async fn handle_create_flow_task(
1336 ddl_manager: &DdlManager,
1337 create_flow_task: CreateFlowTask,
1338 query_context: QueryContext,
1339 procedure_context: ProcedureContext,
1340) -> Result<SubmitDdlTaskResponse> {
1341 let (id, output) = ddl_manager
1342 .submit_create_flow_task(create_flow_task.clone(), query_context, procedure_context)
1343 .await?;
1344
1345 let procedure_id = id.to_string();
1346 let output = output.context(ProcedureOutputSnafu {
1347 procedure_id: &procedure_id,
1348 err_msg: "empty output",
1349 })?;
1350 let flow_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1351 procedure_id: &procedure_id,
1352 err_msg: "downcast to `u32`",
1353 })?);
1354 if !create_flow_task.or_replace {
1355 info!(
1356 "Flow {}.{}({flow_id}) is created via procedure_id {id:?}",
1357 create_flow_task.catalog_name, create_flow_task.flow_name,
1358 );
1359 } else {
1360 info!(
1361 "Flow {}.{}({flow_id}) is replaced via procedure_id {id:?}",
1362 create_flow_task.catalog_name, create_flow_task.flow_name,
1363 );
1364 }
1365
1366 Ok(SubmitDdlTaskResponse {
1367 key: procedure_id.into(),
1368 ..Default::default()
1369 })
1370}
1371
1372#[cfg(feature = "enterprise")]
1373async fn handle_create_trigger_task(
1374 ddl_manager: &DdlManager,
1375 create_trigger_task: CreateTriggerTask,
1376 query_context: QueryContext,
1377 procedure_context: ProcedureContext,
1378) -> Result<SubmitDdlTaskResponse> {
1379 let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1380 use crate::error::UnsupportedSnafu;
1381
1382 return UnsupportedSnafu {
1383 operation: "create trigger",
1384 }
1385 .fail();
1386 };
1387
1388 m.create_trigger(
1389 create_trigger_task,
1390 ddl_manager.procedure_manager.clone(),
1391 ddl_manager.ddl_context.clone(),
1392 query_context,
1393 procedure_context,
1394 )
1395 .await
1396}
1397
1398async fn handle_alter_logical_table_tasks(
1399 ddl_manager: &DdlManager,
1400 alter_table_tasks: Vec<AlterTableTask>,
1401 procedure_context: ProcedureContext,
1402) -> Result<SubmitDdlTaskResponse> {
1403 ensure!(
1404 !alter_table_tasks.is_empty(),
1405 EmptyDdlTasksSnafu {
1406 name: "alter logical tables"
1407 }
1408 );
1409
1410 let first_table = TableNameKey {
1412 catalog: &alter_table_tasks[0].alter_table.catalog_name,
1413 schema: &alter_table_tasks[0].alter_table.schema_name,
1414 table: &alter_table_tasks[0].alter_table.table_name,
1415 };
1416 let physical_table_id =
1417 utils::get_physical_table_id(ddl_manager.table_metadata_manager(), first_table).await?;
1418 let num_logical_tables = alter_table_tasks.len();
1419
1420 let (id, _) = ddl_manager
1421 .submit_alter_logical_table_tasks(alter_table_tasks, physical_table_id, procedure_context)
1422 .await?;
1423
1424 info!(
1425 "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is altered via procedure_id {id:?}"
1426 );
1427
1428 let procedure_id = id.to_string();
1429
1430 Ok(SubmitDdlTaskResponse {
1431 key: procedure_id.into(),
1432 ..Default::default()
1433 })
1434}
1435
1436async fn handle_create_view_task(
1438 ddl_manager: &DdlManager,
1439 create_view_task: CreateViewTask,
1440 procedure_context: ProcedureContext,
1441) -> Result<SubmitDdlTaskResponse> {
1442 let (id, output) = ddl_manager
1443 .submit_create_view_task(create_view_task, procedure_context)
1444 .await?;
1445
1446 let procedure_id = id.to_string();
1447 let output = output.context(ProcedureOutputSnafu {
1448 procedure_id: &procedure_id,
1449 err_msg: "empty output",
1450 })?;
1451 let view_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1452 procedure_id: &procedure_id,
1453 err_msg: "downcast to `u32`",
1454 })?);
1455 info!("View: {view_id} is created via procedure_id {id:?}");
1456
1457 Ok(SubmitDdlTaskResponse {
1458 key: procedure_id.into(),
1459 table_ids: vec![view_id],
1460 })
1461}
1462
1463async fn handle_comment_on_task(
1464 ddl_manager: &DdlManager,
1465 comment_on_task: CommentOnTask,
1466 procedure_context: ProcedureContext,
1467) -> Result<SubmitDdlTaskResponse> {
1468 let (id, _) = ddl_manager
1469 .submit_comment_on_task(comment_on_task.clone(), procedure_context)
1470 .await?;
1471
1472 let procedure_id = id.to_string();
1473 info!(
1474 "Comment on {}.{}.{} is updated via procedure_id {id:?}",
1475 comment_on_task.catalog_name, comment_on_task.schema_name, comment_on_task.object_name
1476 );
1477
1478 Ok(SubmitDdlTaskResponse {
1479 key: procedure_id.into(),
1480 ..Default::default()
1481 })
1482}
1483
1484#[cfg(test)]
1485mod tests {
1486 use std::sync::Arc;
1487 #[cfg(feature = "enterprise")]
1488 use std::sync::Mutex;
1489 use std::time::Duration;
1490
1491 #[cfg(feature = "enterprise")]
1492 use common_base::protocol::Channel;
1493 use common_error::ext::BoxedError;
1494 #[cfg(feature = "enterprise")]
1495 use common_error::ext::ErrorExt;
1496 #[cfg(feature = "enterprise")]
1497 use common_error::status_code::StatusCode;
1498 #[cfg(feature = "enterprise")]
1499 use common_event_recorder::{PersistentEventContext, ProcedureEventInput, TriggerReason};
1500 use common_procedure::local::LocalManager;
1501 use common_procedure::test_util::InMemoryPoisonStore;
1502 use common_procedure::{BoxedProcedure, ProcedureContext, ProcedureManagerRef};
1503 use store_api::storage::TableId;
1504 use table::table_name::TableName;
1505
1506 use super::DdlManager;
1507 use crate::cache_invalidator::DummyCacheInvalidator;
1508 use crate::ddl::alter_table::AlterTableProcedure;
1509 use crate::ddl::create_database::{
1510 AtomicCreateOutcome, CreateDatabaseMetadataCommitter, CreateDatabaseProcedure,
1511 };
1512 use crate::ddl::create_table::CreateTableProcedure;
1513 use crate::ddl::drop_table::DropTableProcedure;
1514 use crate::ddl::flow_meta::FlowMetadataAllocator;
1515 use crate::ddl::table_meta::TableMetadataAllocator;
1516 use crate::ddl::truncate_table::TruncateTableProcedure;
1517 use crate::ddl::{DdlContext, NoopRegionFailureDetectorControl};
1518 use crate::ddl_manager::{RepartitionProcedureFactory, RepartitionSource};
1519 use crate::key::TableMetadataManager;
1520 use crate::key::flow::FlowMetadataManager;
1521 use crate::kv_backend::memory::MemoryKvBackend;
1522 use crate::node_manager::{DatanodeManager, DatanodeRef, FlownodeManager, FlownodeRef};
1523 use crate::peer::Peer;
1524 use crate::procedure_executor::ExecutorContext;
1525 use crate::region_keeper::MemoryRegionKeeper;
1526 use crate::region_registry::LeaderRegionRegistry;
1527 #[cfg(feature = "enterprise")]
1528 use crate::rpc::ddl::trigger::{CreateTriggerTask, DropTriggerTask};
1529 use crate::rpc::ddl::{CreatorGrantIntent, UndropTableTask};
1530 #[cfg(not(feature = "enterprise"))]
1531 use crate::rpc::ddl::{DdlTask, PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest};
1532 #[cfg(feature = "enterprise")]
1533 use crate::rpc::ddl::{DdlTask, QueryContext, SubmitDdlTaskRequest};
1534 use crate::sequence::SequenceBuilder;
1535 use crate::state_store::KvStateStore;
1536 use crate::test_util::{MockDatanodeManager, new_ddl_context};
1537 use crate::wal_provider::WalProvider;
1538
1539 pub struct DummyDatanodeManager;
1541
1542 #[async_trait::async_trait]
1543 impl DatanodeManager for DummyDatanodeManager {
1544 async fn datanode(&self, _datanode: &Peer) -> DatanodeRef {
1545 unimplemented!()
1546 }
1547 }
1548
1549 #[async_trait::async_trait]
1550 impl FlownodeManager for DummyDatanodeManager {
1551 async fn flownode(&self, _node: &Peer) -> FlownodeRef {
1552 unimplemented!()
1553 }
1554 }
1555
1556 struct DummyRepartitionProcedureFactory;
1557
1558 #[async_trait::async_trait]
1559 impl RepartitionProcedureFactory for DummyRepartitionProcedureFactory {
1560 fn create(
1561 &self,
1562 _ddl_ctx: &DdlContext,
1563 _table_name: TableName,
1564 _table_id: TableId,
1565 _source: RepartitionSource,
1566 _to_exprs: Vec<String>,
1567 _timeout: Option<Duration>,
1568 ) -> std::result::Result<BoxedProcedure, BoxedError> {
1569 unimplemented!()
1570 }
1571
1572 fn register_loaders(
1573 &self,
1574 _ddl_ctx: &DdlContext,
1575 _procedure_manager: &ProcedureManagerRef,
1576 ) -> std::result::Result<(), BoxedError> {
1577 Ok(())
1578 }
1579
1580 async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError> {
1581 Ok(())
1582 }
1583 }
1584
1585 struct LoaderCommitter;
1586
1587 #[async_trait::async_trait]
1588 impl CreateDatabaseMetadataCommitter for LoaderCommitter {
1589 async fn commit(
1590 &self,
1591 _catalog: &str,
1592 _schema: &str,
1593 _value: &crate::key::schema_name::SchemaNameValue,
1594 _creator: &CreatorGrantIntent,
1595 ) -> std::result::Result<AtomicCreateOutcome, BoxedError> {
1596 unreachable!()
1597 }
1598 }
1599
1600 #[cfg(feature = "enterprise")]
1601 #[derive(Default)]
1602 struct RecordingTriggerDdlManager {
1603 procedure_contexts: Mutex<Vec<ProcedureContext>>,
1604 }
1605
1606 #[cfg(feature = "enterprise")]
1607 #[async_trait::async_trait]
1608 impl super::TriggerDdlManager for RecordingTriggerDdlManager {
1609 async fn create_trigger(
1610 &self,
1611 _create_trigger_task: CreateTriggerTask,
1612 _procedure_manager: ProcedureManagerRef,
1613 _ddl_context: DdlContext,
1614 _query_context: QueryContext,
1615 procedure_context: ProcedureContext,
1616 ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1617 self.procedure_contexts
1618 .lock()
1619 .unwrap()
1620 .push(procedure_context);
1621 Ok(Default::default())
1622 }
1623
1624 async fn drop_trigger(
1625 &self,
1626 _drop_trigger_task: DropTriggerTask,
1627 _procedure_manager: ProcedureManagerRef,
1628 _ddl_context: DdlContext,
1629 _query_context: QueryContext,
1630 procedure_context: ProcedureContext,
1631 ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1632 self.procedure_contexts
1633 .lock()
1634 .unwrap()
1635 .push(procedure_context);
1636 Ok(Default::default())
1637 }
1638
1639 fn as_any(&self) -> &dyn std::any::Any {
1640 self
1641 }
1642 }
1643
1644 #[test]
1645 fn test_generic_loader_captures_configured_committer() {
1646 let mut context = new_ddl_context(Arc::new(MockDatanodeManager::new(())));
1647 let committer = Arc::new(LoaderCommitter);
1648 context.create_database_metadata_committer = Some(committer.clone());
1649 let loader = procedure_loader!(CreateDatabaseProcedure)[0].1(context);
1650 assert_eq!(Arc::strong_count(&committer), 2);
1651
1652 let loaded = loader(
1653 r#"{
1654 "state":"CreateMetadata",
1655 "catalog":"greptime",
1656 "schema":"metrics",
1657 "create_if_not_exists":false,
1658 "options":{},
1659 "creator":{"username":"alice","created_at_ns":42}
1660 }"#,
1661 )
1662 .unwrap();
1663
1664 assert_eq!(loaded.type_name(), "metasrv-procedure::CreateDatabase");
1665 assert_eq!(Arc::strong_count(&committer), 3);
1666 drop(loaded);
1667 assert_eq!(Arc::strong_count(&committer), 2);
1668 }
1669
1670 #[test]
1671 fn test_register_loaders() {
1672 let kv_backend = Arc::new(MemoryKvBackend::new());
1673 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1674 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1675 Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1676 Arc::new(WalProvider::default()),
1677 ));
1678 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1679 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1680 Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1681 ));
1682
1683 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1684 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1685 let procedure_manager = Arc::new(LocalManager::new(
1686 Default::default(),
1687 state_store,
1688 poison_manager,
1689 None,
1690 None,
1691 ));
1692
1693 let ddl_manager = DdlManager::new(
1694 DdlContext {
1695 node_manager: Arc::new(DummyDatanodeManager),
1696 cache_invalidator: Arc::new(DummyCacheInvalidator),
1697 table_metadata_manager,
1698 table_metadata_allocator,
1699 flow_metadata_manager,
1700 flow_metadata_allocator,
1701 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1702 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1703 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1704 soft_drop_enabled: false,
1705 soft_drop_retention: None,
1706 create_database_metadata_committer: None,
1707 },
1708 procedure_manager.clone(),
1709 Arc::new(DummyRepartitionProcedureFactory),
1710 );
1711 ddl_manager.register_loaders().unwrap();
1712
1713 let expected_loaders = vec![
1714 CreateTableProcedure::TYPE_NAME,
1715 AlterTableProcedure::TYPE_NAME,
1716 DropTableProcedure::TYPE_NAME,
1717 TruncateTableProcedure::TYPE_NAME,
1718 ];
1719
1720 for loader in expected_loaders {
1721 assert!(procedure_manager.contains_loader(loader));
1722 }
1723
1724 let soft_drop_loaders = [
1725 "metasrv-procedure::UndropTable",
1726 "metasrv-procedure::PurgeDroppedTable",
1727 "metasrv-procedure::PurgeExpiredDroppedTable",
1728 ];
1729 for loader in soft_drop_loaders {
1730 assert_eq!(
1731 cfg!(feature = "enterprise"),
1732 procedure_manager.contains_loader(loader)
1733 );
1734 }
1735 }
1736
1737 fn build_soft_drop_test_ddl_manager() -> DdlManager {
1738 let kv_backend = Arc::new(MemoryKvBackend::new());
1739 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1740 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1741 Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1742 Arc::new(WalProvider::default()),
1743 ));
1744 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1745 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1746 Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1747 ));
1748
1749 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1750 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1751 let procedure_manager = Arc::new(LocalManager::new(
1752 Default::default(),
1753 state_store,
1754 poison_manager,
1755 None,
1756 None,
1757 ));
1758
1759 DdlManager::new(
1760 DdlContext {
1761 node_manager: Arc::new(DummyDatanodeManager),
1762 cache_invalidator: Arc::new(DummyCacheInvalidator),
1763 table_metadata_manager,
1764 table_metadata_allocator,
1765 flow_metadata_manager,
1766 flow_metadata_allocator,
1767 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1768 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1769 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1770 soft_drop_enabled: true,
1771 soft_drop_retention: Some(std::time::Duration::from_secs(1)),
1772 create_database_metadata_committer: None,
1773 },
1774 procedure_manager,
1775 Arc::new(DummyRepartitionProcedureFactory),
1776 )
1777 }
1778
1779 #[cfg(feature = "enterprise")]
1780 #[tokio::test]
1781 async fn test_trigger_ddl_forwards_procedure_context() {
1782 let trigger_ddl_manager = Arc::new(RecordingTriggerDdlManager::default());
1783 let ddl_manager = build_soft_drop_test_ddl_manager()
1784 .with_trigger_ddl_manager(trigger_ddl_manager.clone());
1785 let procedure_context = ProcedureContext {
1786 actor: Some("test-user".to_string()),
1787 event_context: Some(
1788 PersistentEventContext::new(TriggerReason::Manual).with_protocol("mysql"),
1789 ),
1790 };
1791 let executor_context = || ExecutorContext {
1792 query_context: Some(QueryContext {
1793 channel: Channel::Mysql as u8,
1794 ..Default::default()
1795 }),
1796 actor: Some("test-user".to_string()),
1797 event_input: Some(ProcedureEventInput::new(TriggerReason::Manual)),
1798 ..Default::default()
1799 };
1800
1801 ddl_manager
1802 .submit_ddl_task(
1803 executor_context(),
1804 SubmitDdlTaskRequest::new(DdlTask::CreateTrigger(CreateTriggerTask {
1805 catalog_name: "greptime".to_string(),
1806 trigger_name: "test_trigger".to_string(),
1807 if_not_exists: false,
1808 sql: "SELECT 1".to_string(),
1809 channels: vec![],
1810 labels: Default::default(),
1811 annotations: Default::default(),
1812 interval: Duration::from_secs(1),
1813 raw_interval_expr: None,
1814 r#for: None,
1815 for_raw_expr: None,
1816 keep_firing_for: None,
1817 keep_firing_for_raw_expr: None,
1818 })),
1819 )
1820 .await
1821 .unwrap();
1822 ddl_manager
1823 .submit_ddl_task(
1824 executor_context(),
1825 SubmitDdlTaskRequest::new(DdlTask::DropTrigger(DropTriggerTask {
1826 catalog_name: "greptime".to_string(),
1827 trigger_name: "test_trigger".to_string(),
1828 drop_if_exists: false,
1829 })),
1830 )
1831 .await
1832 .unwrap();
1833
1834 assert_eq!(
1835 *trigger_ddl_manager.procedure_contexts.lock().unwrap(),
1836 vec![procedure_context.clone(), procedure_context]
1837 );
1838 }
1839
1840 #[cfg(feature = "enterprise")]
1841 #[tokio::test]
1842 async fn test_submit_undrop_missing_tombstone_returns_table_not_found_directly() {
1843 let ddl_manager = build_soft_drop_test_ddl_manager();
1844
1845 let err = ddl_manager
1846 .submit_undrop_table_task(
1847 UndropTableTask { table_id: 1024 },
1848 ProcedureContext::default(),
1849 )
1850 .await
1851 .unwrap_err();
1852
1853 assert_eq!(err.status_code(), StatusCode::TableNotFound);
1854 assert!(matches!(err, crate::error::Error::TableNotFound { .. }));
1855 }
1856
1857 #[cfg(not(feature = "enterprise"))]
1858 #[tokio::test]
1859 async fn test_submit_undrop_and_purge_rejected_in_non_enterprise_build() {
1860 let ddl_manager = build_soft_drop_test_ddl_manager();
1861
1862 let err = ddl_manager
1863 .submit_undrop_table_task(
1864 UndropTableTask { table_id: 1024 },
1865 ProcedureContext::default(),
1866 )
1867 .await
1868 .unwrap_err();
1869 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1870
1871 let err = ddl_manager
1872 .submit_purge_dropped_table_task(
1873 PurgeDroppedTableTask { table_id: 1024 },
1874 ProcedureContext::default(),
1875 )
1876 .await
1877 .unwrap_err();
1878 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1879
1880 let err = ddl_manager
1881 .submit_expired_purge_dropped_table_task(PurgeDroppedTableTask { table_id: 1024 })
1882 .await
1883 .unwrap_err();
1884 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1885
1886 for task in [
1887 DdlTask::UndropTable(UndropTableTask { table_id: 1024 }),
1888 DdlTask::PurgeDroppedTable(PurgeDroppedTableTask { table_id: 1024 }),
1889 ] {
1890 let err = ddl_manager
1891 .submit_ddl_task(
1892 ExecutorContext {
1893 query_context: Some(QueryContext::default()),
1894 ..Default::default()
1895 },
1896 SubmitDdlTaskRequest::new(task),
1897 )
1898 .await
1899 .unwrap_err();
1900 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1901 }
1902 }
1903}