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_sets_skip_wal};
39use crate::ddl::comment_on::CommentOnProcedure;
40use crate::ddl::create_database::{CreateDatabaseMetadataCommitterRef, CreateDatabaseProcedure};
41use crate::ddl::create_flow::{
42 CreateFlowProcedure, FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY,
43 FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_SEQUENCE_RANGE,
44};
45use crate::ddl::create_logical_tables::CreateLogicalTablesProcedure;
46use crate::ddl::create_table::CreateTableProcedure;
47use crate::ddl::create_view::CreateViewProcedure;
48use crate::ddl::drop_database::DropDatabaseProcedure;
49use crate::ddl::drop_flow::DropFlowProcedure;
50use crate::ddl::drop_table::DropTableProcedure;
51use crate::ddl::drop_view::DropViewProcedure;
52#[cfg(feature = "enterprise")]
53use crate::ddl::purge_dropped_table::PurgeDroppedTableProcedure;
54use crate::ddl::truncate_table::TruncateTableProcedure;
55#[cfg(feature = "enterprise")]
56use crate::ddl::undrop_table::UndropTableProcedure;
57use crate::ddl::{DdlContext, utils};
58use crate::error::{
59 self, ConvertAlterTableRequestSnafu, CreateRepartitionProcedureSnafu, EmptyDdlTasksSnafu,
60 PersistRepartitionGcRequirementSnafu, ProcedureOutputSnafu, RegisterProcedureLoaderSnafu,
61 RegisterRepartitionProcedureLoaderSnafu, Result, SubmitProcedureSnafu, TableInfoNotFoundSnafu,
62 TableNotFoundSnafu, TableRouteNotFoundSnafu, UnexpectedLogicalRouteTableSnafu,
63 UnsupportedSnafu, WaitProcedureSnafu,
64};
65use crate::key::table_info::TableInfoValue;
66use crate::key::table_name::TableNameKey;
67use crate::key::{DeserializedValueWithBytes, TableMetadataManagerRef};
68use crate::procedure_executor::ExecutorContext;
69#[cfg(feature = "enterprise")]
70use crate::rpc::ddl::DdlTask::CreateTrigger;
71#[cfg(feature = "enterprise")]
72use crate::rpc::ddl::DdlTask::DropTrigger;
73use crate::rpc::ddl::DdlTask::{
74 AlterDatabase, AlterLogicalTables, AlterTable, CommentOn, CreateDatabase, CreateFlow,
75 CreateLogicalTables, CreateTable, CreateView, DropDatabase, DropFlow, DropLogicalTables,
76 DropTable, DropView, PurgeDroppedTable, TruncateTable, UndropTable,
77};
78#[cfg(feature = "enterprise")]
79use crate::rpc::ddl::trigger::CreateTriggerTask;
80#[cfg(feature = "enterprise")]
81use crate::rpc::ddl::trigger::DropTriggerTask;
82use crate::rpc::ddl::{
83 AlterDatabaseTask, AlterTableTask, CommentOnTask, CreateDatabaseTask, CreateFlowTask,
84 CreateTableTask, CreateViewTask, DropDatabaseTask, DropFlowTask, DropTableTask, DropViewTask,
85 PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest, SubmitDdlTaskResponse,
86 TruncateTableTask, UndropTableTask,
87};
88
89const MAX_REGION_ROUTE_CHANGE_RETRIES: usize = 3;
90
91#[async_trait::async_trait]
93pub trait DdlManagerConfigurator<C>: Send + Sync {
94 async fn configure(
96 &self,
97 ddl_manager: DdlManager,
98 ctx: C,
99 ) -> std::result::Result<DdlManager, BoxedError>;
100}
101
102pub type DdlManagerConfiguratorRef<C> = Arc<dyn DdlManagerConfigurator<C>>;
103
104pub type DdlManagerRef = Arc<DdlManager>;
105
106pub type BoxedProcedureLoaderFactory = dyn Fn(DdlContext) -> BoxedProcedureLoader;
107
108#[derive(Builder)]
110pub struct DdlManager {
111 ddl_context: DdlContext,
112 procedure_manager: ProcedureManagerRef,
113 repartition_procedure_factory: RepartitionProcedureFactoryRef,
114 #[cfg(feature = "enterprise")]
115 trigger_ddl_manager: Option<TriggerDdlManagerRef>,
116 #[cfg(feature = "enterprise")]
117 create_flow_handler: Option<CreateFlowHandlerRef>,
118 #[cfg(feature = "enterprise")]
119 drop_flow_handler: Option<DropFlowHandlerRef>,
120}
121
122#[cfg(feature = "enterprise")]
125#[async_trait::async_trait]
126pub trait TriggerDdlManager: Send + Sync {
127 async fn create_trigger(
128 &self,
129 create_trigger_task: CreateTriggerTask,
130 procedure_manager: ProcedureManagerRef,
131 ddl_context: DdlContext,
132 query_context: QueryContext,
133 procedure_context: ProcedureContext,
134 ) -> Result<SubmitDdlTaskResponse>;
135
136 async fn drop_trigger(
137 &self,
138 drop_trigger_task: DropTriggerTask,
139 procedure_manager: ProcedureManagerRef,
140 ddl_context: DdlContext,
141 query_context: QueryContext,
142 procedure_context: ProcedureContext,
143 ) -> Result<SubmitDdlTaskResponse>;
144
145 fn as_any(&self) -> &dyn std::any::Any;
146}
147
148#[cfg(feature = "enterprise")]
149pub type TriggerDdlManagerRef = Arc<dyn TriggerDdlManager>;
150
151#[cfg(feature = "enterprise")]
153#[async_trait::async_trait]
154pub trait CreateFlowHandler: Send + Sync {
155 async fn create_flow(
156 &self,
157 create_flow_task: CreateFlowTask,
158 procedure_manager: ProcedureManagerRef,
159 ddl_context: DdlContext,
160 query_context: QueryContext,
161 procedure_context: ProcedureContext,
162 ) -> Result<SubmitDdlTaskResponse>;
163}
164
165#[cfg(feature = "enterprise")]
166pub type CreateFlowHandlerRef = Arc<dyn CreateFlowHandler>;
167
168#[cfg(feature = "enterprise")]
170#[async_trait::async_trait]
171pub trait DropFlowHandler: Send + Sync {
172 async fn drop_flow(
173 &self,
174 drop_flow_task: DropFlowTask,
175 procedure_manager: ProcedureManagerRef,
176 ddl_context: DdlContext,
177 procedure_context: ProcedureContext,
178 ) -> Result<SubmitDdlTaskResponse>;
179}
180
181#[cfg(feature = "enterprise")]
182pub type DropFlowHandlerRef = Arc<dyn DropFlowHandler>;
183
184macro_rules! procedure_loader_entry {
185 ($procedure:ident) => {
186 (
187 $procedure::TYPE_NAME,
188 &|context: DdlContext| -> common_procedure::BoxedProcedureLoader {
189 Box::new(move |json: &str| {
190 let context = context.clone();
191 $procedure::from_json(json, context).map(|p| Box::new(p) as _)
192 })
193 },
194 )
195 };
196}
197
198macro_rules! procedure_loader {
199 ($($procedure:ident),*) => {
200 vec![
201 $(procedure_loader_entry!($procedure)),*
202 ]
203 };
204}
205
206pub type RepartitionProcedureFactoryRef = Arc<dyn RepartitionProcedureFactory>;
207
208pub enum RepartitionSource {
209 Partitioned {
210 exprs: Vec<String>,
211 target_partition_columns: Option<Vec<String>>,
217 },
218 Unpartitioned {
219 partition_columns: Vec<String>,
220 },
221}
222
223#[async_trait::async_trait]
224pub trait RepartitionProcedureFactory: Send + Sync {
225 #[allow(clippy::too_many_arguments)]
226 fn create(
227 &self,
228 ddl_ctx: &DdlContext,
229 table_name: TableName,
230 table_id: TableId,
231 source: RepartitionSource,
232 to_exprs: Vec<String>,
233 timeout: Option<Duration>,
234 ) -> std::result::Result<BoxedProcedure, BoxedError>;
235
236 fn register_loaders(
237 &self,
238 ddl_ctx: &DdlContext,
239 procedure_manager: &ProcedureManagerRef,
240 ) -> std::result::Result<(), BoxedError>;
241
242 async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError>;
244}
245
246#[derive(Debug, Clone, Copy)]
251pub struct DdlOptions {
252 pub timeout: Duration,
256 pub wait: bool,
263}
264
265impl DdlManager {
266 pub fn new(
268 ddl_context: DdlContext,
269 procedure_manager: ProcedureManagerRef,
270 repartition_procedure_factory: RepartitionProcedureFactoryRef,
271 ) -> Self {
272 Self {
273 ddl_context,
274 procedure_manager,
275 repartition_procedure_factory,
276 #[cfg(feature = "enterprise")]
277 trigger_ddl_manager: None,
278 #[cfg(feature = "enterprise")]
279 create_flow_handler: None,
280 #[cfg(feature = "enterprise")]
281 drop_flow_handler: None,
282 }
283 }
284
285 #[cfg(feature = "enterprise")]
286 pub fn with_drop_flow_handler(mut self, drop_flow_handler: DropFlowHandlerRef) -> Self {
287 self.drop_flow_handler = Some(drop_flow_handler);
288 self
289 }
290
291 #[cfg(feature = "enterprise")]
292 pub fn with_trigger_ddl_manager(mut self, trigger_ddl_manager: TriggerDdlManagerRef) -> Self {
293 self.trigger_ddl_manager = Some(trigger_ddl_manager);
294 self
295 }
296
297 #[cfg(feature = "enterprise")]
298 pub fn with_create_flow_handler(mut self, create_flow_handler: CreateFlowHandlerRef) -> Self {
299 self.create_flow_handler = Some(create_flow_handler);
300 self
301 }
302
303 pub fn with_create_database_metadata_committer(
304 mut self,
305 committer: CreateDatabaseMetadataCommitterRef,
306 ) -> Self {
307 self.ddl_context.create_database_metadata_committer = Some(committer);
308 self
309 }
310
311 pub fn table_metadata_manager(&self) -> &TableMetadataManagerRef {
313 &self.ddl_context.table_metadata_manager
314 }
315
316 pub fn create_context(&self) -> DdlContext {
318 self.ddl_context.clone()
319 }
320
321 pub fn register_loaders(&self) -> Result<()> {
323 let loaders: Vec<(&str, &BoxedProcedureLoaderFactory)> = procedure_loader!(
324 CreateTableProcedure,
325 CreateLogicalTablesProcedure,
326 CreateViewProcedure,
327 CreateFlowProcedure,
328 CreateDatabaseProcedure,
329 AlterTableProcedure,
330 AlterLogicalTablesProcedure,
331 AlterDatabaseProcedure,
332 DropTableProcedure,
333 DropFlowProcedure,
334 TruncateTableProcedure,
335 DropDatabaseProcedure,
336 DropViewProcedure,
337 CommentOnProcedure
338 );
339 #[cfg(feature = "enterprise")]
340 let loaders = {
341 let soft_drop_loaders: Vec<(&str, &BoxedProcedureLoaderFactory)> =
342 procedure_loader!(UndropTableProcedure, PurgeDroppedTableProcedure);
343 loaders
344 .into_iter()
345 .chain(soft_drop_loaders)
346 .collect::<Vec<_>>()
347 };
348
349 for (type_name, loader_factory) in loaders {
350 let context = self.create_context();
351 self.procedure_manager
352 .register_loader(type_name, loader_factory(context))
353 .context(RegisterProcedureLoaderSnafu { type_name })?;
354 }
355
356 #[cfg(feature = "enterprise")]
357 {
358 let type_name = PurgeDroppedTableProcedure::EXPIRED_TYPE_NAME;
359 let context = self.create_context();
360 self.procedure_manager
361 .register_loader(
362 type_name,
363 Box::new(move |json: &str| {
364 PurgeDroppedTableProcedure::from_json(json, context.clone())
365 .map(|procedure| Box::new(procedure) as _)
366 }),
367 )
368 .context(RegisterProcedureLoaderSnafu { type_name })?;
369 }
370
371 self.repartition_procedure_factory
372 .register_loaders(&self.ddl_context, &self.procedure_manager)
373 .context(RegisterRepartitionProcedureLoaderSnafu)?;
374
375 Ok(())
376 }
377
378 async fn submit_repartition_task(
397 &self,
398 table_id: TableId,
399 table_name: TableName,
400 repartition: Repartition,
401 wait: bool,
402 timeout: Duration,
403 procedure_context: ProcedureContext,
404 ) -> Result<(ProcedureId, Option<Output>)> {
405 let context = self.create_context();
406
407 let into_partition_exprs = repartition.into_partition_exprs;
408 let source = repartition.source;
409
410 let source = match source {
411 Some(PbRepartitionSource::PartitionExprs(source)) => RepartitionSource::Partitioned {
412 exprs: source.exprs,
413 target_partition_columns: source
414 .target_partition_columns
415 .map(|columns| columns.columns),
416 },
417 Some(PbRepartitionSource::Unpartitioned(source)) => RepartitionSource::Unpartitioned {
418 partition_columns: source.partition_columns,
419 },
420 None => {
421 #[allow(deprecated)]
423 RepartitionSource::Partitioned {
424 exprs: repartition.from_partition_exprs,
425 target_partition_columns: None,
426 }
427 }
428 };
429
430 let procedure = self
431 .repartition_procedure_factory
432 .create(
433 &context,
434 table_name,
435 table_id,
436 source,
437 into_partition_exprs,
438 Some(timeout),
439 )
440 .context(CreateRepartitionProcedureSnafu)?;
441 self.repartition_procedure_factory
442 .ensure_gc_requirement()
443 .await
444 .context(PersistRepartitionGcRequirementSnafu)?;
445 let procedure_with_id =
446 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
447 if wait {
448 self.execute_procedure_and_wait(procedure_with_id).await
449 } else {
450 self.submit_procedure(procedure_with_id)
451 .await
452 .map(|p| (p, None))
453 }
454 }
455
456 #[tracing::instrument(skip_all)]
458 pub async fn submit_alter_table_task(
459 &self,
460 table_id: TableId,
461 alter_table_task: AlterTableTask,
462 procedure_context: ProcedureContext,
463 ddl_options: DdlOptions,
464 ) -> Result<(ProcedureId, Option<Output>)> {
465 let mut alter_table_task = alter_table_task;
467 if let Some(Kind::Repartition(_)) = alter_table_task.alter_table.kind.as_ref()
468 && let Kind::Repartition(repartition) =
469 alter_table_task.alter_table.kind.take().unwrap()
470 {
471 let table_name = TableName::new(
472 alter_table_task.alter_table.catalog_name,
473 alter_table_task.alter_table.schema_name,
474 alter_table_task.alter_table.table_name,
475 );
476 return self
477 .submit_repartition_task(
478 table_id,
479 table_name,
480 repartition,
481 ddl_options.wait,
482 ddl_options.timeout,
483 procedure_context,
484 )
485 .await;
486 }
487
488 let lock_regions = alter_table_task
489 .alter_table
490 .kind
491 .as_ref()
492 .is_some_and(only_sets_skip_wal);
493
494 let mut route_change_retries = 0;
495 loop {
496 let region_ids_to_lock = if lock_regions {
499 let (_, route) = self
500 .table_metadata_manager()
501 .table_route_manager()
502 .get_physical_table_route(table_id)
503 .await?;
504 route
505 .region_routes
506 .iter()
507 .map(|route| route.region.id)
508 .collect::<Vec<RegionId>>()
509 } else {
510 vec![]
511 };
512
513 let context = self.create_context();
514 let procedure = AlterTableProcedure::new_with_region_locks(
515 table_id,
516 alter_table_task.clone(),
517 region_ids_to_lock,
518 context,
519 )?;
520
521 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
522 .with_context(procedure_context.clone());
523 let result = self.execute_procedure_and_wait(procedure_with_id).await?;
524 if result
525 .1
526 .as_ref()
527 .is_some_and(|output| output.is::<RegionRouteChanged>())
528 {
529 if route_change_retries == MAX_REGION_ROUTE_CHANGE_RETRIES {
530 let source = error::UnexpectedSnafu {
531 err_msg: format!(
532 "Region route kept changing while altering table {table_id}"
533 ),
534 }
535 .build();
536 return Err(error::Error::retry_later(source));
537 }
538 route_change_retries += 1;
539 continue;
540 }
541 return Ok(result);
542 }
543 }
544
545 #[tracing::instrument(skip_all)]
547 pub async fn submit_create_table_task(
548 &self,
549 create_table_task: CreateTableTask,
550 query_context: QueryContext,
551 procedure_context: ProcedureContext,
552 ) -> Result<(ProcedureId, Option<Output>)> {
553 let context = self.create_context();
554
555 let procedure = CreateTableProcedure::new_with_query_context(
556 create_table_task,
557 query_context,
558 context,
559 )?;
560
561 let procedure_with_id =
562 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
563
564 self.execute_procedure_and_wait(procedure_with_id).await
565 }
566
567 #[tracing::instrument(skip_all)]
569 pub async fn submit_create_view_task(
570 &self,
571 create_view_task: CreateViewTask,
572 procedure_context: ProcedureContext,
573 ) -> Result<(ProcedureId, Option<Output>)> {
574 let context = self.create_context();
575
576 let procedure = CreateViewProcedure::new(create_view_task, context);
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_create_logical_table_tasks(
587 &self,
588 create_table_tasks: Vec<CreateTableTask>,
589 physical_table_id: TableId,
590 procedure_context: ProcedureContext,
591 ) -> Result<(ProcedureId, Option<Output>)> {
592 let context = self.create_context();
593
594 let procedure =
595 CreateLogicalTablesProcedure::new(create_table_tasks, physical_table_id, context);
596
597 let procedure_with_id =
598 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
599
600 self.execute_procedure_and_wait(procedure_with_id).await
601 }
602
603 #[tracing::instrument(skip_all)]
605 pub async fn submit_alter_logical_table_tasks(
606 &self,
607 alter_table_tasks: Vec<AlterTableTask>,
608 physical_table_id: TableId,
609 procedure_context: ProcedureContext,
610 ) -> Result<(ProcedureId, Option<Output>)> {
611 let context = self.create_context();
612
613 let logical_table_ids = {
616 let table_refs = alter_table_tasks
617 .iter()
618 .map(|task| task.table_ref())
619 .collect::<Vec<_>>();
620 utils::table_id::get_all_table_ids_by_names(
621 self.table_metadata_manager().table_name_manager(),
622 &table_refs,
623 )
624 .await?
625 };
626
627 let procedure = AlterLogicalTablesProcedure::new(
628 alter_table_tasks,
629 physical_table_id,
630 logical_table_ids,
631 context,
632 );
633
634 let procedure_with_id =
635 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
636
637 self.execute_procedure_and_wait(procedure_with_id).await
638 }
639
640 #[tracing::instrument(skip_all)]
642 pub async fn submit_drop_table_task(
643 &self,
644 drop_table_task: DropTableTask,
645 procedure_context: ProcedureContext,
646 ) -> Result<(ProcedureId, Option<Output>)> {
647 let context = self.create_context();
648
649 let procedure = DropTableProcedure::new(drop_table_task, context);
650
651 let procedure_with_id =
652 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
653
654 self.execute_procedure_and_wait(procedure_with_id).await
655 }
656
657 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
659 #[tracing::instrument(skip_all)]
660 pub async fn submit_undrop_table_task(
661 &self,
662 undrop_table_task: UndropTableTask,
663 procedure_context: ProcedureContext,
664 ) -> Result<(ProcedureId, Option<Output>)> {
665 #[cfg(not(feature = "enterprise"))]
666 {
667 use crate::error::UnsupportedSnafu;
668
669 return UnsupportedSnafu {
670 operation: "undrop table is only available in GreptimeDB Enterprise Edition",
671 }
672 .fail();
673 }
674 #[cfg(feature = "enterprise")]
675 {
676 let context = self.create_context();
677 let original_table_name = context
678 .table_metadata_manager
679 .get_dropped_table_by_id(undrop_table_task.table_id)
680 .await?
681 .with_context(|| TableNotFoundSnafu {
682 table_name: undrop_table_task.table_id.to_string(),
683 })?
684 .table_name;
685 let procedure = UndropTableProcedure::new_with_original_table_name(
686 undrop_table_task,
687 context,
688 Some(original_table_name),
689 );
690 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
691 .with_context(procedure_context);
692
693 self.execute_procedure_and_wait(procedure_with_id).await
694 }
695 }
696
697 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
699 #[tracing::instrument(skip_all)]
700 pub async fn submit_purge_dropped_table_task(
701 &self,
702 purge_dropped_table_task: PurgeDroppedTableTask,
703 procedure_context: ProcedureContext,
704 ) -> Result<(ProcedureId, Option<Output>)> {
705 #[cfg(not(feature = "enterprise"))]
706 {
707 use crate::error::UnsupportedSnafu;
708
709 return UnsupportedSnafu {
710 operation: "purge dropped table is only available in GreptimeDB Enterprise Edition",
711 }
712 .fail();
713 }
714 #[cfg(feature = "enterprise")]
715 {
716 let context = self.create_context();
717 let procedure = PurgeDroppedTableProcedure::new(purge_dropped_table_task, context);
718 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
719 .with_context(procedure_context);
720
721 self.execute_procedure_and_wait(procedure_with_id).await
722 }
723 }
724
725 #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
727 #[tracing::instrument(skip_all)]
728 pub async fn submit_expired_purge_dropped_table_task(
729 &self,
730 purge_dropped_table_task: PurgeDroppedTableTask,
731 ) -> Result<(ProcedureId, Option<Output>)> {
732 #[cfg(not(feature = "enterprise"))]
733 {
734 use crate::error::UnsupportedSnafu;
735
736 return UnsupportedSnafu {
737 operation:
738 "purge expired dropped table is only available in GreptimeDB Enterprise Edition",
739 }
740 .fail();
741 }
742 #[cfg(feature = "enterprise")]
743 {
744 let context = self.create_context();
745 let procedure_context = ProcedureContext::from_event_context(
746 PersistentEventContext::new(TriggerReason::ScheduledGc),
747 );
748 let procedure =
749 PurgeDroppedTableProcedure::new_if_expired(purge_dropped_table_task, context);
750 let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
751 .with_context(procedure_context);
752
753 self.execute_procedure_and_wait(procedure_with_id).await
754 }
755 }
756
757 #[tracing::instrument(skip_all)]
759 pub async fn submit_create_database(
760 &self,
761 CreateDatabaseTask {
762 catalog,
763 schema,
764 create_if_not_exists,
765 options,
766 creator,
767 }: CreateDatabaseTask,
768 procedure_context: ProcedureContext,
769 ) -> Result<(ProcedureId, Option<Output>)> {
770 let context = self.create_context();
771 let procedure = CreateDatabaseProcedure::new(
772 catalog,
773 schema,
774 create_if_not_exists,
775 options,
776 creator,
777 context,
778 );
779 let procedure_with_id =
780 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
781
782 self.execute_procedure_and_wait(procedure_with_id).await
783 }
784
785 #[tracing::instrument(skip_all)]
787 pub async fn submit_drop_database(
788 &self,
789 DropDatabaseTask {
790 catalog,
791 schema,
792 drop_if_exists,
793 }: DropDatabaseTask,
794 procedure_context: ProcedureContext,
795 ) -> Result<(ProcedureId, Option<Output>)> {
796 let context = self.create_context();
797 let procedure = DropDatabaseProcedure::new(catalog, schema, drop_if_exists, context);
798 let procedure_with_id =
799 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
800
801 self.execute_procedure_and_wait(procedure_with_id).await
802 }
803
804 pub async fn submit_alter_database(
805 &self,
806 alter_database_task: AlterDatabaseTask,
807 procedure_context: ProcedureContext,
808 ) -> Result<(ProcedureId, Option<Output>)> {
809 let context = self.create_context();
810 let procedure = AlterDatabaseProcedure::new(alter_database_task, context)?;
811 let procedure_with_id =
812 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
813
814 self.execute_procedure_and_wait(procedure_with_id).await
815 }
816
817 #[tracing::instrument(skip_all)]
819 pub async fn submit_create_flow_task(
820 &self,
821 create_flow: CreateFlowTask,
822 query_context: QueryContext,
823 procedure_context: ProcedureContext,
824 ) -> Result<(ProcedureId, Option<Output>)> {
825 let context = self.create_context();
826 let procedure = CreateFlowProcedure::new(create_flow, query_context, context);
827 let procedure_with_id =
828 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
829
830 self.execute_procedure_and_wait(procedure_with_id).await
831 }
832
833 #[tracing::instrument(skip_all)]
835 pub async fn submit_drop_flow_task(
836 &self,
837 drop_flow: DropFlowTask,
838 procedure_context: ProcedureContext,
839 ) -> Result<(ProcedureId, Option<Output>)> {
840 let context = self.create_context();
841 let procedure = DropFlowProcedure::new(drop_flow, context);
842 let procedure_with_id =
843 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
844
845 self.execute_procedure_and_wait(procedure_with_id).await
846 }
847
848 #[tracing::instrument(skip_all)]
850 pub async fn submit_drop_view_task(
851 &self,
852 drop_view: DropViewTask,
853 procedure_context: ProcedureContext,
854 ) -> Result<(ProcedureId, Option<Output>)> {
855 let context = self.create_context();
856 let procedure = DropViewProcedure::new(drop_view, context);
857 let procedure_with_id =
858 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
859
860 self.execute_procedure_and_wait(procedure_with_id).await
861 }
862
863 #[tracing::instrument(skip_all)]
865 pub async fn submit_truncate_table_task(
866 &self,
867 truncate_table_task: TruncateTableTask,
868 table_info_value: DeserializedValueWithBytes<TableInfoValue>,
869 procedure_context: ProcedureContext,
870 ) -> Result<(ProcedureId, Option<Output>)> {
871 let context = self.create_context();
872 let procedure = TruncateTableProcedure::new(truncate_table_task, table_info_value, context);
873
874 let procedure_with_id =
875 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
876
877 self.execute_procedure_and_wait(procedure_with_id).await
878 }
879
880 #[tracing::instrument(skip_all)]
882 pub async fn submit_comment_on_task(
883 &self,
884 mut comment_on_task: CommentOnTask,
885 procedure_context: ProcedureContext,
886 ) -> Result<(ProcedureId, Option<Output>)> {
887 let context = self.create_context();
888 comment_on_task
889 .enrich_object_id(
890 context.table_metadata_manager.table_name_manager(),
891 context.flow_metadata_manager.flow_name_manager(),
892 )
893 .await?;
894 let procedure = CommentOnProcedure::new(comment_on_task, context);
895 let procedure_with_id =
896 ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
897
898 self.execute_procedure_and_wait(procedure_with_id).await
899 }
900
901 async fn execute_procedure_and_wait(
903 &self,
904 procedure_with_id: ProcedureWithId,
905 ) -> Result<(ProcedureId, Option<Output>)> {
906 let procedure_id = procedure_with_id.id;
907
908 let mut watcher = self
909 .procedure_manager
910 .submit(procedure_with_id)
911 .await
912 .context(SubmitProcedureSnafu)?;
913
914 let output = watcher::wait(&mut watcher)
915 .await
916 .context(WaitProcedureSnafu)?;
917
918 Ok((procedure_id, output))
919 }
920
921 async fn submit_procedure(&self, procedure_with_id: ProcedureWithId) -> Result<ProcedureId> {
923 let procedure_id = procedure_with_id.id;
924 let _ = self
925 .procedure_manager
926 .submit(procedure_with_id)
927 .await
928 .context(SubmitProcedureSnafu)?;
929
930 Ok(procedure_id)
931 }
932
933 pub async fn submit_ddl_task(
934 &self,
935 context: ExecutorContext,
936 request: SubmitDdlTaskRequest,
937 ) -> Result<SubmitDdlTaskResponse> {
938 let ExecutorContext {
939 tracing_context,
940 query_context,
941 actor,
942 event_input,
943 } = context;
944 let query_context = query_context.context(UnsupportedSnafu {
945 operation: "submit_ddl_task without query context",
946 })?;
947 let procedure_context = ProcedureContext {
948 actor,
949 event_context: event_input
950 .map(|input| PersistentEventContext::from((input, query_context.protocol()))),
951 };
952 let span = tracing_context
953 .as_ref()
954 .map(TracingContext::from_w3c)
955 .unwrap_or_else(TracingContext::from_current_span)
956 .attach(tracing::info_span!("DdlManager::submit_ddl_task"));
957 let SubmitDdlTaskRequest {
958 wait,
959 timeout,
960 task,
961 } = request;
962 let ddl_options = DdlOptions { wait, timeout };
963 async move {
964 debug!("Submitting Ddl task: {:?}", task);
965 match task {
966 CreateTable(create_table_task) => {
967 handle_create_table_task(
968 self,
969 create_table_task,
970 query_context,
971 procedure_context,
972 )
973 .await
974 }
975 DropTable(drop_table_task) => {
976 handle_drop_table_task(self, drop_table_task, procedure_context).await
977 }
978 UndropTable(undrop_table_task) => {
979 handle_undrop_table_task(self, undrop_table_task, procedure_context).await
980 }
981 PurgeDroppedTable(purge_dropped_table_task) => {
982 handle_purge_dropped_table_task(
983 self,
984 purge_dropped_table_task,
985 procedure_context,
986 )
987 .await
988 }
989 AlterTable(alter_table_task) => {
990 handle_alter_table_task(self, alter_table_task, ddl_options, procedure_context)
991 .await
992 }
993 TruncateTable(truncate_table_task) => {
994 handle_truncate_table_task(self, truncate_table_task, procedure_context).await
995 }
996 CreateLogicalTables(create_table_tasks) => {
997 handle_create_logical_table_tasks(self, create_table_tasks, procedure_context)
998 .await
999 }
1000 AlterLogicalTables(alter_table_tasks) => {
1001 handle_alter_logical_table_tasks(self, alter_table_tasks, procedure_context)
1002 .await
1003 }
1004 DropLogicalTables(_) => todo!(),
1005 CreateDatabase(create_database_task) => {
1006 handle_create_database_task(self, create_database_task, procedure_context).await
1007 }
1008 DropDatabase(drop_database_task) => {
1009 handle_drop_database_task(self, drop_database_task, procedure_context).await
1010 }
1011 AlterDatabase(alter_database_task) => {
1012 handle_alter_database_task(self, alter_database_task, procedure_context).await
1013 }
1014 CreateFlow(create_flow_task) => {
1015 handle_create_flow_task(
1016 self,
1017 create_flow_task,
1018 query_context,
1019 procedure_context,
1020 )
1021 .await
1022 }
1023 DropFlow(drop_flow_task) => {
1024 #[cfg(feature = "enterprise")]
1025 if let Some(handler) = self.drop_flow_handler.as_ref() {
1026 return handler
1027 .drop_flow(
1028 drop_flow_task,
1029 self.procedure_manager.clone(),
1030 self.ddl_context.clone(),
1031 procedure_context,
1032 )
1033 .await;
1034 }
1035 handle_drop_flow_task(self, drop_flow_task, procedure_context).await
1036 }
1037 CreateView(create_view_task) => {
1038 handle_create_view_task(self, create_view_task, procedure_context).await
1039 }
1040 DropView(drop_view_task) => {
1041 handle_drop_view_task(self, drop_view_task, procedure_context).await
1042 }
1043 CommentOn(comment_on_task) => {
1044 handle_comment_on_task(self, comment_on_task, procedure_context).await
1045 }
1046 #[cfg(feature = "enterprise")]
1047 CreateTrigger(create_trigger_task) => {
1048 handle_create_trigger_task(
1049 self,
1050 create_trigger_task,
1051 query_context,
1052 procedure_context,
1053 )
1054 .await
1055 }
1056 #[cfg(feature = "enterprise")]
1057 DropTrigger(drop_trigger_task) => {
1058 handle_drop_trigger_task(
1059 self,
1060 drop_trigger_task,
1061 query_context,
1062 procedure_context,
1063 )
1064 .await
1065 }
1066 }
1067 }
1068 .trace(span)
1069 .await
1070 }
1071}
1072
1073async fn handle_truncate_table_task(
1074 ddl_manager: &DdlManager,
1075 truncate_table_task: TruncateTableTask,
1076 procedure_context: ProcedureContext,
1077) -> Result<SubmitDdlTaskResponse> {
1078 let table_id = truncate_table_task.table_id;
1079 let table_metadata_manager = &ddl_manager.table_metadata_manager();
1080 let table_ref = truncate_table_task.table_ref();
1081
1082 let table_info_value = table_metadata_manager
1083 .table_info_manager()
1084 .get(table_id)
1085 .await?
1086 .with_context(|| TableInfoNotFoundSnafu {
1087 table: table_ref.to_string(),
1088 })?;
1089 let physical_table_id = table_metadata_manager
1090 .table_route_manager()
1091 .get_physical_table_id(table_id)
1092 .await?;
1093 ensure!(
1094 physical_table_id == table_id,
1095 error::UnexpectedSnafu {
1096 err_msg: "Truncate table is only supported for physical tables."
1097 }
1098 );
1099
1100 let (id, _) = ddl_manager
1101 .submit_truncate_table_task(truncate_table_task, table_info_value, procedure_context)
1102 .await?;
1103
1104 info!("Table: {table_id} is truncated via procedure_id {id:?}");
1105
1106 Ok(SubmitDdlTaskResponse {
1107 key: id.to_string().into(),
1108 ..Default::default()
1109 })
1110}
1111
1112async fn handle_alter_table_task(
1113 ddl_manager: &DdlManager,
1114 alter_table_task: AlterTableTask,
1115 ddl_options: DdlOptions,
1116 procedure_context: ProcedureContext,
1117) -> Result<SubmitDdlTaskResponse> {
1118 let table_ref = alter_table_task.table_ref();
1119
1120 let table_id = ddl_manager
1121 .table_metadata_manager()
1122 .table_name_manager()
1123 .get(TableNameKey::new(
1124 table_ref.catalog,
1125 table_ref.schema,
1126 table_ref.table,
1127 ))
1128 .await?
1129 .with_context(|| TableNotFoundSnafu {
1130 table_name: table_ref.to_string(),
1131 })?
1132 .table_id();
1133
1134 let table_route_value = ddl_manager
1135 .table_metadata_manager()
1136 .table_route_manager()
1137 .table_route_storage()
1138 .get(table_id)
1139 .await?
1140 .context(TableRouteNotFoundSnafu { table_id })?;
1141 let annotation_family = match alter_table_task.alter_table.kind.as_ref() {
1145 Some(kind) => common_grpc_expr::annotation_alter_family(kind)
1146 .context(ConvertAlterTableRequestSnafu)?,
1147 None => None,
1148 };
1149 ensure!(
1150 table_route_value.is_physical()
1151 || annotation_family.is_some_and(|family| family.allows_logical_tables()),
1152 UnexpectedLogicalRouteTableSnafu {
1153 err_msg: format!("{:?} is a non-physical TableRouteValue.", table_ref),
1154 }
1155 );
1156
1157 let (id, _) = ddl_manager
1158 .submit_alter_table_task(table_id, alter_table_task, procedure_context, ddl_options)
1159 .await?;
1160
1161 info!("Table: {table_id} is altered via procedure_id {id:?}");
1162
1163 Ok(SubmitDdlTaskResponse {
1164 key: id.to_string().into(),
1165 ..Default::default()
1166 })
1167}
1168
1169async fn handle_drop_table_task(
1170 ddl_manager: &DdlManager,
1171 drop_table_task: DropTableTask,
1172 procedure_context: ProcedureContext,
1173) -> Result<SubmitDdlTaskResponse> {
1174 let table_id = drop_table_task.table_id;
1175 let (id, _) = ddl_manager
1176 .submit_drop_table_task(drop_table_task, procedure_context)
1177 .await?;
1178
1179 info!("Table: {table_id} is dropped via procedure_id {id:?}");
1180
1181 Ok(SubmitDdlTaskResponse {
1182 key: id.to_string().into(),
1183 ..Default::default()
1184 })
1185}
1186
1187async fn handle_undrop_table_task(
1188 ddl_manager: &DdlManager,
1189 undrop_table_task: UndropTableTask,
1190 procedure_context: ProcedureContext,
1191) -> Result<SubmitDdlTaskResponse> {
1192 let table_id = undrop_table_task.table_id;
1193 let (id, _) = ddl_manager
1194 .submit_undrop_table_task(undrop_table_task, procedure_context)
1195 .await?;
1196
1197 info!("Table: {table_id} is undropped via procedure_id {id:?}");
1198
1199 Ok(SubmitDdlTaskResponse {
1200 key: id.to_string().into(),
1201 ..Default::default()
1202 })
1203}
1204
1205async fn handle_purge_dropped_table_task(
1206 ddl_manager: &DdlManager,
1207 purge_dropped_table_task: PurgeDroppedTableTask,
1208 procedure_context: ProcedureContext,
1209) -> Result<SubmitDdlTaskResponse> {
1210 let (id, _) = ddl_manager
1211 .submit_purge_dropped_table_task(purge_dropped_table_task, procedure_context)
1212 .await?;
1213
1214 info!("Dropped table is purged via procedure_id {id:?}");
1215
1216 Ok(SubmitDdlTaskResponse {
1217 key: id.to_string().into(),
1218 ..Default::default()
1219 })
1220}
1221
1222async fn handle_create_table_task(
1223 ddl_manager: &DdlManager,
1224 create_table_task: CreateTableTask,
1225 query_context: QueryContext,
1226 procedure_context: ProcedureContext,
1227) -> Result<SubmitDdlTaskResponse> {
1228 let (id, output) = ddl_manager
1229 .submit_create_table_task(create_table_task, query_context, procedure_context)
1230 .await?;
1231
1232 let procedure_id = id.to_string();
1233 let output = output.context(ProcedureOutputSnafu {
1234 procedure_id: &procedure_id,
1235 err_msg: "empty output",
1236 })?;
1237 let table_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1238 procedure_id: &procedure_id,
1239 err_msg: "downcast to `u32`",
1240 })?);
1241 info!("Table: {table_id} is created via procedure_id {id:?}");
1242
1243 Ok(SubmitDdlTaskResponse {
1244 key: procedure_id.into(),
1245 table_ids: vec![table_id],
1246 })
1247}
1248
1249async fn handle_create_logical_table_tasks(
1250 ddl_manager: &DdlManager,
1251 create_table_tasks: Vec<CreateTableTask>,
1252 procedure_context: ProcedureContext,
1253) -> Result<SubmitDdlTaskResponse> {
1254 ensure!(
1255 !create_table_tasks.is_empty(),
1256 EmptyDdlTasksSnafu {
1257 name: "create logical tables"
1258 }
1259 );
1260 let physical_table_id = utils::check_and_get_physical_table_id(
1261 ddl_manager.table_metadata_manager(),
1262 &create_table_tasks,
1263 )
1264 .await?;
1265 let num_logical_tables = create_table_tasks.len();
1266
1267 let (id, output) = ddl_manager
1268 .submit_create_logical_table_tasks(create_table_tasks, physical_table_id, procedure_context)
1269 .await?;
1270
1271 info!(
1272 "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is created via procedure_id {id:?}"
1273 );
1274
1275 let procedure_id = id.to_string();
1276 let output = output.context(ProcedureOutputSnafu {
1277 procedure_id: &procedure_id,
1278 err_msg: "empty output",
1279 })?;
1280 let table_ids = output
1281 .downcast_ref::<Vec<TableId>>()
1282 .context(ProcedureOutputSnafu {
1283 procedure_id: &procedure_id,
1284 err_msg: "downcast to `Vec<TableId>`",
1285 })?
1286 .clone();
1287
1288 Ok(SubmitDdlTaskResponse {
1289 key: procedure_id.into(),
1290 table_ids,
1291 })
1292}
1293
1294async fn handle_create_database_task(
1295 ddl_manager: &DdlManager,
1296 create_database_task: CreateDatabaseTask,
1297 procedure_context: ProcedureContext,
1298) -> Result<SubmitDdlTaskResponse> {
1299 let catalog = create_database_task.catalog.clone();
1300 let schema = create_database_task.schema.clone();
1301 let (id, _) = ddl_manager
1302 .submit_create_database(create_database_task, procedure_context)
1303 .await?;
1304
1305 let procedure_id = id.to_string();
1306 info!(
1307 "Database {}.{} is created via procedure_id {id:?}",
1308 catalog, schema
1309 );
1310
1311 Ok(SubmitDdlTaskResponse {
1312 key: procedure_id.into(),
1313 ..Default::default()
1314 })
1315}
1316
1317async fn handle_drop_database_task(
1318 ddl_manager: &DdlManager,
1319 drop_database_task: DropDatabaseTask,
1320 procedure_context: ProcedureContext,
1321) -> Result<SubmitDdlTaskResponse> {
1322 let (id, _) = ddl_manager
1323 .submit_drop_database(drop_database_task.clone(), procedure_context)
1324 .await?;
1325
1326 let procedure_id = id.to_string();
1327 info!(
1328 "Database {}.{} is dropped via procedure_id {id:?}",
1329 drop_database_task.catalog, drop_database_task.schema
1330 );
1331
1332 Ok(SubmitDdlTaskResponse {
1333 key: procedure_id.into(),
1334 ..Default::default()
1335 })
1336}
1337
1338async fn handle_alter_database_task(
1339 ddl_manager: &DdlManager,
1340 alter_database_task: AlterDatabaseTask,
1341 procedure_context: ProcedureContext,
1342) -> Result<SubmitDdlTaskResponse> {
1343 let (id, _) = ddl_manager
1344 .submit_alter_database(alter_database_task.clone(), procedure_context)
1345 .await?;
1346
1347 let procedure_id = id.to_string();
1348 info!(
1349 "Database {}.{} is altered via procedure_id {id:?}",
1350 alter_database_task.catalog(),
1351 alter_database_task.schema()
1352 );
1353
1354 Ok(SubmitDdlTaskResponse {
1355 key: procedure_id.into(),
1356 ..Default::default()
1357 })
1358}
1359
1360async fn handle_drop_flow_task(
1361 ddl_manager: &DdlManager,
1362 drop_flow_task: DropFlowTask,
1363 procedure_context: ProcedureContext,
1364) -> Result<SubmitDdlTaskResponse> {
1365 let (id, _) = ddl_manager
1366 .submit_drop_flow_task(drop_flow_task.clone(), procedure_context)
1367 .await?;
1368
1369 let procedure_id = id.to_string();
1370 info!(
1371 "Flow {}.{}({}) is dropped via procedure_id {id:?}",
1372 drop_flow_task.catalog_name, drop_flow_task.flow_name, drop_flow_task.flow_id,
1373 );
1374
1375 Ok(SubmitDdlTaskResponse {
1376 key: procedure_id.into(),
1377 ..Default::default()
1378 })
1379}
1380
1381#[cfg(feature = "enterprise")]
1382async fn handle_drop_trigger_task(
1383 ddl_manager: &DdlManager,
1384 drop_trigger_task: DropTriggerTask,
1385 query_context: QueryContext,
1386 procedure_context: ProcedureContext,
1387) -> Result<SubmitDdlTaskResponse> {
1388 let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1389 use crate::error::UnsupportedSnafu;
1390
1391 return UnsupportedSnafu {
1392 operation: "drop trigger",
1393 }
1394 .fail();
1395 };
1396
1397 m.drop_trigger(
1398 drop_trigger_task,
1399 ddl_manager.procedure_manager.clone(),
1400 ddl_manager.ddl_context.clone(),
1401 query_context,
1402 procedure_context,
1403 )
1404 .await
1405}
1406
1407async fn handle_drop_view_task(
1408 ddl_manager: &DdlManager,
1409 drop_view_task: DropViewTask,
1410 procedure_context: ProcedureContext,
1411) -> Result<SubmitDdlTaskResponse> {
1412 let (id, _) = ddl_manager
1413 .submit_drop_view_task(drop_view_task.clone(), procedure_context)
1414 .await?;
1415
1416 let procedure_id = id.to_string();
1417 info!(
1418 "View {}({}) is dropped via procedure_id {id:?}",
1419 drop_view_task.table_ref(),
1420 drop_view_task.view_id,
1421 );
1422
1423 Ok(SubmitDdlTaskResponse {
1424 key: procedure_id.into(),
1425 ..Default::default()
1426 })
1427}
1428
1429async fn handle_create_flow_task(
1430 ddl_manager: &DdlManager,
1431 create_flow_task: CreateFlowTask,
1432 query_context: QueryContext,
1433 procedure_context: ProcedureContext,
1434) -> Result<SubmitDdlTaskResponse> {
1435 if create_flow_task
1436 .flow_options
1437 .get(FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY)
1438 .is_some_and(|value| value == FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_SEQUENCE_RANGE)
1439 {
1440 return error::UnexpectedSnafu {
1441 err_msg: format!(
1442 "reserved flow option value for {FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY} is internal"
1443 ),
1444 }
1445 .fail();
1446 }
1447
1448 #[cfg(feature = "enterprise")]
1449 if let Some(handler) = ddl_manager.create_flow_handler.as_ref() {
1450 return handler
1451 .create_flow(
1452 create_flow_task,
1453 ddl_manager.procedure_manager.clone(),
1454 ddl_manager.ddl_context.clone(),
1455 query_context,
1456 procedure_context,
1457 )
1458 .await;
1459 }
1460
1461 let (id, output) = ddl_manager
1462 .submit_create_flow_task(create_flow_task.clone(), query_context, procedure_context)
1463 .await?;
1464
1465 let procedure_id = id.to_string();
1466 let output = output.context(ProcedureOutputSnafu {
1467 procedure_id: &procedure_id,
1468 err_msg: "empty output",
1469 })?;
1470 let flow_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1471 procedure_id: &procedure_id,
1472 err_msg: "downcast to `u32`",
1473 })?);
1474 if !create_flow_task.or_replace {
1475 info!(
1476 "Flow {}.{}({flow_id}) is created via procedure_id {id:?}",
1477 create_flow_task.catalog_name, create_flow_task.flow_name,
1478 );
1479 } else {
1480 info!(
1481 "Flow {}.{}({flow_id}) is replaced via procedure_id {id:?}",
1482 create_flow_task.catalog_name, create_flow_task.flow_name,
1483 );
1484 }
1485
1486 Ok(SubmitDdlTaskResponse {
1487 key: procedure_id.into(),
1488 ..Default::default()
1489 })
1490}
1491
1492#[cfg(feature = "enterprise")]
1493async fn handle_create_trigger_task(
1494 ddl_manager: &DdlManager,
1495 create_trigger_task: CreateTriggerTask,
1496 query_context: QueryContext,
1497 procedure_context: ProcedureContext,
1498) -> Result<SubmitDdlTaskResponse> {
1499 let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1500 use crate::error::UnsupportedSnafu;
1501
1502 return UnsupportedSnafu {
1503 operation: "create trigger",
1504 }
1505 .fail();
1506 };
1507
1508 m.create_trigger(
1509 create_trigger_task,
1510 ddl_manager.procedure_manager.clone(),
1511 ddl_manager.ddl_context.clone(),
1512 query_context,
1513 procedure_context,
1514 )
1515 .await
1516}
1517
1518async fn handle_alter_logical_table_tasks(
1519 ddl_manager: &DdlManager,
1520 alter_table_tasks: Vec<AlterTableTask>,
1521 procedure_context: ProcedureContext,
1522) -> Result<SubmitDdlTaskResponse> {
1523 ensure!(
1524 !alter_table_tasks.is_empty(),
1525 EmptyDdlTasksSnafu {
1526 name: "alter logical tables"
1527 }
1528 );
1529
1530 let first_table = TableNameKey {
1532 catalog: &alter_table_tasks[0].alter_table.catalog_name,
1533 schema: &alter_table_tasks[0].alter_table.schema_name,
1534 table: &alter_table_tasks[0].alter_table.table_name,
1535 };
1536 let physical_table_id =
1537 utils::get_physical_table_id(ddl_manager.table_metadata_manager(), first_table).await?;
1538 let num_logical_tables = alter_table_tasks.len();
1539
1540 let (id, _) = ddl_manager
1541 .submit_alter_logical_table_tasks(alter_table_tasks, physical_table_id, procedure_context)
1542 .await?;
1543
1544 info!(
1545 "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is altered via procedure_id {id:?}"
1546 );
1547
1548 let procedure_id = id.to_string();
1549
1550 Ok(SubmitDdlTaskResponse {
1551 key: procedure_id.into(),
1552 ..Default::default()
1553 })
1554}
1555
1556async fn handle_create_view_task(
1558 ddl_manager: &DdlManager,
1559 create_view_task: CreateViewTask,
1560 procedure_context: ProcedureContext,
1561) -> Result<SubmitDdlTaskResponse> {
1562 let (id, output) = ddl_manager
1563 .submit_create_view_task(create_view_task, procedure_context)
1564 .await?;
1565
1566 let procedure_id = id.to_string();
1567 let output = output.context(ProcedureOutputSnafu {
1568 procedure_id: &procedure_id,
1569 err_msg: "empty output",
1570 })?;
1571 let view_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1572 procedure_id: &procedure_id,
1573 err_msg: "downcast to `u32`",
1574 })?);
1575 info!("View: {view_id} is created via procedure_id {id:?}");
1576
1577 Ok(SubmitDdlTaskResponse {
1578 key: procedure_id.into(),
1579 table_ids: vec![view_id],
1580 })
1581}
1582
1583async fn handle_comment_on_task(
1584 ddl_manager: &DdlManager,
1585 comment_on_task: CommentOnTask,
1586 procedure_context: ProcedureContext,
1587) -> Result<SubmitDdlTaskResponse> {
1588 let (id, _) = ddl_manager
1589 .submit_comment_on_task(comment_on_task.clone(), procedure_context)
1590 .await?;
1591
1592 let procedure_id = id.to_string();
1593 info!(
1594 "Comment on {}.{}.{} is updated via procedure_id {id:?}",
1595 comment_on_task.catalog_name, comment_on_task.schema_name, comment_on_task.object_name
1596 );
1597
1598 Ok(SubmitDdlTaskResponse {
1599 key: procedure_id.into(),
1600 ..Default::default()
1601 })
1602}
1603
1604#[cfg(test)]
1605mod tests {
1606 use std::sync::Arc;
1607 #[cfg(feature = "enterprise")]
1608 use std::sync::Mutex;
1609 use std::time::Duration;
1610
1611 #[cfg(feature = "enterprise")]
1612 use common_base::protocol::Channel;
1613 use common_error::ext::BoxedError;
1614 #[cfg(feature = "enterprise")]
1615 use common_error::ext::ErrorExt;
1616 #[cfg(feature = "enterprise")]
1617 use common_error::status_code::StatusCode;
1618 #[cfg(feature = "enterprise")]
1619 use common_event_recorder::{PersistentEventContext, ProcedureEventInput, TriggerReason};
1620 use common_procedure::local::LocalManager;
1621 use common_procedure::test_util::InMemoryPoisonStore;
1622 use common_procedure::{
1623 BoxedProcedure, ProcedureContext, ProcedureManager, ProcedureManagerRef,
1624 };
1625 use store_api::storage::TableId;
1626 use table::table_name::TableName;
1627
1628 use super::DdlManager;
1629 use crate::cache_invalidator::DummyCacheInvalidator;
1630 use crate::ddl::alter_table::AlterTableProcedure;
1631 use crate::ddl::create_database::{
1632 AtomicCreateOutcome, CreateDatabaseMetadataCommitter, CreateDatabaseProcedure,
1633 };
1634 #[cfg(feature = "enterprise")]
1635 use crate::ddl::create_flow::{
1636 CreateFlowProcedure, FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY,
1637 FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_SEQUENCE_RANGE,
1638 };
1639 use crate::ddl::create_table::CreateTableProcedure;
1640 use crate::ddl::drop_table::DropTableProcedure;
1641 use crate::ddl::flow_meta::FlowMetadataAllocator;
1642 use crate::ddl::table_meta::TableMetadataAllocator;
1643 use crate::ddl::truncate_table::TruncateTableProcedure;
1644 use crate::ddl::{DdlContext, NoopRegionFailureDetectorControl};
1645 use crate::ddl_manager::{RepartitionProcedureFactory, RepartitionSource};
1646 use crate::key::TableMetadataManager;
1647 use crate::key::flow::FlowMetadataManager;
1648 use crate::kv_backend::memory::MemoryKvBackend;
1649 use crate::node_manager::{DatanodeManager, DatanodeRef, FlownodeManager, FlownodeRef};
1650 use crate::peer::Peer;
1651 use crate::procedure_executor::ExecutorContext;
1652 use crate::region_keeper::MemoryRegionKeeper;
1653 use crate::region_registry::LeaderRegionRegistry;
1654 #[cfg(feature = "enterprise")]
1655 use crate::rpc::ddl::trigger::{CreateTriggerTask, DropTriggerTask};
1656 #[cfg(feature = "enterprise")]
1657 use crate::rpc::ddl::{
1658 CreateFlowTask, DdlTask, QueryContext, SubmitDdlTaskRequest, SubmitDdlTaskResponse,
1659 };
1660 use crate::rpc::ddl::{CreatorGrantIntent, UndropTableTask};
1661 #[cfg(not(feature = "enterprise"))]
1662 use crate::rpc::ddl::{DdlTask, PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest};
1663 use crate::sequence::SequenceBuilder;
1664 use crate::state_store::KvStateStore;
1665 use crate::test_util::{MockDatanodeManager, new_ddl_context};
1666 use crate::wal_provider::WalProvider;
1667
1668 pub struct DummyDatanodeManager;
1670
1671 #[async_trait::async_trait]
1672 impl DatanodeManager for DummyDatanodeManager {
1673 async fn datanode(&self, _datanode: &Peer) -> DatanodeRef {
1674 unimplemented!()
1675 }
1676 }
1677
1678 #[async_trait::async_trait]
1679 impl FlownodeManager for DummyDatanodeManager {
1680 async fn flownode(&self, _node: &Peer) -> FlownodeRef {
1681 unimplemented!()
1682 }
1683 }
1684
1685 struct DummyRepartitionProcedureFactory;
1686
1687 #[async_trait::async_trait]
1688 impl RepartitionProcedureFactory for DummyRepartitionProcedureFactory {
1689 fn create(
1690 &self,
1691 _ddl_ctx: &DdlContext,
1692 _table_name: TableName,
1693 _table_id: TableId,
1694 _source: RepartitionSource,
1695 _to_exprs: Vec<String>,
1696 _timeout: Option<Duration>,
1697 ) -> std::result::Result<BoxedProcedure, BoxedError> {
1698 unimplemented!()
1699 }
1700
1701 fn register_loaders(
1702 &self,
1703 _ddl_ctx: &DdlContext,
1704 _procedure_manager: &ProcedureManagerRef,
1705 ) -> std::result::Result<(), BoxedError> {
1706 Ok(())
1707 }
1708
1709 async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError> {
1710 Ok(())
1711 }
1712 }
1713
1714 struct LoaderCommitter;
1715
1716 #[async_trait::async_trait]
1717 impl CreateDatabaseMetadataCommitter for LoaderCommitter {
1718 async fn commit(
1719 &self,
1720 _catalog: &str,
1721 _schema: &str,
1722 _value: &crate::key::schema_name::SchemaNameValue,
1723 _creator: &CreatorGrantIntent,
1724 ) -> std::result::Result<AtomicCreateOutcome, BoxedError> {
1725 unreachable!()
1726 }
1727 }
1728
1729 #[cfg(feature = "enterprise")]
1730 #[derive(Default)]
1731 struct RecordingCreateFlowHandler {
1732 tasks: Mutex<Vec<CreateFlowTask>>,
1733 contexts: Mutex<Vec<(QueryContext, ProcedureContext)>>,
1734 response: Mutex<Option<SubmitDdlTaskResponse>>,
1735 fail: bool,
1736 }
1737
1738 #[cfg(feature = "enterprise")]
1739 #[async_trait::async_trait]
1740 impl super::CreateFlowHandler for RecordingCreateFlowHandler {
1741 async fn create_flow(
1742 &self,
1743 create_flow_task: CreateFlowTask,
1744 _procedure_manager: ProcedureManagerRef,
1745 _ddl_context: DdlContext,
1746 query_context: QueryContext,
1747 procedure_context: ProcedureContext,
1748 ) -> crate::error::Result<SubmitDdlTaskResponse> {
1749 self.tasks.lock().unwrap().push(create_flow_task);
1750 self.contexts
1751 .lock()
1752 .unwrap()
1753 .push((query_context, procedure_context));
1754 if self.fail {
1755 return crate::error::UnsupportedSnafu {
1756 operation: "test create flow handler",
1757 }
1758 .fail();
1759 }
1760 Ok(self.response.lock().unwrap().take().unwrap_or_default())
1761 }
1762 }
1763
1764 #[cfg(feature = "enterprise")]
1765 fn test_create_flow_task() -> CreateFlowTask {
1766 CreateFlowTask {
1767 catalog_name: "greptime".to_string(),
1768 flow_name: "test_flow".to_string(),
1769 source_table_names: vec![],
1770 sink_table_name: TableName::new("greptime", "public", "sink"),
1771 or_replace: false,
1772 create_if_not_exists: false,
1773 expire_after: None,
1774 eval_interval_secs: None,
1775 eval_offset_secs: None,
1776 comment: String::new(),
1777 sql: "select 1".to_string(),
1778 flow_options: Default::default(),
1779 eval_schedule: None,
1780 }
1781 }
1782
1783 #[cfg(feature = "enterprise")]
1784 #[derive(Default)]
1785 struct RecordingTriggerDdlManager {
1786 procedure_contexts: Mutex<Vec<ProcedureContext>>,
1787 }
1788
1789 #[cfg(feature = "enterprise")]
1790 #[async_trait::async_trait]
1791 impl super::TriggerDdlManager for RecordingTriggerDdlManager {
1792 async fn create_trigger(
1793 &self,
1794 _create_trigger_task: CreateTriggerTask,
1795 _procedure_manager: ProcedureManagerRef,
1796 _ddl_context: DdlContext,
1797 _query_context: QueryContext,
1798 procedure_context: ProcedureContext,
1799 ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1800 self.procedure_contexts
1801 .lock()
1802 .unwrap()
1803 .push(procedure_context);
1804 Ok(Default::default())
1805 }
1806
1807 async fn drop_trigger(
1808 &self,
1809 _drop_trigger_task: DropTriggerTask,
1810 _procedure_manager: ProcedureManagerRef,
1811 _ddl_context: DdlContext,
1812 _query_context: QueryContext,
1813 procedure_context: ProcedureContext,
1814 ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1815 self.procedure_contexts
1816 .lock()
1817 .unwrap()
1818 .push(procedure_context);
1819 Ok(Default::default())
1820 }
1821
1822 fn as_any(&self) -> &dyn std::any::Any {
1823 self
1824 }
1825 }
1826
1827 #[test]
1828 fn test_generic_loader_captures_configured_committer() {
1829 let mut context = new_ddl_context(Arc::new(MockDatanodeManager::new(())));
1830 let committer = Arc::new(LoaderCommitter);
1831 context.create_database_metadata_committer = Some(committer.clone());
1832 let loader = procedure_loader!(CreateDatabaseProcedure)[0].1(context);
1833 assert_eq!(Arc::strong_count(&committer), 2);
1834
1835 let loaded = loader(
1836 r#"{
1837 "state":"CreateMetadata",
1838 "catalog":"greptime",
1839 "schema":"metrics",
1840 "create_if_not_exists":false,
1841 "options":{},
1842 "creator":{"username":"alice","created_at_ns":42}
1843 }"#,
1844 )
1845 .unwrap();
1846
1847 assert_eq!(loaded.type_name(), "metasrv-procedure::CreateDatabase");
1848 assert_eq!(Arc::strong_count(&committer), 3);
1849 drop(loaded);
1850 assert_eq!(Arc::strong_count(&committer), 2);
1851 }
1852
1853 #[test]
1854 fn test_register_loaders() {
1855 let kv_backend = Arc::new(MemoryKvBackend::new());
1856 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1857 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1858 Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1859 Arc::new(WalProvider::default()),
1860 ));
1861 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1862 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1863 Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1864 ));
1865
1866 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1867 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1868 let procedure_manager = Arc::new(LocalManager::new(
1869 Default::default(),
1870 state_store,
1871 poison_manager,
1872 None,
1873 None,
1874 ));
1875
1876 let ddl_manager = DdlManager::new(
1877 DdlContext {
1878 node_manager: Arc::new(DummyDatanodeManager),
1879 cache_invalidator: Arc::new(DummyCacheInvalidator),
1880 table_metadata_manager,
1881 table_metadata_allocator,
1882 flow_metadata_manager,
1883 flow_metadata_allocator,
1884 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1885 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1886 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1887 soft_drop_enabled: false,
1888 soft_drop_retention: None,
1889 create_database_metadata_committer: None,
1890 },
1891 procedure_manager.clone(),
1892 Arc::new(DummyRepartitionProcedureFactory),
1893 );
1894 ddl_manager.register_loaders().unwrap();
1895
1896 let expected_loaders = vec![
1897 CreateTableProcedure::TYPE_NAME,
1898 AlterTableProcedure::TYPE_NAME,
1899 DropTableProcedure::TYPE_NAME,
1900 TruncateTableProcedure::TYPE_NAME,
1901 ];
1902
1903 for loader in expected_loaders {
1904 assert!(procedure_manager.contains_loader(loader));
1905 }
1906
1907 let soft_drop_loaders = [
1908 "metasrv-procedure::UndropTable",
1909 "metasrv-procedure::PurgeDroppedTable",
1910 "metasrv-procedure::PurgeExpiredDroppedTable",
1911 ];
1912 for loader in soft_drop_loaders {
1913 assert_eq!(
1914 cfg!(feature = "enterprise"),
1915 procedure_manager.contains_loader(loader)
1916 );
1917 }
1918 }
1919
1920 fn build_soft_drop_test_ddl_manager() -> DdlManager {
1921 let kv_backend = Arc::new(MemoryKvBackend::new());
1922 let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1923 let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1924 Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1925 Arc::new(WalProvider::default()),
1926 ));
1927 let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1928 let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1929 Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1930 ));
1931
1932 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1933 let poison_manager = Arc::new(InMemoryPoisonStore::default());
1934 let procedure_manager = Arc::new(LocalManager::new(
1935 Default::default(),
1936 state_store,
1937 poison_manager,
1938 None,
1939 None,
1940 ));
1941
1942 DdlManager::new(
1943 DdlContext {
1944 node_manager: Arc::new(DummyDatanodeManager),
1945 cache_invalidator: Arc::new(DummyCacheInvalidator),
1946 table_metadata_manager,
1947 table_metadata_allocator,
1948 flow_metadata_manager,
1949 flow_metadata_allocator,
1950 memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1951 leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1952 region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1953 soft_drop_enabled: true,
1954 soft_drop_retention: Some(std::time::Duration::from_secs(1)),
1955 create_database_metadata_committer: None,
1956 },
1957 procedure_manager,
1958 Arc::new(DummyRepartitionProcedureFactory),
1959 )
1960 }
1961
1962 #[cfg(feature = "enterprise")]
1963 #[tokio::test]
1964 async fn test_reserved_sequence_range_is_rejected_before_custom_handler() {
1965 let handler = Arc::new(RecordingCreateFlowHandler::default());
1966 let ddl_manager =
1967 build_soft_drop_test_ddl_manager().with_create_flow_handler(handler.clone());
1968 let mut task = test_create_flow_task();
1969 task.flow_options.insert(
1970 FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_KEY.to_string(),
1971 FLOW_EXPERIMENTAL_ENABLE_INCREMENTAL_READ_SEQUENCE_RANGE.to_string(),
1972 );
1973
1974 let result = ddl_manager
1975 .submit_ddl_task(
1976 ExecutorContext {
1977 query_context: Some(QueryContext::default()),
1978 ..Default::default()
1979 },
1980 SubmitDdlTaskRequest::new(DdlTask::new_create_flow(task)),
1981 )
1982 .await;
1983
1984 assert!(result.is_err());
1985 assert!(handler.tasks.lock().unwrap().is_empty());
1986 }
1987
1988 #[cfg(feature = "enterprise")]
1989 #[tokio::test]
1990 async fn test_create_flow_handler_dispatches_without_procedure() {
1991 let response = SubmitDdlTaskResponse {
1992 key: b"custom-handler".to_vec(),
1993 ..Default::default()
1994 };
1995 let handler = Arc::new(RecordingCreateFlowHandler {
1996 response: Mutex::new(Some(response)),
1997 ..Default::default()
1998 });
1999 let ddl_manager =
2000 build_soft_drop_test_ddl_manager().with_create_flow_handler(handler.clone());
2001 let procedure_context = ProcedureContext {
2002 actor: Some("test-user".to_string()),
2003 event_context: Some(
2004 PersistentEventContext::new(TriggerReason::Manual).with_protocol("mysql"),
2005 ),
2006 };
2007 let query_context = QueryContext {
2008 channel: Channel::Mysql as u8,
2009 ..Default::default()
2010 };
2011 let actual = ddl_manager
2012 .submit_ddl_task(
2013 ExecutorContext {
2014 query_context: Some(query_context.clone()),
2015 actor: procedure_context.actor.clone(),
2016 event_input: Some(ProcedureEventInput::new(TriggerReason::Manual)),
2017 ..Default::default()
2018 },
2019 SubmitDdlTaskRequest::new(DdlTask::new_create_flow(test_create_flow_task())),
2020 )
2021 .await
2022 .unwrap();
2023 assert_eq!(actual.key, b"custom-handler");
2024 assert_eq!(handler.tasks.lock().unwrap().len(), 1);
2025 assert_eq!(
2026 handler.contexts.lock().unwrap().as_slice(),
2027 &[(query_context, procedure_context)]
2028 );
2029 }
2030
2031 #[cfg(feature = "enterprise")]
2032 #[tokio::test]
2033 async fn test_create_flow_handler_error_does_not_fallback() {
2034 let handler = Arc::new(RecordingCreateFlowHandler {
2035 fail: true,
2036 ..Default::default()
2037 });
2038 let ddl_manager =
2039 build_soft_drop_test_ddl_manager().with_create_flow_handler(handler.clone());
2040 let err = ddl_manager
2041 .submit_ddl_task(
2042 ExecutorContext {
2043 query_context: Some(QueryContext::default()),
2044 ..Default::default()
2045 },
2046 SubmitDdlTaskRequest::new(DdlTask::new_create_flow(test_create_flow_task())),
2047 )
2048 .await
2049 .unwrap_err();
2050 assert!(err.to_string().contains("test create flow handler"));
2051 assert_eq!(handler.tasks.lock().unwrap().len(), 1);
2052 }
2053
2054 #[cfg(feature = "enterprise")]
2055 #[tokio::test]
2056 async fn test_create_flow_without_handler_uses_ordinary_procedure() {
2057 let ddl_manager = build_soft_drop_test_ddl_manager();
2058 ddl_manager.procedure_manager.start().await.unwrap();
2059 let err = ddl_manager
2060 .submit_ddl_task(
2061 ExecutorContext {
2062 query_context: Some(QueryContext::default()),
2063 ..Default::default()
2064 },
2065 SubmitDdlTaskRequest::new(DdlTask::new_create_flow(test_create_flow_task())),
2066 )
2067 .await
2068 .unwrap_err();
2069 assert!(!err.to_string().contains("test create flow handler"));
2070 assert_eq!(
2071 ddl_manager
2072 .procedure_manager
2073 .list_procedures()
2074 .await
2075 .unwrap()
2076 .iter()
2077 .filter(|procedure| procedure.type_name == CreateFlowProcedure::TYPE_NAME)
2078 .count(),
2079 1
2080 );
2081 }
2082
2083 #[cfg(feature = "enterprise")]
2084 #[tokio::test]
2085 async fn test_trigger_ddl_forwards_procedure_context() {
2086 let trigger_ddl_manager = Arc::new(RecordingTriggerDdlManager::default());
2087 let ddl_manager = build_soft_drop_test_ddl_manager()
2088 .with_trigger_ddl_manager(trigger_ddl_manager.clone());
2089 let procedure_context = ProcedureContext {
2090 actor: Some("test-user".to_string()),
2091 event_context: Some(
2092 PersistentEventContext::new(TriggerReason::Manual).with_protocol("mysql"),
2093 ),
2094 };
2095 let executor_context = || ExecutorContext {
2096 query_context: Some(QueryContext {
2097 channel: Channel::Mysql as u8,
2098 ..Default::default()
2099 }),
2100 actor: Some("test-user".to_string()),
2101 event_input: Some(ProcedureEventInput::new(TriggerReason::Manual)),
2102 ..Default::default()
2103 };
2104
2105 ddl_manager
2106 .submit_ddl_task(
2107 executor_context(),
2108 SubmitDdlTaskRequest::new(DdlTask::CreateTrigger(CreateTriggerTask {
2109 catalog_name: "greptime".to_string(),
2110 trigger_name: "test_trigger".to_string(),
2111 if_not_exists: false,
2112 sql: "SELECT 1".to_string(),
2113 channels: vec![],
2114 labels: Default::default(),
2115 annotations: Default::default(),
2116 interval: Duration::from_secs(1),
2117 raw_interval_expr: None,
2118 r#for: None,
2119 for_raw_expr: None,
2120 keep_firing_for: None,
2121 keep_firing_for_raw_expr: None,
2122 })),
2123 )
2124 .await
2125 .unwrap();
2126 ddl_manager
2127 .submit_ddl_task(
2128 executor_context(),
2129 SubmitDdlTaskRequest::new(DdlTask::DropTrigger(DropTriggerTask {
2130 catalog_name: "greptime".to_string(),
2131 trigger_name: "test_trigger".to_string(),
2132 drop_if_exists: false,
2133 })),
2134 )
2135 .await
2136 .unwrap();
2137
2138 assert_eq!(
2139 *trigger_ddl_manager.procedure_contexts.lock().unwrap(),
2140 vec![procedure_context.clone(), procedure_context]
2141 );
2142 }
2143
2144 #[cfg(feature = "enterprise")]
2145 #[tokio::test]
2146 async fn test_submit_undrop_missing_tombstone_returns_table_not_found_directly() {
2147 let ddl_manager = build_soft_drop_test_ddl_manager();
2148
2149 let err = ddl_manager
2150 .submit_undrop_table_task(
2151 UndropTableTask { table_id: 1024 },
2152 ProcedureContext::default(),
2153 )
2154 .await
2155 .unwrap_err();
2156
2157 assert_eq!(err.status_code(), StatusCode::TableNotFound);
2158 assert!(matches!(err, crate::error::Error::TableNotFound { .. }));
2159 }
2160
2161 #[cfg(not(feature = "enterprise"))]
2162 #[tokio::test]
2163 async fn test_submit_undrop_and_purge_rejected_in_non_enterprise_build() {
2164 let ddl_manager = build_soft_drop_test_ddl_manager();
2165
2166 let err = ddl_manager
2167 .submit_undrop_table_task(
2168 UndropTableTask { table_id: 1024 },
2169 ProcedureContext::default(),
2170 )
2171 .await
2172 .unwrap_err();
2173 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
2174
2175 let err = ddl_manager
2176 .submit_purge_dropped_table_task(
2177 PurgeDroppedTableTask { table_id: 1024 },
2178 ProcedureContext::default(),
2179 )
2180 .await
2181 .unwrap_err();
2182 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
2183
2184 let err = ddl_manager
2185 .submit_expired_purge_dropped_table_task(PurgeDroppedTableTask { table_id: 1024 })
2186 .await
2187 .unwrap_err();
2188 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
2189
2190 for task in [
2191 DdlTask::UndropTable(UndropTableTask { table_id: 1024 }),
2192 DdlTask::PurgeDroppedTable(PurgeDroppedTableTask { table_id: 1024 }),
2193 ] {
2194 let err = ddl_manager
2195 .submit_ddl_task(
2196 ExecutorContext {
2197 query_context: Some(QueryContext::default()),
2198 ..Default::default()
2199 },
2200 SubmitDdlTaskRequest::new(task),
2201 )
2202 .await
2203 .unwrap_err();
2204 assert!(matches!(err, crate::error::Error::Unsupported { .. }));
2205 }
2206 }
2207
2208 async fn ddl_manager_with_context(ddl_context: DdlContext) -> DdlManager {
2209 let kv_backend = Arc::new(MemoryKvBackend::new());
2210 let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
2211 let poison_manager = Arc::new(InMemoryPoisonStore::default());
2212 let procedure_manager = Arc::new(LocalManager::new(
2213 Default::default(),
2214 state_store,
2215 poison_manager,
2216 None,
2217 None,
2218 ));
2219 procedure_manager.start().await.unwrap();
2220 let ddl_manager = DdlManager::new(
2221 ddl_context,
2222 procedure_manager,
2223 Arc::new(DummyRepartitionProcedureFactory),
2224 );
2225 ddl_manager.register_loaders().unwrap();
2226 ddl_manager
2227 }
2228
2229 fn set_options_expr(table_name: &str, options: &[(&str, &str)]) -> api::v1::AlterTableExpr {
2230 api::v1::AlterTableExpr {
2231 catalog_name: common_catalog::consts::DEFAULT_CATALOG_NAME.to_string(),
2232 schema_name: common_catalog::consts::DEFAULT_SCHEMA_NAME.to_string(),
2233 table_name: table_name.to_string(),
2234 kind: Some(api::v1::alter_table_expr::Kind::SetTableOptions(
2235 api::v1::SetTableOptions {
2236 table_options: options
2237 .iter()
2238 .map(|(key, value)| api::v1::Option {
2239 key: key.to_string(),
2240 value: value.to_string(),
2241 })
2242 .collect(),
2243 },
2244 )),
2245 }
2246 }
2247
2248 #[tokio::test]
2249 async fn test_logical_table_annotation_alter_routing() {
2250 let (tx, mut rx) = tokio::sync::mpsc::channel(8);
2251 let node_manager = Arc::new(crate::test_util::MockDatanodeManager::new(
2252 crate::ddl::test_util::datanode_handler::DatanodeWatcher::new(tx),
2253 ));
2254 let ddl_context = crate::test_util::new_ddl_context(node_manager);
2255 let phy_id = crate::ddl::test_util::create_physical_table(&ddl_context, "phy").await;
2256 let logical_id =
2257 crate::ddl::test_util::create_logical_table(ddl_context.clone(), phy_id, "logical")
2258 .await;
2259 let ddl_manager = ddl_manager_with_context(ddl_context.clone()).await;
2260
2261 ddl_manager
2264 .submit_ddl_task(
2265 ExecutorContext {
2266 query_context: Some(QueryContext::default()),
2267 ..Default::default()
2268 },
2269 SubmitDdlTaskRequest::new(DdlTask::new_alter_table(set_options_expr(
2270 "logical",
2271 &[("greptime.semantic.signal_type", "metric")],
2272 ))),
2273 )
2274 .await
2275 .unwrap();
2276 rx.try_recv().unwrap_err();
2277 let table_info = ddl_manager
2278 .table_metadata_manager()
2279 .table_info_manager()
2280 .get(logical_id)
2281 .await
2282 .unwrap()
2283 .unwrap()
2284 .into_inner()
2285 .table_info;
2286 assert_eq!(
2287 table_info
2288 .meta
2289 .options
2290 .extra_options
2291 .get("greptime.semantic.signal_type"),
2292 Some(&"metric".to_string())
2293 );
2294
2295 let err = ddl_manager
2297 .submit_ddl_task(
2298 ExecutorContext {
2299 query_context: Some(QueryContext::default()),
2300 ..Default::default()
2301 },
2302 SubmitDdlTaskRequest::new(DdlTask::new_alter_table(set_options_expr(
2303 "logical",
2304 &[("greptime.semantic.source", "prometheus"), ("ttl", "7d")],
2305 ))),
2306 )
2307 .await
2308 .unwrap_err();
2309 let msg = common_error::ext::ErrorExt::output_msg(&err);
2310 assert!(msg.contains("must be altered separately"), "{msg}");
2311
2312 for (key, value) in [
2315 (table::requests::REPARTITION_COLUMN_HINT_KEY, "host"),
2316 (table::requests::REPARTITION_PARTITION_NUM_HINT_KEY, "8"),
2317 ] {
2318 let err = ddl_manager
2319 .submit_ddl_task(
2320 ExecutorContext {
2321 query_context: Some(QueryContext::default()),
2322 ..Default::default()
2323 },
2324 SubmitDdlTaskRequest::new(DdlTask::new_alter_table(set_options_expr(
2325 "logical",
2326 &[(key, value)],
2327 ))),
2328 )
2329 .await
2330 .unwrap_err();
2331 let msg = common_error::ext::ErrorExt::output_msg(&err);
2332 assert!(msg.contains("non-physical TableRouteValue"), "{msg}");
2333 }
2334 }
2335}