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, ConvertAlterTableRequestSnafu, 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 logical_table_ids = {
560 let table_refs = alter_table_tasks
561 .iter()
562 .map(|task| task.table_ref())
563 .collect::<Vec<_>>();
564 utils::table_id::get_all_table_ids_by_names(
565 self.table_metadata_manager().table_name_manager(),
566 &table_refs,
567 )
568 .await?
569 };
570
571 let procedure = AlterLogicalTablesProcedure::new(
572 alter_table_tasks,
573 physical_table_id,
574 logical_table_ids,
575 context,
576 );
577
578 let procedure_with_id =
579 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
580
581 self.execute_procedure_and_wait(procedure_with_id).await
582 }
583
584 #[tracing::instrument(skip_all)]
586 pub async fn submit_drop_table_task(
587 &self,
588 drop_table_task: DropTableTask,
589 procedure_context: ProcedureContext,
590 ) -> Result<(ProcedureId, Option<Output>)> {
591 let context = self.create_context();
592
593 let procedure = DropTableProcedure::new(drop_table_task, context);
594
595 let procedure_with_id =
596 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
597
598 self.execute_procedure_and_wait(procedure_with_id).await
599 }
600
601 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
603 #[tracing::instrument(skip_all)]
604 pub async fn submit_undrop_table_task(
605 &self,
606 undrop_table_task: UndropTableTask,
607 procedure_context: ProcedureContext,
608 ) -> Result<(ProcedureId, Option<Output>)> {
609 #[cfg(not(feature = "enterprise"))]
610 {
611 use crate::error::UnsupportedSnafu;
612
613 return UnsupportedSnafu {
614 operation: "undrop table is only available in GreptimeDB Enterprise Edition",
615 }
616 .fail();
617 }
618 #[cfg(feature = "enterprise")]
619 {
620 let context = self.create_context();
621 let original_table_name = context
622 .table_metadata_manager
623 .get_dropped_table_by_id(undrop_table_task.table_id)
624 .await?
625 .with_context(|| TableNotFoundSnafu {
626 table_name: undrop_table_task.table_id.to_string(),
627 })?
628 .table_name;
629 let procedure = UndropTableProcedure::new_with_original_table_name(
630 undrop_table_task,
631 context,
632 Some(original_table_name),
633 );
634 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
635 .with_context(procedure_context);
636
637 self.execute_procedure_and_wait(procedure_with_id).await
638 }
639 }
640
641 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
643 #[tracing::instrument(skip_all)]
644 pub async fn submit_purge_dropped_table_task(
645 &self,
646 purge_dropped_table_task: PurgeDroppedTableTask,
647 procedure_context: ProcedureContext,
648 ) -> Result<(ProcedureId, Option<Output>)> {
649 #[cfg(not(feature = "enterprise"))]
650 {
651 use crate::error::UnsupportedSnafu;
652
653 return UnsupportedSnafu {
654 operation: "purge dropped table is only available in GreptimeDB Enterprise Edition",
655 }
656 .fail();
657 }
658 #[cfg(feature = "enterprise")]
659 {
660 let context = self.create_context();
661 let procedure = PurgeDroppedTableProcedure::new(purge_dropped_table_task, context);
662 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
663 .with_context(procedure_context);
664
665 self.execute_procedure_and_wait(procedure_with_id).await
666 }
667 }
668
669 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
671 #[tracing::instrument(skip_all)]
672 pub async fn submit_expired_purge_dropped_table_task(
673 &self,
674 purge_dropped_table_task: PurgeDroppedTableTask,
675 ) -> Result<(ProcedureId, Option<Output>)> {
676 #[cfg(not(feature = "enterprise"))]
677 {
678 use crate::error::UnsupportedSnafu;
679
680 return UnsupportedSnafu {
681 operation:
682 "purge expired dropped table is only available in GreptimeDB Enterprise Edition",
683 }
684 .fail();
685 }
686 #[cfg(feature = "enterprise")]
687 {
688 let context = self.create_context();
689 let procedure_context = ProcedureContext::from_event_context(
690 PersistentEventContext::new(TriggerReason::ScheduledGc),
691 );
692 let procedure =
693 PurgeDroppedTableProcedure::new_if_expired(purge_dropped_table_task, context);
694 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
695 .with_context(procedure_context);
696
697 self.execute_procedure_and_wait(procedure_with_id).await
698 }
699 }
700
701 #[tracing::instrument(skip_all)]
703 pub async fn submit_create_database(
704 &self,
705 CreateDatabaseTask {
706 catalog,
707 schema,
708 create_if_not_exists,
709 options,
710 creator,
711 }: CreateDatabaseTask,
712 procedure_context: ProcedureContext,
713 ) -> Result<(ProcedureId, Option<Output>)> {
714 let context = self.create_context();
715 let procedure = CreateDatabaseProcedure::new(
716 catalog,
717 schema,
718 create_if_not_exists,
719 options,
720 creator,
721 context,
722 );
723 let procedure_with_id =
724 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
725
726 self.execute_procedure_and_wait(procedure_with_id).await
727 }
728
729 #[tracing::instrument(skip_all)]
731 pub async fn submit_drop_database(
732 &self,
733 DropDatabaseTask {
734 catalog,
735 schema,
736 drop_if_exists,
737 }: DropDatabaseTask,
738 procedure_context: ProcedureContext,
739 ) -> Result<(ProcedureId, Option<Output>)> {
740 let context = self.create_context();
741 let procedure = DropDatabaseProcedure::new(catalog, schema, drop_if_exists, context);
742 let procedure_with_id =
743 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
744
745 self.execute_procedure_and_wait(procedure_with_id).await
746 }
747
748 pub async fn submit_alter_database(
749 &self,
750 alter_database_task: AlterDatabaseTask,
751 procedure_context: ProcedureContext,
752 ) -> Result<(ProcedureId, Option<Output>)> {
753 let context = self.create_context();
754 let procedure = AlterDatabaseProcedure::new(alter_database_task, context)?;
755 let procedure_with_id =
756 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
757
758 self.execute_procedure_and_wait(procedure_with_id).await
759 }
760
761 #[tracing::instrument(skip_all)]
763 pub async fn submit_create_flow_task(
764 &self,
765 create_flow: CreateFlowTask,
766 query_context: QueryContext,
767 procedure_context: ProcedureContext,
768 ) -> Result<(ProcedureId, Option<Output>)> {
769 let context = self.create_context();
770 let procedure = CreateFlowProcedure::new(create_flow, query_context, context);
771 let procedure_with_id =
772 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
773
774 self.execute_procedure_and_wait(procedure_with_id).await
775 }
776
777 #[tracing::instrument(skip_all)]
779 pub async fn submit_drop_flow_task(
780 &self,
781 drop_flow: DropFlowTask,
782 procedure_context: ProcedureContext,
783 ) -> Result<(ProcedureId, Option<Output>)> {
784 let context = self.create_context();
785 let procedure = DropFlowProcedure::new(drop_flow, context);
786 let procedure_with_id =
787 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
788
789 self.execute_procedure_and_wait(procedure_with_id).await
790 }
791
792 #[tracing::instrument(skip_all)]
794 pub async fn submit_drop_view_task(
795 &self,
796 drop_view: DropViewTask,
797 procedure_context: ProcedureContext,
798 ) -> Result<(ProcedureId, Option<Output>)> {
799 let context = self.create_context();
800 let procedure = DropViewProcedure::new(drop_view, context);
801 let procedure_with_id =
802 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
803
804 self.execute_procedure_and_wait(procedure_with_id).await
805 }
806
807 #[tracing::instrument(skip_all)]
809 pub async fn submit_truncate_table_task(
810 &self,
811 truncate_table_task: TruncateTableTask,
812 table_info_value: DeserializedValueWithBytes<TableInfoValue>,
813 procedure_context: ProcedureContext,
814 ) -> Result<(ProcedureId, Option<Output>)> {
815 let context = self.create_context();
816 let procedure = TruncateTableProcedure::new(truncate_table_task, table_info_value, context);
817
818 let procedure_with_id =
819 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
820
821 self.execute_procedure_and_wait(procedure_with_id).await
822 }
823
824 #[tracing::instrument(skip_all)]
826 pub async fn submit_comment_on_task(
827 &self,
828 mut comment_on_task: CommentOnTask,
829 procedure_context: ProcedureContext,
830 ) -> Result<(ProcedureId, Option<Output>)> {
831 let context = self.create_context();
832 comment_on_task
833 .enrich_object_id(
834 context.table_metadata_manager.table_name_manager(),
835 context.flow_metadata_manager.flow_name_manager(),
836 )
837 .await?;
838 let procedure = CommentOnProcedure::new(comment_on_task, context);
839 let procedure_with_id =
840 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
841
842 self.execute_procedure_and_wait(procedure_with_id).await
843 }
844
845 async fn execute_procedure_and_wait(
847 &self,
848 procedure_with_id: ProcedureWithId,
849 ) -> Result<(ProcedureId, Option<Output>)> {
850 let procedure_id = procedure_with_id.id;
851
852 let mut watcher = self
853 .procedure_manager
854 .submit(procedure_with_id)
855 .await
856 .context(SubmitProcedureSnafu)?;
857
858 let output = watcher::wait(&mut watcher)
859 .await
860 .context(WaitProcedureSnafu)?;
861
862 Ok((procedure_id, output))
863 }
864
865 async fn submit_procedure(&self, procedure_with_id: ProcedureWithId) -> Result<ProcedureId> {
867 let procedure_id = procedure_with_id.id;
868 let _ = self
869 .procedure_manager
870 .submit(procedure_with_id)
871 .await
872 .context(SubmitProcedureSnafu)?;
873
874 Ok(procedure_id)
875 }
876
877 pub async fn submit_ddl_task(
878 &self,
879 context: ExecutorContext,
880 request: SubmitDdlTaskRequest,
881 ) -> Result<SubmitDdlTaskResponse> {
882 let ExecutorContext {
883 tracing_context,
884 query_context,
885 actor,
886 event_input,
887 } = context;
888 let query_context = query_context.context(UnsupportedSnafu {
889 operation: "submit_ddl_task without query context",
890 })?;
891 let procedure_context = ProcedureContext {
892 actor,
893 event_context: event_input
894 .map(|input| PersistentEventContext::from((input, query_context.protocol()))),
895 };
896 let span = tracing_context
897 .as_ref()
898 .map(TracingContext::from_w3c)
899 .unwrap_or_else(TracingContext::from_current_span)
900 .attach(tracing::info_span!("DdlManager::submit_ddl_task"));
901 let SubmitDdlTaskRequest {
902 wait,
903 timeout,
904 task,
905 } = request;
906 let ddl_options = DdlOptions { wait, timeout };
907 async move {
908 debug!("Submitting Ddl task: {:?}", task);
909 match task {
910 CreateTable(create_table_task) => {
911 handle_create_table_task(
912 self,
913 create_table_task,
914 query_context,
915 procedure_context,
916 )
917 .await
918 }
919 DropTable(drop_table_task) => {
920 handle_drop_table_task(self, drop_table_task, procedure_context).await
921 }
922 UndropTable(undrop_table_task) => {
923 handle_undrop_table_task(self, undrop_table_task, procedure_context).await
924 }
925 PurgeDroppedTable(purge_dropped_table_task) => {
926 handle_purge_dropped_table_task(
927 self,
928 purge_dropped_table_task,
929 procedure_context,
930 )
931 .await
932 }
933 AlterTable(alter_table_task) => {
934 handle_alter_table_task(self, alter_table_task, ddl_options, procedure_context)
935 .await
936 }
937 TruncateTable(truncate_table_task) => {
938 handle_truncate_table_task(self, truncate_table_task, procedure_context).await
939 }
940 CreateLogicalTables(create_table_tasks) => {
941 handle_create_logical_table_tasks(self, create_table_tasks, procedure_context)
942 .await
943 }
944 AlterLogicalTables(alter_table_tasks) => {
945 handle_alter_logical_table_tasks(self, alter_table_tasks, procedure_context)
946 .await
947 }
948 DropLogicalTables(_) => todo!(),
949 CreateDatabase(create_database_task) => {
950 handle_create_database_task(self, create_database_task, procedure_context).await
951 }
952 DropDatabase(drop_database_task) => {
953 handle_drop_database_task(self, drop_database_task, procedure_context).await
954 }
955 AlterDatabase(alter_database_task) => {
956 handle_alter_database_task(self, alter_database_task, procedure_context).await
957 }
958 CreateFlow(create_flow_task) => {
959 handle_create_flow_task(
960 self,
961 create_flow_task,
962 query_context,
963 procedure_context,
964 )
965 .await
966 }
967 DropFlow(drop_flow_task) => {
968 handle_drop_flow_task(self, drop_flow_task, procedure_context).await
969 }
970 CreateView(create_view_task) => {
971 handle_create_view_task(self, create_view_task, procedure_context).await
972 }
973 DropView(drop_view_task) => {
974 handle_drop_view_task(self, drop_view_task, procedure_context).await
975 }
976 CommentOn(comment_on_task) => {
977 handle_comment_on_task(self, comment_on_task, procedure_context).await
978 }
979 #[cfg(feature = "enterprise")]
980 CreateTrigger(create_trigger_task) => {
981 handle_create_trigger_task(
982 self,
983 create_trigger_task,
984 query_context,
985 procedure_context,
986 )
987 .await
988 }
989 #[cfg(feature = "enterprise")]
990 DropTrigger(drop_trigger_task) => {
991 handle_drop_trigger_task(
992 self,
993 drop_trigger_task,
994 query_context,
995 procedure_context,
996 )
997 .await
998 }
999 }
1000 }
1001 .trace(span)
1002 .await
1003 }
1004}
1005
1006async fn handle_truncate_table_task(
1007 ddl_manager: &DdlManager,
1008 truncate_table_task: TruncateTableTask,
1009 procedure_context: ProcedureContext,
1010) -> Result<SubmitDdlTaskResponse> {
1011 let table_id = truncate_table_task.table_id;
1012 let table_metadata_manager = &ddl_manager.table_metadata_manager();
1013 let table_ref = truncate_table_task.table_ref();
1014
1015 let table_info_value = table_metadata_manager
1016 .table_info_manager()
1017 .get(table_id)
1018 .await?
1019 .with_context(|| TableInfoNotFoundSnafu {
1020 table: table_ref.to_string(),
1021 })?;
1022 let physical_table_id = table_metadata_manager
1023 .table_route_manager()
1024 .get_physical_table_id(table_id)
1025 .await?;
1026 ensure!(
1027 physical_table_id == table_id,
1028 error::UnexpectedSnafu {
1029 err_msg: "Truncate table is only supported for physical tables."
1030 }
1031 );
1032
1033 let (id, _) = ddl_manager
1034 .submit_truncate_table_task(truncate_table_task, table_info_value, procedure_context)
1035 .await?;
1036
1037 info!("Table: {table_id} is truncated via procedure_id {id:?}");
1038
1039 Ok(SubmitDdlTaskResponse {
1040 key: id.to_string().into(),
1041 ..Default::default()
1042 })
1043}
1044
1045async fn handle_alter_table_task(
1046 ddl_manager: &DdlManager,
1047 alter_table_task: AlterTableTask,
1048 ddl_options: DdlOptions,
1049 procedure_context: ProcedureContext,
1050) -> Result<SubmitDdlTaskResponse> {
1051 let table_ref = alter_table_task.table_ref();
1052
1053 let table_id = ddl_manager
1054 .table_metadata_manager()
1055 .table_name_manager()
1056 .get(TableNameKey::new(
1057 table_ref.catalog,
1058 table_ref.schema,
1059 table_ref.table,
1060 ))
1061 .await?
1062 .with_context(|| TableNotFoundSnafu {
1063 table_name: table_ref.to_string(),
1064 })?
1065 .table_id();
1066
1067 let table_route_value = ddl_manager
1068 .table_metadata_manager()
1069 .table_route_manager()
1070 .table_route_storage()
1071 .get(table_id)
1072 .await?
1073 .context(TableRouteNotFoundSnafu { table_id })?;
1074 let annotation_family = match alter_table_task.alter_table.kind.as_ref() {
1078 Some(kind) => common_grpc_expr::annotation_alter_family(kind)
1079 .context(ConvertAlterTableRequestSnafu)?,
1080 None => None,
1081 };
1082 ensure!(
1083 table_route_value.is_physical()
1084 || annotation_family.is_some_and(|family| family.allows_logical_tables()),
1085 UnexpectedLogicalRouteTableSnafu {
1086 err_msg: format!("{:?} is a non-physical TableRouteValue.", table_ref),
1087 }
1088 );
1089
1090 let (id, _) = ddl_manager
1091 .submit_alter_table_task(table_id, alter_table_task, procedure_context, ddl_options)
1092 .await?;
1093
1094 info!("Table: {table_id} is altered via procedure_id {id:?}");
1095
1096 Ok(SubmitDdlTaskResponse {
1097 key: id.to_string().into(),
1098 ..Default::default()
1099 })
1100}
1101
1102async fn handle_drop_table_task(
1103 ddl_manager: &DdlManager,
1104 drop_table_task: DropTableTask,
1105 procedure_context: ProcedureContext,
1106) -> Result<SubmitDdlTaskResponse> {
1107 let table_id = drop_table_task.table_id;
1108 let (id, _) = ddl_manager
1109 .submit_drop_table_task(drop_table_task, procedure_context)
1110 .await?;
1111
1112 info!("Table: {table_id} is dropped via procedure_id {id:?}");
1113
1114 Ok(SubmitDdlTaskResponse {
1115 key: id.to_string().into(),
1116 ..Default::default()
1117 })
1118}
1119
1120async fn handle_undrop_table_task(
1121 ddl_manager: &DdlManager,
1122 undrop_table_task: UndropTableTask,
1123 procedure_context: ProcedureContext,
1124) -> Result<SubmitDdlTaskResponse> {
1125 let table_id = undrop_table_task.table_id;
1126 let (id, _) = ddl_manager
1127 .submit_undrop_table_task(undrop_table_task, procedure_context)
1128 .await?;
1129
1130 info!("Table: {table_id} is undropped via procedure_id {id:?}");
1131
1132 Ok(SubmitDdlTaskResponse {
1133 key: id.to_string().into(),
1134 ..Default::default()
1135 })
1136}
1137
1138async fn handle_purge_dropped_table_task(
1139 ddl_manager: &DdlManager,
1140 purge_dropped_table_task: PurgeDroppedTableTask,
1141 procedure_context: ProcedureContext,
1142) -> Result<SubmitDdlTaskResponse> {
1143 let (id, _) = ddl_manager
1144 .submit_purge_dropped_table_task(purge_dropped_table_task, procedure_context)
1145 .await?;
1146
1147 info!("Dropped table is purged via procedure_id {id:?}");
1148
1149 Ok(SubmitDdlTaskResponse {
1150 key: id.to_string().into(),
1151 ..Default::default()
1152 })
1153}
1154
1155async fn handle_create_table_task(
1156 ddl_manager: &DdlManager,
1157 create_table_task: CreateTableTask,
1158 query_context: QueryContext,
1159 procedure_context: ProcedureContext,
1160) -> Result<SubmitDdlTaskResponse> {
1161 let (id, output) = ddl_manager
1162 .submit_create_table_task(create_table_task, query_context, procedure_context)
1163 .await?;
1164
1165 let procedure_id = id.to_string();
1166 let output = output.context(ProcedureOutputSnafu {
1167 procedure_id: &procedure_id,
1168 err_msg: "empty output",
1169 })?;
1170 let table_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1171 procedure_id: &procedure_id,
1172 err_msg: "downcast to `u32`",
1173 })?);
1174 info!("Table: {table_id} is created via procedure_id {id:?}");
1175
1176 Ok(SubmitDdlTaskResponse {
1177 key: procedure_id.into(),
1178 table_ids: vec![table_id],
1179 })
1180}
1181
1182async fn handle_create_logical_table_tasks(
1183 ddl_manager: &DdlManager,
1184 create_table_tasks: Vec<CreateTableTask>,
1185 procedure_context: ProcedureContext,
1186) -> Result<SubmitDdlTaskResponse> {
1187 ensure!(
1188 !create_table_tasks.is_empty(),
1189 EmptyDdlTasksSnafu {
1190 name: "create logical tables"
1191 }
1192 );
1193 let physical_table_id = utils::check_and_get_physical_table_id(
1194 ddl_manager.table_metadata_manager(),
1195 &create_table_tasks,
1196 )
1197 .await?;
1198 let num_logical_tables = create_table_tasks.len();
1199
1200 let (id, output) = ddl_manager
1201 .submit_create_logical_table_tasks(create_table_tasks, physical_table_id, procedure_context)
1202 .await?;
1203
1204 info!(
1205 "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is created via procedure_id {id:?}"
1206 );
1207
1208 let procedure_id = id.to_string();
1209 let output = output.context(ProcedureOutputSnafu {
1210 procedure_id: &procedure_id,
1211 err_msg: "empty output",
1212 })?;
1213 let table_ids = output
1214 .downcast_ref::<Vec<TableId>>()
1215 .context(ProcedureOutputSnafu {
1216 procedure_id: &procedure_id,
1217 err_msg: "downcast to `Vec<TableId>`",
1218 })?
1219 .clone();
1220
1221 Ok(SubmitDdlTaskResponse {
1222 key: procedure_id.into(),
1223 table_ids,
1224 })
1225}
1226
1227async fn handle_create_database_task(
1228 ddl_manager: &DdlManager,
1229 create_database_task: CreateDatabaseTask,
1230 procedure_context: ProcedureContext,
1231) -> Result<SubmitDdlTaskResponse> {
1232 let catalog = create_database_task.catalog.clone();
1233 let schema = create_database_task.schema.clone();
1234 let (id, _) = ddl_manager
1235 .submit_create_database(create_database_task, procedure_context)
1236 .await?;
1237
1238 let procedure_id = id.to_string();
1239 info!(
1240 "Database {}.{} is created via procedure_id {id:?}",
1241 catalog, schema
1242 );
1243
1244 Ok(SubmitDdlTaskResponse {
1245 key: procedure_id.into(),
1246 ..Default::default()
1247 })
1248}
1249
1250async fn handle_drop_database_task(
1251 ddl_manager: &DdlManager,
1252 drop_database_task: DropDatabaseTask,
1253 procedure_context: ProcedureContext,
1254) -> Result<SubmitDdlTaskResponse> {
1255 let (id, _) = ddl_manager
1256 .submit_drop_database(drop_database_task.clone(), procedure_context)
1257 .await?;
1258
1259 let procedure_id = id.to_string();
1260 info!(
1261 "Database {}.{} is dropped via procedure_id {id:?}",
1262 drop_database_task.catalog, drop_database_task.schema
1263 );
1264
1265 Ok(SubmitDdlTaskResponse {
1266 key: procedure_id.into(),
1267 ..Default::default()
1268 })
1269}
1270
1271async fn handle_alter_database_task(
1272 ddl_manager: &DdlManager,
1273 alter_database_task: AlterDatabaseTask,
1274 procedure_context: ProcedureContext,
1275) -> Result<SubmitDdlTaskResponse> {
1276 let (id, _) = ddl_manager
1277 .submit_alter_database(alter_database_task.clone(), procedure_context)
1278 .await?;
1279
1280 let procedure_id = id.to_string();
1281 info!(
1282 "Database {}.{} is altered via procedure_id {id:?}",
1283 alter_database_task.catalog(),
1284 alter_database_task.schema()
1285 );
1286
1287 Ok(SubmitDdlTaskResponse {
1288 key: procedure_id.into(),
1289 ..Default::default()
1290 })
1291}
1292
1293async fn handle_drop_flow_task(
1294 ddl_manager: &DdlManager,
1295 drop_flow_task: DropFlowTask,
1296 procedure_context: ProcedureContext,
1297) -> Result<SubmitDdlTaskResponse> {
1298 let (id, _) = ddl_manager
1299 .submit_drop_flow_task(drop_flow_task.clone(), procedure_context)
1300 .await?;
1301
1302 let procedure_id = id.to_string();
1303 info!(
1304 "Flow {}.{}({}) is dropped via procedure_id {id:?}",
1305 drop_flow_task.catalog_name, drop_flow_task.flow_name, drop_flow_task.flow_id,
1306 );
1307
1308 Ok(SubmitDdlTaskResponse {
1309 key: procedure_id.into(),
1310 ..Default::default()
1311 })
1312}
1313
1314#[cfg(feature = "enterprise")]
1315async fn handle_drop_trigger_task(
1316 ddl_manager: &DdlManager,
1317 drop_trigger_task: DropTriggerTask,
1318 query_context: QueryContext,
1319 procedure_context: ProcedureContext,
1320) -> Result<SubmitDdlTaskResponse> {
1321 let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1322 use crate::error::UnsupportedSnafu;
1323
1324 return UnsupportedSnafu {
1325 operation: "drop trigger",
1326 }
1327 .fail();
1328 };
1329
1330 m.drop_trigger(
1331 drop_trigger_task,
1332 ddl_manager.procedure_manager.clone(),
1333 ddl_manager.ddl_context.clone(),
1334 query_context,
1335 procedure_context,
1336 )
1337 .await
1338}
1339
1340async fn handle_drop_view_task(
1341 ddl_manager: &DdlManager,
1342 drop_view_task: DropViewTask,
1343 procedure_context: ProcedureContext,
1344) -> Result<SubmitDdlTaskResponse> {
1345 let (id, _) = ddl_manager
1346 .submit_drop_view_task(drop_view_task.clone(), procedure_context)
1347 .await?;
1348
1349 let procedure_id = id.to_string();
1350 info!(
1351 "View {}({}) is dropped via procedure_id {id:?}",
1352 drop_view_task.table_ref(),
1353 drop_view_task.view_id,
1354 );
1355
1356 Ok(SubmitDdlTaskResponse {
1357 key: procedure_id.into(),
1358 ..Default::default()
1359 })
1360}
1361
1362async fn handle_create_flow_task(
1363 ddl_manager: &DdlManager,
1364 create_flow_task: CreateFlowTask,
1365 query_context: QueryContext,
1366 procedure_context: ProcedureContext,
1367) -> Result<SubmitDdlTaskResponse> {
1368 let (id, output) = ddl_manager
1369 .submit_create_flow_task(create_flow_task.clone(), query_context, procedure_context)
1370 .await?;
1371
1372 let procedure_id = id.to_string();
1373 let output = output.context(ProcedureOutputSnafu {
1374 procedure_id: &procedure_id,
1375 err_msg: "empty output",
1376 })?;
1377 let flow_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1378 procedure_id: &procedure_id,
1379 err_msg: "downcast to `u32`",
1380 })?);
1381 if !create_flow_task.or_replace {
1382 info!(
1383 "Flow {}.{}({flow_id}) is created via procedure_id {id:?}",
1384 create_flow_task.catalog_name, create_flow_task.flow_name,
1385 );
1386 } else {
1387 info!(
1388 "Flow {}.{}({flow_id}) is replaced via procedure_id {id:?}",
1389 create_flow_task.catalog_name, create_flow_task.flow_name,
1390 );
1391 }
1392
1393 Ok(SubmitDdlTaskResponse {
1394 key: procedure_id.into(),
1395 ..Default::default()
1396 })
1397}
1398
1399#[cfg(feature = "enterprise")]
1400async fn handle_create_trigger_task(
1401 ddl_manager: &DdlManager,
1402 create_trigger_task: CreateTriggerTask,
1403 query_context: QueryContext,
1404 procedure_context: ProcedureContext,
1405) -> Result<SubmitDdlTaskResponse> {
1406 let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1407 use crate::error::UnsupportedSnafu;
1408
1409 return UnsupportedSnafu {
1410 operation: "create trigger",
1411 }
1412 .fail();
1413 };
1414
1415 m.create_trigger(
1416 create_trigger_task,
1417 ddl_manager.procedure_manager.clone(),
1418 ddl_manager.ddl_context.clone(),
1419 query_context,
1420 procedure_context,
1421 )
1422 .await
1423}
1424
1425async fn handle_alter_logical_table_tasks(
1426 ddl_manager: &DdlManager,
1427 alter_table_tasks: Vec<AlterTableTask>,
1428 procedure_context: ProcedureContext,
1429) -> Result<SubmitDdlTaskResponse> {
1430 ensure!(
1431 !alter_table_tasks.is_empty(),
1432 EmptyDdlTasksSnafu {
1433 name: "alter logical tables"
1434 }
1435 );
1436
1437 let first_table = TableNameKey {
1439 catalog: &alter_table_tasks[0].alter_table.catalog_name,
1440 schema: &alter_table_tasks[0].alter_table.schema_name,
1441 table: &alter_table_tasks[0].alter_table.table_name,
1442 };
1443 let physical_table_id =
1444 utils::get_physical_table_id(ddl_manager.table_metadata_manager(), first_table).await?;
1445 let num_logical_tables = alter_table_tasks.len();
1446
1447 let (id, _) = ddl_manager
1448 .submit_alter_logical_table_tasks(alter_table_tasks, physical_table_id, procedure_context)
1449 .await?;
1450
1451 info!(
1452 "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is altered via procedure_id {id:?}"
1453 );
1454
1455 let procedure_id = id.to_string();
1456
1457 Ok(SubmitDdlTaskResponse {
1458 key: procedure_id.into(),
1459 ..Default::default()
1460 })
1461}
1462
1463async fn handle_create_view_task(
1465 ddl_manager: &DdlManager,
1466 create_view_task: CreateViewTask,
1467 procedure_context: ProcedureContext,
1468) -> Result<SubmitDdlTaskResponse> {
1469 let (id, output) = ddl_manager
1470 .submit_create_view_task(create_view_task, procedure_context)
1471 .await?;
1472
1473 let procedure_id = id.to_string();
1474 let output = output.context(ProcedureOutputSnafu {
1475 procedure_id: &procedure_id,
1476 err_msg: "empty output",
1477 })?;
1478 let view_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1479 procedure_id: &procedure_id,
1480 err_msg: "downcast to `u32`",
1481 })?);
1482 info!("View: {view_id} is created via procedure_id {id:?}");
1483
1484 Ok(SubmitDdlTaskResponse {
1485 key: procedure_id.into(),
1486 table_ids: vec![view_id],
1487 })
1488}
1489
1490async fn handle_comment_on_task(
1491 ddl_manager: &DdlManager,
1492 comment_on_task: CommentOnTask,
1493 procedure_context: ProcedureContext,
1494) -> Result<SubmitDdlTaskResponse> {
1495 let (id, _) = ddl_manager
1496 .submit_comment_on_task(comment_on_task.clone(), procedure_context)
1497 .await?;
1498
1499 let procedure_id = id.to_string();
1500 info!(
1501 "Comment on {}.{}.{} is updated via procedure_id {id:?}",
1502 comment_on_task.catalog_name, comment_on_task.schema_name, comment_on_task.object_name
1503 );
1504
1505 Ok(SubmitDdlTaskResponse {
1506 key: procedure_id.into(),
1507 ..Default::default()
1508 })
1509}
1510
1511#[cfg(test)]
1512mod tests {
1513 use std::sync::Arc;
1514 #[cfg(feature = "enterprise")]
1515 use std::sync::Mutex;
1516 use std::time::Duration;
1517
1518 #[cfg(feature = "enterprise")]
1519 use common_base::protocol::Channel;
1520 use common_error::ext::BoxedError;
1521 #[cfg(feature = "enterprise")]
1522 use common_error::ext::ErrorExt;
1523 #[cfg(feature = "enterprise")]
1524 use common_error::status_code::StatusCode;
1525 #[cfg(feature = "enterprise")]
1526 use common_event_recorder::{PersistentEventContext, ProcedureEventInput, TriggerReason};
1527 use common_procedure::local::LocalManager;
1528 use common_procedure::test_util::InMemoryPoisonStore;
1529 use common_procedure::{
1530 BoxedProcedure, ProcedureContext, ProcedureManager, ProcedureManagerRef,
1531 };
1532 use store_api::storage::TableId;
1533 use table::table_name::TableName;
1534
1535 use super::DdlManager;
1536 use crate::cache_invalidator::DummyCacheInvalidator;
1537 use crate::ddl::alter_table::AlterTableProcedure;
1538 use crate::ddl::create_database::{
1539 AtomicCreateOutcome, CreateDatabaseMetadataCommitter, CreateDatabaseProcedure,
1540 };
1541 use crate::ddl::create_table::CreateTableProcedure;
1542 use crate::ddl::drop_table::DropTableProcedure;
1543 use crate::ddl::flow_meta::FlowMetadataAllocator;
1544 use crate::ddl::table_meta::TableMetadataAllocator;
1545 use crate::ddl::truncate_table::TruncateTableProcedure;
1546 use crate::ddl::{DdlContext, NoopRegionFailureDetectorControl};
1547 use crate::ddl_manager::{RepartitionProcedureFactory, RepartitionSource};
1548 use crate::key::TableMetadataManager;
1549 use crate::key::flow::FlowMetadataManager;
1550 use crate::kv_backend::memory::MemoryKvBackend;
1551 use crate::node_manager::{DatanodeManager, DatanodeRef, FlownodeManager, FlownodeRef};
1552 use crate::peer::Peer;
1553 use crate::procedure_executor::ExecutorContext;
1554 use crate::region_keeper::MemoryRegionKeeper;
1555 use crate::region_registry::LeaderRegionRegistry;
1556 #[cfg(feature = "enterprise")]
1557 use crate::rpc::ddl::trigger::{CreateTriggerTask, DropTriggerTask};
1558 use crate::rpc::ddl::{CreatorGrantIntent, UndropTableTask};
1559 #[cfg(not(feature = "enterprise"))]
1560 use crate::rpc::ddl::{DdlTask, PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest};
1561 #[cfg(feature = "enterprise")]
1562 use crate::rpc::ddl::{DdlTask, QueryContext, SubmitDdlTaskRequest};
1563 use crate::sequence::SequenceBuilder;
1564 use crate::state_store::KvStateStore;
1565 use crate::test_util::{MockDatanodeManager, new_ddl_context};
1566 use crate::wal_provider::WalProvider;
1567
1568 pub struct DummyDatanodeManager;
1570
1571 #[async_trait::async_trait]
1572 impl DatanodeManager for DummyDatanodeManager {
1573 async fn datanode(&self, _datanode: &Peer) -> DatanodeRef {
1574 unimplemented!()
1575 }
1576 }
1577
1578 #[async_trait::async_trait]
1579 impl FlownodeManager for DummyDatanodeManager {
1580 async fn flownode(&self, _node: &Peer) -> FlownodeRef {
1581 unimplemented!()
1582 }
1583 }
1584
1585 struct DummyRepartitionProcedureFactory;
1586
1587 #[async_trait::async_trait]
1588 impl RepartitionProcedureFactory for DummyRepartitionProcedureFactory {
1589 fn create(
1590 &self,
1591 _ddl_ctx: &DdlContext,
1592 _table_name: TableName,
1593 _table_id: TableId,
1594 _source: RepartitionSource,
1595 _to_exprs: Vec<String>,
1596 _timeout: Option<Duration>,
1597 ) -> std::result::Result<BoxedProcedure, BoxedError> {
1598 unimplemented!()
1599 }
1600
1601 fn register_loaders(
1602 &self,
1603 _ddl_ctx: &DdlContext,
1604 _procedure_manager: &ProcedureManagerRef,
1605 ) -> std::result::Result<(), BoxedError> {
1606 Ok(())
1607 }
1608
1609 async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError> {
1610 Ok(())
1611 }
1612 }
1613
1614 struct LoaderCommitter;
1615
1616 #[async_trait::async_trait]
1617 impl CreateDatabaseMetadataCommitter for LoaderCommitter {
1618 async fn commit(
1619 &self,
1620 _catalog: &str,
1621 _schema: &str,
1622 _value: &crate::key::schema_name::SchemaNameValue,
1623 _creator: &CreatorGrantIntent,
1624 ) -> std::result::Result<AtomicCreateOutcome, BoxedError> {
1625 unreachable!()
1626 }
1627 }
1628
1629 #[cfg(feature = "enterprise")]
1630 #[derive(Default)]
1631 struct RecordingTriggerDdlManager {
1632 procedure_contexts: Mutex<Vec<ProcedureContext>>,
1633 }
1634
1635 #[cfg(feature = "enterprise")]
1636 #[async_trait::async_trait]
1637 impl super::TriggerDdlManager for RecordingTriggerDdlManager {
1638 async fn create_trigger(
1639 &self,
1640 _create_trigger_task: CreateTriggerTask,
1641 _procedure_manager: ProcedureManagerRef,
1642 _ddl_context: DdlContext,
1643 _query_context: QueryContext,
1644 procedure_context: ProcedureContext,
1645 ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1646 self.procedure_contexts
1647 .lock()
1648 .unwrap()
1649 .push(procedure_context);
1650 Ok(Default::default())
1651 }
1652
1653 async fn drop_trigger(
1654 &self,
1655 _drop_trigger_task: DropTriggerTask,
1656 _procedure_manager: ProcedureManagerRef,
1657 _ddl_context: DdlContext,
1658 _query_context: QueryContext,
1659 procedure_context: ProcedureContext,
1660 ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1661 self.procedure_contexts
1662 .lock()
1663 .unwrap()
1664 .push(procedure_context);
1665 Ok(Default::default())
1666 }
1667
1668 fn as_any(&self) -> &dyn std::any::Any {
1669 self
1670 }
1671 }
1672
1673 #[test]
1674 fn test_generic_loader_captures_configured_committer() {
1675 let mut context = new_ddl_context(Arc::new(MockDatanodeManager::new(())));
1676 let committer = Arc::new(LoaderCommitter);
1677 context.create_database_metadata_committer = Some(committer.clone());
1678 let loader = procedure_loader!(CreateDatabaseProcedure)[0].1(context);
1679 assert_eq!(Arc::strong_count(&committer), 2);
1680
1681 let loaded = loader(
1682 r#"{
1683 "state":"CreateMetadata",
1684 "catalog":"greptime",
1685 "schema":"metrics",
1686 "create_if_not_exists":false,
1687 "options":{},
1688 "creator":{"username":"alice","created_at_ns":42}
1689 }"#,
1690 )
1691 .unwrap();
1692
1693 assert_eq!(loaded.type_name(), "metasrv-procedure::CreateDatabase");
1694 assert_eq!(Arc::strong_count(&committer), 3);
1695 drop(loaded);
1696 assert_eq!(Arc::strong_count(&committer), 2);
1697 }
1698
1699 #[test]
1700 fn test_register_loaders() {
1701 let kv_backend = Arc::new(MemoryKvBackend::new());
1702 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1703 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1704 Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1705 Arc::new(WalProvider::default()),
1706 ));
1707 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1708 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1709 Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1710 ));
1711
1712 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1713 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1714 let procedure_manager = Arc::new(LocalManager::new(
1715 Default::default(),
1716 state_store,
1717 poison_manager,
1718 None,
1719 None,
1720 ));
1721
1722 let ddl_manager = DdlManager::new(
1723 DdlContext {
1724 node_manager: Arc::new(DummyDatanodeManager),
1725 cache_invalidator: Arc::new(DummyCacheInvalidator),
1726 table_metadata_manager,
1727 table_metadata_allocator,
1728 flow_metadata_manager,
1729 flow_metadata_allocator,
1730 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1731 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1732 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1733 soft_drop_enabled: false,
1734 soft_drop_retention: None,
1735 create_database_metadata_committer: None,
1736 },
1737 procedure_manager.clone(),
1738 Arc::new(DummyRepartitionProcedureFactory),
1739 );
1740 ddl_manager.register_loaders().unwrap();
1741
1742 let expected_loaders = vec![
1743 CreateTableProcedure::TYPE_NAME,
1744 AlterTableProcedure::TYPE_NAME,
1745 DropTableProcedure::TYPE_NAME,
1746 TruncateTableProcedure::TYPE_NAME,
1747 ];
1748
1749 for loader in expected_loaders {
1750 assert!(procedure_manager.contains_loader(loader));
1751 }
1752
1753 let soft_drop_loaders = [
1754 "metasrv-procedure::UndropTable",
1755 "metasrv-procedure::PurgeDroppedTable",
1756 "metasrv-procedure::PurgeExpiredDroppedTable",
1757 ];
1758 for loader in soft_drop_loaders {
1759 assert_eq!(
1760 cfg!(feature = "enterprise"),
1761 procedure_manager.contains_loader(loader)
1762 );
1763 }
1764 }
1765
1766 fn build_soft_drop_test_ddl_manager() -> DdlManager {
1767 let kv_backend = Arc::new(MemoryKvBackend::new());
1768 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1769 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1770 Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1771 Arc::new(WalProvider::default()),
1772 ));
1773 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1774 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1775 Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1776 ));
1777
1778 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1779 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1780 let procedure_manager = Arc::new(LocalManager::new(
1781 Default::default(),
1782 state_store,
1783 poison_manager,
1784 None,
1785 None,
1786 ));
1787
1788 DdlManager::new(
1789 DdlContext {
1790 node_manager: Arc::new(DummyDatanodeManager),
1791 cache_invalidator: Arc::new(DummyCacheInvalidator),
1792 table_metadata_manager,
1793 table_metadata_allocator,
1794 flow_metadata_manager,
1795 flow_metadata_allocator,
1796 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1797 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1798 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1799 soft_drop_enabled: true,
1800 soft_drop_retention: Some(std::time::Duration::from_secs(1)),
1801 create_database_metadata_committer: None,
1802 },
1803 procedure_manager,
1804 Arc::new(DummyRepartitionProcedureFactory),
1805 )
1806 }
1807
1808 #[cfg(feature = "enterprise")]
1809 #[tokio::test]
1810 async fn test_trigger_ddl_forwards_procedure_context() {
1811 let trigger_ddl_manager = Arc::new(RecordingTriggerDdlManager::default());
1812 let ddl_manager = build_soft_drop_test_ddl_manager()
1813 .with_trigger_ddl_manager(trigger_ddl_manager.clone());
1814 let procedure_context = ProcedureContext {
1815 actor: Some("test-user".to_string()),
1816 event_context: Some(
1817 PersistentEventContext::new(TriggerReason::Manual).with_protocol("mysql"),
1818 ),
1819 };
1820 let executor_context = || ExecutorContext {
1821 query_context: Some(QueryContext {
1822 channel: Channel::Mysql as u8,
1823 ..Default::default()
1824 }),
1825 actor: Some("test-user".to_string()),
1826 event_input: Some(ProcedureEventInput::new(TriggerReason::Manual)),
1827 ..Default::default()
1828 };
1829
1830 ddl_manager
1831 .submit_ddl_task(
1832 executor_context(),
1833 SubmitDdlTaskRequest::new(DdlTask::CreateTrigger(CreateTriggerTask {
1834 catalog_name: "greptime".to_string(),
1835 trigger_name: "test_trigger".to_string(),
1836 if_not_exists: false,
1837 sql: "SELECT 1".to_string(),
1838 channels: vec![],
1839 labels: Default::default(),
1840 annotations: Default::default(),
1841 interval: Duration::from_secs(1),
1842 raw_interval_expr: None,
1843 r#for: None,
1844 for_raw_expr: None,
1845 keep_firing_for: None,
1846 keep_firing_for_raw_expr: None,
1847 })),
1848 )
1849 .await
1850 .unwrap();
1851 ddl_manager
1852 .submit_ddl_task(
1853 executor_context(),
1854 SubmitDdlTaskRequest::new(DdlTask::DropTrigger(DropTriggerTask {
1855 catalog_name: "greptime".to_string(),
1856 trigger_name: "test_trigger".to_string(),
1857 drop_if_exists: false,
1858 })),
1859 )
1860 .await
1861 .unwrap();
1862
1863 assert_eq!(
1864 *trigger_ddl_manager.procedure_contexts.lock().unwrap(),
1865 vec![procedure_context.clone(), procedure_context]
1866 );
1867 }
1868
1869 #[cfg(feature = "enterprise")]
1870 #[tokio::test]
1871 async fn test_submit_undrop_missing_tombstone_returns_table_not_found_directly() {
1872 let ddl_manager = build_soft_drop_test_ddl_manager();
1873
1874 let err = ddl_manager
1875 .submit_undrop_table_task(
1876 UndropTableTask { table_id: 1024 },
1877 ProcedureContext::default(),
1878 )
1879 .await
1880 .unwrap_err();
1881
1882 assert_eq!(err.status_code(), StatusCode::TableNotFound);
1883 assert!(matches!(err, crate::error::Error::TableNotFound { .. }));
1884 }
1885
1886 #[cfg(not(feature = "enterprise"))]
1887 #[tokio::test]
1888 async fn test_submit_undrop_and_purge_rejected_in_non_enterprise_build() {
1889 let ddl_manager = build_soft_drop_test_ddl_manager();
1890
1891 let err = ddl_manager
1892 .submit_undrop_table_task(
1893 UndropTableTask { table_id: 1024 },
1894 ProcedureContext::default(),
1895 )
1896 .await
1897 .unwrap_err();
1898 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1899
1900 let err = ddl_manager
1901 .submit_purge_dropped_table_task(
1902 PurgeDroppedTableTask { table_id: 1024 },
1903 ProcedureContext::default(),
1904 )
1905 .await
1906 .unwrap_err();
1907 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1908
1909 let err = ddl_manager
1910 .submit_expired_purge_dropped_table_task(PurgeDroppedTableTask { table_id: 1024 })
1911 .await
1912 .unwrap_err();
1913 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1914
1915 for task in [
1916 DdlTask::UndropTable(UndropTableTask { table_id: 1024 }),
1917 DdlTask::PurgeDroppedTable(PurgeDroppedTableTask { table_id: 1024 }),
1918 ] {
1919 let err = ddl_manager
1920 .submit_ddl_task(
1921 ExecutorContext {
1922 query_context: Some(QueryContext::default()),
1923 ..Default::default()
1924 },
1925 SubmitDdlTaskRequest::new(task),
1926 )
1927 .await
1928 .unwrap_err();
1929 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1930 }
1931 }
1932
1933 async fn ddl_manager_with_context(ddl_context: DdlContext) -> DdlManager {
1934 let kv_backend = Arc::new(MemoryKvBackend::new());
1935 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1936 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1937 let procedure_manager = Arc::new(LocalManager::new(
1938 Default::default(),
1939 state_store,
1940 poison_manager,
1941 None,
1942 None,
1943 ));
1944 procedure_manager.start().await.unwrap();
1945 let ddl_manager = DdlManager::new(
1946 ddl_context,
1947 procedure_manager,
1948 Arc::new(DummyRepartitionProcedureFactory),
1949 );
1950 ddl_manager.register_loaders().unwrap();
1951 ddl_manager
1952 }
1953
1954 fn set_options_expr(table_name: &str, options: &[(&str, &str)]) -> api::v1::AlterTableExpr {
1955 api::v1::AlterTableExpr {
1956 catalog_name: common_catalog::consts::DEFAULT_CATALOG_NAME.to_string(),
1957 schema_name: common_catalog::consts::DEFAULT_SCHEMA_NAME.to_string(),
1958 table_name: table_name.to_string(),
1959 kind: Some(api::v1::alter_table_expr::Kind::SetTableOptions(
1960 api::v1::SetTableOptions {
1961 table_options: options
1962 .iter()
1963 .map(|(key, value)| api::v1::Option {
1964 key: key.to_string(),
1965 value: value.to_string(),
1966 })
1967 .collect(),
1968 },
1969 )),
1970 }
1971 }
1972
1973 #[tokio::test]
1974 async fn test_logical_table_annotation_alter_routing() {
1975 let (tx, mut rx) = tokio::sync::mpsc::channel(8);
1976 let node_manager = Arc::new(crate::test_util::MockDatanodeManager::new(
1977 crate::ddl::test_util::datanode_handler::DatanodeWatcher::new(tx),
1978 ));
1979 let ddl_context = crate::test_util::new_ddl_context(node_manager);
1980 let phy_id = crate::ddl::test_util::create_physical_table(&ddl_context, "phy").await;
1981 let logical_id =
1982 crate::ddl::test_util::create_logical_table(ddl_context.clone(), phy_id, "logical")
1983 .await;
1984 let ddl_manager = ddl_manager_with_context(ddl_context.clone()).await;
1985
1986 ddl_manager
1989 .submit_ddl_task(
1990 ExecutorContext {
1991 query_context: Some(QueryContext::default()),
1992 ..Default::default()
1993 },
1994 SubmitDdlTaskRequest::new(DdlTask::new_alter_table(set_options_expr(
1995 "logical",
1996 &[("greptime.semantic.signal_type", "metric")],
1997 ))),
1998 )
1999 .await
2000 .unwrap();
2001 rx.try_recv().unwrap_err();
2002 let table_info = ddl_manager
2003 .table_metadata_manager()
2004 .table_info_manager()
2005 .get(logical_id)
2006 .await
2007 .unwrap()
2008 .unwrap()
2009 .into_inner()
2010 .table_info;
2011 assert_eq!(
2012 table_info
2013 .meta
2014 .options
2015 .extra_options
2016 .get("greptime.semantic.signal_type"),
2017 Some(&"metric".to_string())
2018 );
2019
2020 let err = ddl_manager
2022 .submit_ddl_task(
2023 ExecutorContext {
2024 query_context: Some(QueryContext::default()),
2025 ..Default::default()
2026 },
2027 SubmitDdlTaskRequest::new(DdlTask::new_alter_table(set_options_expr(
2028 "logical",
2029 &[("greptime.semantic.source", "prometheus"), ("ttl", "7d")],
2030 ))),
2031 )
2032 .await
2033 .unwrap_err();
2034 let msg = common_error::ext::ErrorExt::output_msg(&err);
2035 assert!(msg.contains("must be altered separately"), "{msg}");
2036
2037 let err = ddl_manager
2040 .submit_ddl_task(
2041 ExecutorContext {
2042 query_context: Some(QueryContext::default()),
2043 ..Default::default()
2044 },
2045 SubmitDdlTaskRequest::new(DdlTask::new_alter_table(set_options_expr(
2046 "logical",
2047 &[("repartition.column.hint", "host")],
2048 ))),
2049 )
2050 .await
2051 .unwrap_err();
2052 let msg = common_error::ext::ErrorExt::output_msg(&err);
2053 assert!(msg.contains("non-physical TableRouteValue"), "{msg}");
2054 }
2055}