Skip to main content

common_meta/
ddl_manager.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::sync::Arc;
16use std::time::Duration;
17
18use api::v1::Repartition;
19use api::v1::alter_table_expr::Kind;
20use api::v1::repartition::Source as PbRepartitionSource;
21use common_error::ext::BoxedError;
22use common_event_recorder::PersistentEventContext;
23#[cfg(feature = "enterprise")]
24use common_event_recorder::TriggerReason;
25use common_procedure::{
26    BoxedProcedure, BoxedProcedureLoader, Output, ProcedureContext, ProcedureId,
27    ProcedureManagerRef, ProcedureWithId, watcher,
28};
29use common_telemetry::tracing_context::{FutureExt, TracingContext};
30use common_telemetry::{debug, info, tracing};
31use derive_builder::Builder;
32use snafu::{OptionExt, ResultExt, ensure};
33use store_api::storage::{RegionId, TableId};
34use table::table_name::TableName;
35
36use crate::ddl::alter_database::AlterDatabaseProcedure;
37use crate::ddl::alter_logical_tables::AlterLogicalTablesProcedure;
38use crate::ddl::alter_table::{AlterTableProcedure, RegionRouteChanged, only_enables_skip_wal};
39use crate::ddl::comment_on::CommentOnProcedure;
40use crate::ddl::create_database::{CreateDatabaseMetadataCommitterRef, CreateDatabaseProcedure};
41use crate::ddl::create_flow::CreateFlowProcedure;
42use crate::ddl::create_logical_tables::CreateLogicalTablesProcedure;
43use crate::ddl::create_table::CreateTableProcedure;
44use crate::ddl::create_view::CreateViewProcedure;
45use crate::ddl::drop_database::DropDatabaseProcedure;
46use crate::ddl::drop_flow::DropFlowProcedure;
47use crate::ddl::drop_table::DropTableProcedure;
48use crate::ddl::drop_view::DropViewProcedure;
49#[cfg(feature = "enterprise")]
50use crate::ddl::purge_dropped_table::PurgeDroppedTableProcedure;
51use crate::ddl::truncate_table::TruncateTableProcedure;
52#[cfg(feature = "enterprise")]
53use crate::ddl::undrop_table::UndropTableProcedure;
54use crate::ddl::{DdlContext, utils};
55use crate::error::{
56    self, CreateRepartitionProcedureSnafu, EmptyDdlTasksSnafu,
57    PersistRepartitionGcRequirementSnafu, ProcedureOutputSnafu, RegisterProcedureLoaderSnafu,
58    RegisterRepartitionProcedureLoaderSnafu, Result, SubmitProcedureSnafu, TableInfoNotFoundSnafu,
59    TableNotFoundSnafu, TableRouteNotFoundSnafu, UnexpectedLogicalRouteTableSnafu,
60    UnsupportedSnafu, WaitProcedureSnafu,
61};
62use crate::key::table_info::TableInfoValue;
63use crate::key::table_name::TableNameKey;
64use crate::key::{DeserializedValueWithBytes, TableMetadataManagerRef};
65use crate::procedure_executor::ExecutorContext;
66#[cfg(feature = "enterprise")]
67use crate::rpc::ddl::DdlTask::CreateTrigger;
68#[cfg(feature = "enterprise")]
69use crate::rpc::ddl::DdlTask::DropTrigger;
70use crate::rpc::ddl::DdlTask::{
71    AlterDatabase, AlterLogicalTables, AlterTable, CommentOn, CreateDatabase, CreateFlow,
72    CreateLogicalTables, CreateTable, CreateView, DropDatabase, DropFlow, DropLogicalTables,
73    DropTable, DropView, PurgeDroppedTable, TruncateTable, UndropTable,
74};
75#[cfg(feature = "enterprise")]
76use crate::rpc::ddl::trigger::CreateTriggerTask;
77#[cfg(feature = "enterprise")]
78use crate::rpc::ddl::trigger::DropTriggerTask;
79use crate::rpc::ddl::{
80    AlterDatabaseTask, AlterTableTask, CommentOnTask, CreateDatabaseTask, CreateFlowTask,
81    CreateTableTask, CreateViewTask, DropDatabaseTask, DropFlowTask, DropTableTask, DropViewTask,
82    PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest, SubmitDdlTaskResponse,
83    TruncateTableTask, UndropTableTask,
84};
85
86const MAX_REGION_ROUTE_CHANGE_RETRIES: usize = 3;
87
88/// A configurator that customizes or enhances a [`DdlManager`].
89#[async_trait::async_trait]
90pub trait DdlManagerConfigurator<C>: Send + Sync {
91    /// Configures the given [`DdlManager`] using the provided [`DdlManagerConfigureContext`].
92    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/// The [DdlManager] provides the ability to execute Ddl.
106#[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/// This trait is responsible for handling DDL tasks about triggers. e.g.,
116/// create trigger, drop trigger, etc.
117#[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        /// Full target partition columns to overwrite table metadata.
172        ///
173        /// `None` means the repartition keeps using the current table
174        /// partition columns, so the procedure won't update
175        /// `partition_key_indices`.
176        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    /// Persists the cluster-level GC requirement before submitting a repartition.
203    async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError>;
204}
205
206/// The options for DDL tasks.
207///
208/// Note: These options may not be utilized by all procedures.
209/// At present, they are specifically applied in `RepartitionProcedure`.
210#[derive(Debug, Clone, Copy)]
211pub struct DdlOptions {
212    /// The timeout will be passed to the procedure.
213    ///
214    /// Note: Each procedure may implement its own timeout handling mechanism.
215    pub timeout: Duration,
216    /// The flag that controls whether to wait for the procedure to complete.
217    ///
218    /// If wait is `true`, the procedure will wait for completion(success or failure) and the result will be returned.
219    /// Otherwise, the procedure will be submitted and return the [ProcedureId](common_procedure::ProcedureId) immediately.
220    ///
221    /// Note: The value of `wait` is independent of the `timeout` option. If a procedure ignores the `timeout` and `wait` is set to true, the operation returns until the procedure completes.
222    pub wait: bool,
223}
224
225impl DdlManager {
226    /// Returns a new [DdlManager].
227    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    /// Returns the [TableMetadataManagerRef].
256    pub fn table_metadata_manager(&self) -> &TableMetadataManagerRef {
257        &self.ddl_context.table_metadata_manager
258    }
259
260    /// Returns the [DdlContext]
261    pub fn create_context(&self) -> DdlContext {
262        self.ddl_context.clone()
263    }
264
265    /// Registers all Ddl loaders.
266    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    /// Submits a repartition procedure for the specified table.
323    ///
324    /// This creates a repartition procedure using the provided `table_id`,
325    /// `table_name`, and `Repartition` configuration, and then either executes it
326    /// to completion or just submits it for asynchronous execution.
327    ///
328    /// The `Repartition` argument contains the original (`from_partition_exprs`)
329    /// and target (`into_partition_exprs`) partition expressions that define how
330    /// the table should be repartitioned.
331    ///
332    /// The `wait` flag controls whether this method waits for the repartition
333    /// procedure to finish:
334    /// - If `wait` is `true`, the procedure is executed and this method awaits
335    ///   its completion, returning both the generated `ProcedureId` and the
336    ///   final `Output` of the procedure.
337    /// - If `wait` is `false`, the procedure is only submitted to the procedure
338    ///   manager for asynchronous execution, and this method returns the
339    ///   `ProcedureId` along with `None` as the output.
340    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                // Reads the deprecated field for backward compatibility with old persisted DDL tasks.
366                #[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    /// Submits and executes an alter table task.
401    #[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        // make alter_table_task mutable so we can call .take() on its field
410        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            // Hold the same physical region locks as migration while validating that
441            // this route snapshot is still the one the procedure will mutate.
442            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    /// Submits and executes a create table task.
490    #[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    /// Submits and executes a `[CreateViewTask]`.
512    #[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    /// Submits and executes a create multiple logical table tasks.
529    #[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    /// Submits and executes alter multiple table tasks.
548    #[tracing::instrument(skip_all)]
549    pub async fn submit_alter_logical_table_tasks(
550        &self,
551        alter_table_tasks: Vec<AlterTableTask>,
552        physical_table_id: TableId,
553        procedure_context: ProcedureContext,
554    ) -> Result<(ProcedureId, Option<Output>)> {
555        let context = self.create_context();
556
557        let procedure =
558            AlterLogicalTablesProcedure::new(alter_table_tasks, physical_table_id, context);
559
560        let procedure_with_id =
561            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
562
563        self.execute_procedure_and_wait(procedure_with_id).await
564    }
565
566    /// Submits and executes a drop table task.
567    #[tracing::instrument(skip_all)]
568    pub async fn submit_drop_table_task(
569        &self,
570        drop_table_task: DropTableTask,
571        procedure_context: ProcedureContext,
572    ) -> Result<(ProcedureId, Option<Output>)> {
573        let context = self.create_context();
574
575        let procedure = DropTableProcedure::new(drop_table_task, context);
576
577        let procedure_with_id =
578            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
579
580        self.execute_procedure_and_wait(procedure_with_id).await
581    }
582
583    /// Submits and executes an undrop table task.
584    #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
585    #[tracing::instrument(skip_all)]
586    pub async fn submit_undrop_table_task(
587        &self,
588        undrop_table_task: UndropTableTask,
589        procedure_context: ProcedureContext,
590    ) -> Result<(ProcedureId, Option<Output>)> {
591        #[cfg(not(feature = "enterprise"))]
592        {
593            use crate::error::UnsupportedSnafu;
594
595            return UnsupportedSnafu {
596                operation: "undrop table is only available in GreptimeDB Enterprise Edition",
597            }
598            .fail();
599        }
600        #[cfg(feature = "enterprise")]
601        {
602            let context = self.create_context();
603            let original_table_name = context
604                .table_metadata_manager
605                .get_dropped_table_by_id(undrop_table_task.table_id)
606                .await?
607                .with_context(|| TableNotFoundSnafu {
608                    table_name: undrop_table_task.table_id.to_string(),
609                })?
610                .table_name;
611            let procedure = UndropTableProcedure::new_with_original_table_name(
612                undrop_table_task,
613                context,
614                Some(original_table_name),
615            );
616            let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
617                .with_context(procedure_context);
618
619            self.execute_procedure_and_wait(procedure_with_id).await
620        }
621    }
622
623    /// Submits and executes a purge dropped table task.
624    #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
625    #[tracing::instrument(skip_all)]
626    pub async fn submit_purge_dropped_table_task(
627        &self,
628        purge_dropped_table_task: PurgeDroppedTableTask,
629        procedure_context: ProcedureContext,
630    ) -> Result<(ProcedureId, Option<Output>)> {
631        #[cfg(not(feature = "enterprise"))]
632        {
633            use crate::error::UnsupportedSnafu;
634
635            return UnsupportedSnafu {
636                operation: "purge dropped table is only available in GreptimeDB Enterprise Edition",
637            }
638            .fail();
639        }
640        #[cfg(feature = "enterprise")]
641        {
642            let context = self.create_context();
643            let procedure = PurgeDroppedTableProcedure::new(purge_dropped_table_task, context);
644            let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
645                .with_context(procedure_context);
646
647            self.execute_procedure_and_wait(procedure_with_id).await
648        }
649    }
650
651    /// Submits and executes a purge task that first rechecks the tombstone deadline.
652    #[cfg_attr(not(feature = "enterprise"), allow(unused_variables))]
653    #[tracing::instrument(skip_all)]
654    pub async fn submit_expired_purge_dropped_table_task(
655        &self,
656        purge_dropped_table_task: PurgeDroppedTableTask,
657    ) -> Result<(ProcedureId, Option<Output>)> {
658        #[cfg(not(feature = "enterprise"))]
659        {
660            use crate::error::UnsupportedSnafu;
661
662            return UnsupportedSnafu {
663                operation:
664                    "purge expired dropped table is only available in GreptimeDB Enterprise Edition",
665            }
666            .fail();
667        }
668        #[cfg(feature = "enterprise")]
669        {
670            let context = self.create_context();
671            let procedure_context = ProcedureContext::from_event_context(
672                PersistentEventContext::new(TriggerReason::ScheduledGc),
673            );
674            let procedure =
675                PurgeDroppedTableProcedure::new_if_expired(purge_dropped_table_task, context);
676            let procedure_with_id = ProcedureWithId::with_random_id(Box::new(procedure))
677                .with_context(procedure_context);
678
679            self.execute_procedure_and_wait(procedure_with_id).await
680        }
681    }
682
683    /// Submits and executes a create database task.
684    #[tracing::instrument(skip_all)]
685    pub async fn submit_create_database(
686        &self,
687        CreateDatabaseTask {
688            catalog,
689            schema,
690            create_if_not_exists,
691            options,
692            creator,
693        }: CreateDatabaseTask,
694        procedure_context: ProcedureContext,
695    ) -> Result<(ProcedureId, Option<Output>)> {
696        let context = self.create_context();
697        let procedure = CreateDatabaseProcedure::new(
698            catalog,
699            schema,
700            create_if_not_exists,
701            options,
702            creator,
703            context,
704        );
705        let procedure_with_id =
706            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
707
708        self.execute_procedure_and_wait(procedure_with_id).await
709    }
710
711    /// Submits and executes a drop table task.
712    #[tracing::instrument(skip_all)]
713    pub async fn submit_drop_database(
714        &self,
715        DropDatabaseTask {
716            catalog,
717            schema,
718            drop_if_exists,
719        }: DropDatabaseTask,
720        procedure_context: ProcedureContext,
721    ) -> Result<(ProcedureId, Option<Output>)> {
722        let context = self.create_context();
723        let procedure = DropDatabaseProcedure::new(catalog, schema, drop_if_exists, context);
724        let procedure_with_id =
725            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
726
727        self.execute_procedure_and_wait(procedure_with_id).await
728    }
729
730    pub async fn submit_alter_database(
731        &self,
732        alter_database_task: AlterDatabaseTask,
733        procedure_context: ProcedureContext,
734    ) -> Result<(ProcedureId, Option<Output>)> {
735        let context = self.create_context();
736        let procedure = AlterDatabaseProcedure::new(alter_database_task, context)?;
737        let procedure_with_id =
738            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
739
740        self.execute_procedure_and_wait(procedure_with_id).await
741    }
742
743    /// Submits and executes a create flow task.
744    #[tracing::instrument(skip_all)]
745    pub async fn submit_create_flow_task(
746        &self,
747        create_flow: CreateFlowTask,
748        query_context: QueryContext,
749        procedure_context: ProcedureContext,
750    ) -> Result<(ProcedureId, Option<Output>)> {
751        let context = self.create_context();
752        let procedure = CreateFlowProcedure::new(create_flow, query_context, context);
753        let procedure_with_id =
754            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
755
756        self.execute_procedure_and_wait(procedure_with_id).await
757    }
758
759    /// Submits and executes a drop flow task.
760    #[tracing::instrument(skip_all)]
761    pub async fn submit_drop_flow_task(
762        &self,
763        drop_flow: DropFlowTask,
764        procedure_context: ProcedureContext,
765    ) -> Result<(ProcedureId, Option<Output>)> {
766        let context = self.create_context();
767        let procedure = DropFlowProcedure::new(drop_flow, context);
768        let procedure_with_id =
769            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
770
771        self.execute_procedure_and_wait(procedure_with_id).await
772    }
773
774    /// Submits and executes a drop view task.
775    #[tracing::instrument(skip_all)]
776    pub async fn submit_drop_view_task(
777        &self,
778        drop_view: DropViewTask,
779        procedure_context: ProcedureContext,
780    ) -> Result<(ProcedureId, Option<Output>)> {
781        let context = self.create_context();
782        let procedure = DropViewProcedure::new(drop_view, context);
783        let procedure_with_id =
784            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
785
786        self.execute_procedure_and_wait(procedure_with_id).await
787    }
788
789    /// Submits and executes a truncate table task.
790    #[tracing::instrument(skip_all)]
791    pub async fn submit_truncate_table_task(
792        &self,
793        truncate_table_task: TruncateTableTask,
794        table_info_value: DeserializedValueWithBytes<TableInfoValue>,
795        procedure_context: ProcedureContext,
796    ) -> Result<(ProcedureId, Option<Output>)> {
797        let context = self.create_context();
798        let procedure = TruncateTableProcedure::new(truncate_table_task, table_info_value, context);
799
800        let procedure_with_id =
801            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
802
803        self.execute_procedure_and_wait(procedure_with_id).await
804    }
805
806    /// Submits and executes a comment on task.
807    #[tracing::instrument(skip_all)]
808    pub async fn submit_comment_on_task(
809        &self,
810        mut comment_on_task: CommentOnTask,
811        procedure_context: ProcedureContext,
812    ) -> Result<(ProcedureId, Option<Output>)> {
813        let context = self.create_context();
814        comment_on_task
815            .enrich_object_id(
816                context.table_metadata_manager.table_name_manager(),
817                context.flow_metadata_manager.flow_name_manager(),
818            )
819            .await?;
820        let procedure = CommentOnProcedure::new(comment_on_task, context);
821        let procedure_with_id =
822            ProcedureWithId::with_random_id(Box::new(procedure)).with_context(procedure_context);
823
824        self.execute_procedure_and_wait(procedure_with_id).await
825    }
826
827    /// Executes a procedure and waits for the result.
828    async fn execute_procedure_and_wait(
829        &self,
830        procedure_with_id: ProcedureWithId,
831    ) -> Result<(ProcedureId, Option<Output>)> {
832        let procedure_id = procedure_with_id.id;
833
834        let mut watcher = self
835            .procedure_manager
836            .submit(procedure_with_id)
837            .await
838            .context(SubmitProcedureSnafu)?;
839
840        let output = watcher::wait(&mut watcher)
841            .await
842            .context(WaitProcedureSnafu)?;
843
844        Ok((procedure_id, output))
845    }
846
847    /// Submits a procedure and returns the procedure id.
848    async fn submit_procedure(&self, procedure_with_id: ProcedureWithId) -> Result<ProcedureId> {
849        let procedure_id = procedure_with_id.id;
850        let _ = self
851            .procedure_manager
852            .submit(procedure_with_id)
853            .await
854            .context(SubmitProcedureSnafu)?;
855
856        Ok(procedure_id)
857    }
858
859    pub async fn submit_ddl_task(
860        &self,
861        context: ExecutorContext,
862        request: SubmitDdlTaskRequest,
863    ) -> Result<SubmitDdlTaskResponse> {
864        let ExecutorContext {
865            tracing_context,
866            query_context,
867            actor,
868            event_input,
869        } = context;
870        let query_context = query_context.context(UnsupportedSnafu {
871            operation: "submit_ddl_task without query context",
872        })?;
873        let procedure_context = ProcedureContext {
874            actor,
875            event_context: event_input
876                .map(|input| PersistentEventContext::from((input, query_context.protocol()))),
877        };
878        let span = tracing_context
879            .as_ref()
880            .map(TracingContext::from_w3c)
881            .unwrap_or_else(TracingContext::from_current_span)
882            .attach(tracing::info_span!("DdlManager::submit_ddl_task"));
883        let SubmitDdlTaskRequest {
884            wait,
885            timeout,
886            task,
887        } = request;
888        let ddl_options = DdlOptions { wait, timeout };
889        async move {
890            debug!("Submitting Ddl task: {:?}", task);
891            match task {
892                CreateTable(create_table_task) => {
893                    handle_create_table_task(
894                        self,
895                        create_table_task,
896                        query_context,
897                        procedure_context,
898                    )
899                    .await
900                }
901                DropTable(drop_table_task) => {
902                    handle_drop_table_task(self, drop_table_task, procedure_context).await
903                }
904                UndropTable(undrop_table_task) => {
905                    handle_undrop_table_task(self, undrop_table_task, procedure_context).await
906                }
907                PurgeDroppedTable(purge_dropped_table_task) => {
908                    handle_purge_dropped_table_task(
909                        self,
910                        purge_dropped_table_task,
911                        procedure_context,
912                    )
913                    .await
914                }
915                AlterTable(alter_table_task) => {
916                    handle_alter_table_task(self, alter_table_task, ddl_options, procedure_context)
917                        .await
918                }
919                TruncateTable(truncate_table_task) => {
920                    handle_truncate_table_task(self, truncate_table_task, procedure_context).await
921                }
922                CreateLogicalTables(create_table_tasks) => {
923                    handle_create_logical_table_tasks(self, create_table_tasks, procedure_context)
924                        .await
925                }
926                AlterLogicalTables(alter_table_tasks) => {
927                    handle_alter_logical_table_tasks(self, alter_table_tasks, procedure_context)
928                        .await
929                }
930                DropLogicalTables(_) => todo!(),
931                CreateDatabase(create_database_task) => {
932                    handle_create_database_task(self, create_database_task, procedure_context).await
933                }
934                DropDatabase(drop_database_task) => {
935                    handle_drop_database_task(self, drop_database_task, procedure_context).await
936                }
937                AlterDatabase(alter_database_task) => {
938                    handle_alter_database_task(self, alter_database_task, procedure_context).await
939                }
940                CreateFlow(create_flow_task) => {
941                    handle_create_flow_task(
942                        self,
943                        create_flow_task,
944                        query_context,
945                        procedure_context,
946                    )
947                    .await
948                }
949                DropFlow(drop_flow_task) => {
950                    handle_drop_flow_task(self, drop_flow_task, procedure_context).await
951                }
952                CreateView(create_view_task) => {
953                    handle_create_view_task(self, create_view_task, procedure_context).await
954                }
955                DropView(drop_view_task) => {
956                    handle_drop_view_task(self, drop_view_task, procedure_context).await
957                }
958                CommentOn(comment_on_task) => {
959                    handle_comment_on_task(self, comment_on_task, procedure_context).await
960                }
961                #[cfg(feature = "enterprise")]
962                CreateTrigger(create_trigger_task) => {
963                    handle_create_trigger_task(
964                        self,
965                        create_trigger_task,
966                        query_context,
967                        procedure_context,
968                    )
969                    .await
970                }
971                #[cfg(feature = "enterprise")]
972                DropTrigger(drop_trigger_task) => {
973                    handle_drop_trigger_task(
974                        self,
975                        drop_trigger_task,
976                        query_context,
977                        procedure_context,
978                    )
979                    .await
980                }
981            }
982        }
983        .trace(span)
984        .await
985    }
986}
987
988async fn handle_truncate_table_task(
989    ddl_manager: &DdlManager,
990    truncate_table_task: TruncateTableTask,
991    procedure_context: ProcedureContext,
992) -> Result<SubmitDdlTaskResponse> {
993    let table_id = truncate_table_task.table_id;
994    let table_metadata_manager = &ddl_manager.table_metadata_manager();
995    let table_ref = truncate_table_task.table_ref();
996
997    let table_info_value = table_metadata_manager
998        .table_info_manager()
999        .get(table_id)
1000        .await?
1001        .with_context(|| TableInfoNotFoundSnafu {
1002            table: table_ref.to_string(),
1003        })?;
1004    let physical_table_id = table_metadata_manager
1005        .table_route_manager()
1006        .get_physical_table_id(table_id)
1007        .await?;
1008    ensure!(
1009        physical_table_id == table_id,
1010        error::UnexpectedSnafu {
1011            err_msg: "Truncate table is only supported for physical tables."
1012        }
1013    );
1014
1015    let (id, _) = ddl_manager
1016        .submit_truncate_table_task(truncate_table_task, table_info_value, procedure_context)
1017        .await?;
1018
1019    info!("Table: {table_id} is truncated via procedure_id {id:?}");
1020
1021    Ok(SubmitDdlTaskResponse {
1022        key: id.to_string().into(),
1023        ..Default::default()
1024    })
1025}
1026
1027async fn handle_alter_table_task(
1028    ddl_manager: &DdlManager,
1029    alter_table_task: AlterTableTask,
1030    ddl_options: DdlOptions,
1031    procedure_context: ProcedureContext,
1032) -> Result<SubmitDdlTaskResponse> {
1033    let table_ref = alter_table_task.table_ref();
1034
1035    let table_id = ddl_manager
1036        .table_metadata_manager()
1037        .table_name_manager()
1038        .get(TableNameKey::new(
1039            table_ref.catalog,
1040            table_ref.schema,
1041            table_ref.table,
1042        ))
1043        .await?
1044        .with_context(|| TableNotFoundSnafu {
1045            table_name: table_ref.to_string(),
1046        })?
1047        .table_id();
1048
1049    let table_route_value = ddl_manager
1050        .table_metadata_manager()
1051        .table_route_manager()
1052        .table_route_storage()
1053        .get(table_id)
1054        .await?
1055        .context(TableRouteNotFoundSnafu { table_id })?;
1056    ensure!(
1057        table_route_value.is_physical(),
1058        UnexpectedLogicalRouteTableSnafu {
1059            err_msg: format!("{:?} is a non-physical TableRouteValue.", table_ref),
1060        }
1061    );
1062
1063    let (id, _) = ddl_manager
1064        .submit_alter_table_task(table_id, alter_table_task, procedure_context, ddl_options)
1065        .await?;
1066
1067    info!("Table: {table_id} is altered via procedure_id {id:?}");
1068
1069    Ok(SubmitDdlTaskResponse {
1070        key: id.to_string().into(),
1071        ..Default::default()
1072    })
1073}
1074
1075async fn handle_drop_table_task(
1076    ddl_manager: &DdlManager,
1077    drop_table_task: DropTableTask,
1078    procedure_context: ProcedureContext,
1079) -> Result<SubmitDdlTaskResponse> {
1080    let table_id = drop_table_task.table_id;
1081    let (id, _) = ddl_manager
1082        .submit_drop_table_task(drop_table_task, procedure_context)
1083        .await?;
1084
1085    info!("Table: {table_id} is dropped via procedure_id {id:?}");
1086
1087    Ok(SubmitDdlTaskResponse {
1088        key: id.to_string().into(),
1089        ..Default::default()
1090    })
1091}
1092
1093async fn handle_undrop_table_task(
1094    ddl_manager: &DdlManager,
1095    undrop_table_task: UndropTableTask,
1096    procedure_context: ProcedureContext,
1097) -> Result<SubmitDdlTaskResponse> {
1098    let table_id = undrop_table_task.table_id;
1099    let (id, _) = ddl_manager
1100        .submit_undrop_table_task(undrop_table_task, procedure_context)
1101        .await?;
1102
1103    info!("Table: {table_id} is undropped via procedure_id {id:?}");
1104
1105    Ok(SubmitDdlTaskResponse {
1106        key: id.to_string().into(),
1107        ..Default::default()
1108    })
1109}
1110
1111async fn handle_purge_dropped_table_task(
1112    ddl_manager: &DdlManager,
1113    purge_dropped_table_task: PurgeDroppedTableTask,
1114    procedure_context: ProcedureContext,
1115) -> Result<SubmitDdlTaskResponse> {
1116    let (id, _) = ddl_manager
1117        .submit_purge_dropped_table_task(purge_dropped_table_task, procedure_context)
1118        .await?;
1119
1120    info!("Dropped table is purged via procedure_id {id:?}");
1121
1122    Ok(SubmitDdlTaskResponse {
1123        key: id.to_string().into(),
1124        ..Default::default()
1125    })
1126}
1127
1128async fn handle_create_table_task(
1129    ddl_manager: &DdlManager,
1130    create_table_task: CreateTableTask,
1131    query_context: QueryContext,
1132    procedure_context: ProcedureContext,
1133) -> Result<SubmitDdlTaskResponse> {
1134    let (id, output) = ddl_manager
1135        .submit_create_table_task(create_table_task, query_context, procedure_context)
1136        .await?;
1137
1138    let procedure_id = id.to_string();
1139    let output = output.context(ProcedureOutputSnafu {
1140        procedure_id: &procedure_id,
1141        err_msg: "empty output",
1142    })?;
1143    let table_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1144        procedure_id: &procedure_id,
1145        err_msg: "downcast to `u32`",
1146    })?);
1147    info!("Table: {table_id} is created via procedure_id {id:?}");
1148
1149    Ok(SubmitDdlTaskResponse {
1150        key: procedure_id.into(),
1151        table_ids: vec![table_id],
1152    })
1153}
1154
1155async fn handle_create_logical_table_tasks(
1156    ddl_manager: &DdlManager,
1157    create_table_tasks: Vec<CreateTableTask>,
1158    procedure_context: ProcedureContext,
1159) -> Result<SubmitDdlTaskResponse> {
1160    ensure!(
1161        !create_table_tasks.is_empty(),
1162        EmptyDdlTasksSnafu {
1163            name: "create logical tables"
1164        }
1165    );
1166    let physical_table_id = utils::check_and_get_physical_table_id(
1167        ddl_manager.table_metadata_manager(),
1168        &create_table_tasks,
1169    )
1170    .await?;
1171    let num_logical_tables = create_table_tasks.len();
1172
1173    let (id, output) = ddl_manager
1174        .submit_create_logical_table_tasks(create_table_tasks, physical_table_id, procedure_context)
1175        .await?;
1176
1177    info!(
1178        "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is created via procedure_id {id:?}"
1179    );
1180
1181    let procedure_id = id.to_string();
1182    let output = output.context(ProcedureOutputSnafu {
1183        procedure_id: &procedure_id,
1184        err_msg: "empty output",
1185    })?;
1186    let table_ids = output
1187        .downcast_ref::<Vec<TableId>>()
1188        .context(ProcedureOutputSnafu {
1189            procedure_id: &procedure_id,
1190            err_msg: "downcast to `Vec<TableId>`",
1191        })?
1192        .clone();
1193
1194    Ok(SubmitDdlTaskResponse {
1195        key: procedure_id.into(),
1196        table_ids,
1197    })
1198}
1199
1200async fn handle_create_database_task(
1201    ddl_manager: &DdlManager,
1202    create_database_task: CreateDatabaseTask,
1203    procedure_context: ProcedureContext,
1204) -> Result<SubmitDdlTaskResponse> {
1205    let catalog = create_database_task.catalog.clone();
1206    let schema = create_database_task.schema.clone();
1207    let (id, _) = ddl_manager
1208        .submit_create_database(create_database_task, procedure_context)
1209        .await?;
1210
1211    let procedure_id = id.to_string();
1212    info!(
1213        "Database {}.{} is created via procedure_id {id:?}",
1214        catalog, schema
1215    );
1216
1217    Ok(SubmitDdlTaskResponse {
1218        key: procedure_id.into(),
1219        ..Default::default()
1220    })
1221}
1222
1223async fn handle_drop_database_task(
1224    ddl_manager: &DdlManager,
1225    drop_database_task: DropDatabaseTask,
1226    procedure_context: ProcedureContext,
1227) -> Result<SubmitDdlTaskResponse> {
1228    let (id, _) = ddl_manager
1229        .submit_drop_database(drop_database_task.clone(), procedure_context)
1230        .await?;
1231
1232    let procedure_id = id.to_string();
1233    info!(
1234        "Database {}.{} is dropped via procedure_id {id:?}",
1235        drop_database_task.catalog, drop_database_task.schema
1236    );
1237
1238    Ok(SubmitDdlTaskResponse {
1239        key: procedure_id.into(),
1240        ..Default::default()
1241    })
1242}
1243
1244async fn handle_alter_database_task(
1245    ddl_manager: &DdlManager,
1246    alter_database_task: AlterDatabaseTask,
1247    procedure_context: ProcedureContext,
1248) -> Result<SubmitDdlTaskResponse> {
1249    let (id, _) = ddl_manager
1250        .submit_alter_database(alter_database_task.clone(), procedure_context)
1251        .await?;
1252
1253    let procedure_id = id.to_string();
1254    info!(
1255        "Database {}.{} is altered via procedure_id {id:?}",
1256        alter_database_task.catalog(),
1257        alter_database_task.schema()
1258    );
1259
1260    Ok(SubmitDdlTaskResponse {
1261        key: procedure_id.into(),
1262        ..Default::default()
1263    })
1264}
1265
1266async fn handle_drop_flow_task(
1267    ddl_manager: &DdlManager,
1268    drop_flow_task: DropFlowTask,
1269    procedure_context: ProcedureContext,
1270) -> Result<SubmitDdlTaskResponse> {
1271    let (id, _) = ddl_manager
1272        .submit_drop_flow_task(drop_flow_task.clone(), procedure_context)
1273        .await?;
1274
1275    let procedure_id = id.to_string();
1276    info!(
1277        "Flow {}.{}({}) is dropped via procedure_id {id:?}",
1278        drop_flow_task.catalog_name, drop_flow_task.flow_name, drop_flow_task.flow_id,
1279    );
1280
1281    Ok(SubmitDdlTaskResponse {
1282        key: procedure_id.into(),
1283        ..Default::default()
1284    })
1285}
1286
1287#[cfg(feature = "enterprise")]
1288async fn handle_drop_trigger_task(
1289    ddl_manager: &DdlManager,
1290    drop_trigger_task: DropTriggerTask,
1291    query_context: QueryContext,
1292    procedure_context: ProcedureContext,
1293) -> Result<SubmitDdlTaskResponse> {
1294    let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1295        use crate::error::UnsupportedSnafu;
1296
1297        return UnsupportedSnafu {
1298            operation: "drop trigger",
1299        }
1300        .fail();
1301    };
1302
1303    m.drop_trigger(
1304        drop_trigger_task,
1305        ddl_manager.procedure_manager.clone(),
1306        ddl_manager.ddl_context.clone(),
1307        query_context,
1308        procedure_context,
1309    )
1310    .await
1311}
1312
1313async fn handle_drop_view_task(
1314    ddl_manager: &DdlManager,
1315    drop_view_task: DropViewTask,
1316    procedure_context: ProcedureContext,
1317) -> Result<SubmitDdlTaskResponse> {
1318    let (id, _) = ddl_manager
1319        .submit_drop_view_task(drop_view_task.clone(), procedure_context)
1320        .await?;
1321
1322    let procedure_id = id.to_string();
1323    info!(
1324        "View {}({}) is dropped via procedure_id {id:?}",
1325        drop_view_task.table_ref(),
1326        drop_view_task.view_id,
1327    );
1328
1329    Ok(SubmitDdlTaskResponse {
1330        key: procedure_id.into(),
1331        ..Default::default()
1332    })
1333}
1334
1335async fn handle_create_flow_task(
1336    ddl_manager: &DdlManager,
1337    create_flow_task: CreateFlowTask,
1338    query_context: QueryContext,
1339    procedure_context: ProcedureContext,
1340) -> Result<SubmitDdlTaskResponse> {
1341    let (id, output) = ddl_manager
1342        .submit_create_flow_task(create_flow_task.clone(), query_context, procedure_context)
1343        .await?;
1344
1345    let procedure_id = id.to_string();
1346    let output = output.context(ProcedureOutputSnafu {
1347        procedure_id: &procedure_id,
1348        err_msg: "empty output",
1349    })?;
1350    let flow_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1351        procedure_id: &procedure_id,
1352        err_msg: "downcast to `u32`",
1353    })?);
1354    if !create_flow_task.or_replace {
1355        info!(
1356            "Flow {}.{}({flow_id}) is created via procedure_id {id:?}",
1357            create_flow_task.catalog_name, create_flow_task.flow_name,
1358        );
1359    } else {
1360        info!(
1361            "Flow {}.{}({flow_id}) is replaced via procedure_id {id:?}",
1362            create_flow_task.catalog_name, create_flow_task.flow_name,
1363        );
1364    }
1365
1366    Ok(SubmitDdlTaskResponse {
1367        key: procedure_id.into(),
1368        ..Default::default()
1369    })
1370}
1371
1372#[cfg(feature = "enterprise")]
1373async fn handle_create_trigger_task(
1374    ddl_manager: &DdlManager,
1375    create_trigger_task: CreateTriggerTask,
1376    query_context: QueryContext,
1377    procedure_context: ProcedureContext,
1378) -> Result<SubmitDdlTaskResponse> {
1379    let Some(m) = ddl_manager.trigger_ddl_manager.as_ref() else {
1380        use crate::error::UnsupportedSnafu;
1381
1382        return UnsupportedSnafu {
1383            operation: "create trigger",
1384        }
1385        .fail();
1386    };
1387
1388    m.create_trigger(
1389        create_trigger_task,
1390        ddl_manager.procedure_manager.clone(),
1391        ddl_manager.ddl_context.clone(),
1392        query_context,
1393        procedure_context,
1394    )
1395    .await
1396}
1397
1398async fn handle_alter_logical_table_tasks(
1399    ddl_manager: &DdlManager,
1400    alter_table_tasks: Vec<AlterTableTask>,
1401    procedure_context: ProcedureContext,
1402) -> Result<SubmitDdlTaskResponse> {
1403    ensure!(
1404        !alter_table_tasks.is_empty(),
1405        EmptyDdlTasksSnafu {
1406            name: "alter logical tables"
1407        }
1408    );
1409
1410    // Use the physical table id in the first logical table, then it will be checked in the procedure.
1411    let first_table = TableNameKey {
1412        catalog: &alter_table_tasks[0].alter_table.catalog_name,
1413        schema: &alter_table_tasks[0].alter_table.schema_name,
1414        table: &alter_table_tasks[0].alter_table.table_name,
1415    };
1416    let physical_table_id =
1417        utils::get_physical_table_id(ddl_manager.table_metadata_manager(), first_table).await?;
1418    let num_logical_tables = alter_table_tasks.len();
1419
1420    let (id, _) = ddl_manager
1421        .submit_alter_logical_table_tasks(alter_table_tasks, physical_table_id, procedure_context)
1422        .await?;
1423
1424    info!(
1425        "{num_logical_tables} logical tables on physical table: {physical_table_id:?} is altered via procedure_id {id:?}"
1426    );
1427
1428    let procedure_id = id.to_string();
1429
1430    Ok(SubmitDdlTaskResponse {
1431        key: procedure_id.into(),
1432        ..Default::default()
1433    })
1434}
1435
1436/// Handle the `[CreateViewTask]` and returns the DDL response when success.
1437async fn handle_create_view_task(
1438    ddl_manager: &DdlManager,
1439    create_view_task: CreateViewTask,
1440    procedure_context: ProcedureContext,
1441) -> Result<SubmitDdlTaskResponse> {
1442    let (id, output) = ddl_manager
1443        .submit_create_view_task(create_view_task, procedure_context)
1444        .await?;
1445
1446    let procedure_id = id.to_string();
1447    let output = output.context(ProcedureOutputSnafu {
1448        procedure_id: &procedure_id,
1449        err_msg: "empty output",
1450    })?;
1451    let view_id = *(output.downcast_ref::<u32>().context(ProcedureOutputSnafu {
1452        procedure_id: &procedure_id,
1453        err_msg: "downcast to `u32`",
1454    })?);
1455    info!("View: {view_id} is created via procedure_id {id:?}");
1456
1457    Ok(SubmitDdlTaskResponse {
1458        key: procedure_id.into(),
1459        table_ids: vec![view_id],
1460    })
1461}
1462
1463async fn handle_comment_on_task(
1464    ddl_manager: &DdlManager,
1465    comment_on_task: CommentOnTask,
1466    procedure_context: ProcedureContext,
1467) -> Result<SubmitDdlTaskResponse> {
1468    let (id, _) = ddl_manager
1469        .submit_comment_on_task(comment_on_task.clone(), procedure_context)
1470        .await?;
1471
1472    let procedure_id = id.to_string();
1473    info!(
1474        "Comment on {}.{}.{} is updated via procedure_id {id:?}",
1475        comment_on_task.catalog_name, comment_on_task.schema_name, comment_on_task.object_name
1476    );
1477
1478    Ok(SubmitDdlTaskResponse {
1479        key: procedure_id.into(),
1480        ..Default::default()
1481    })
1482}
1483
1484#[cfg(test)]
1485mod tests {
1486    use std::sync::Arc;
1487    #[cfg(feature = "enterprise")]
1488    use std::sync::Mutex;
1489    use std::time::Duration;
1490
1491    #[cfg(feature = "enterprise")]
1492    use common_base::protocol::Channel;
1493    use common_error::ext::BoxedError;
1494    #[cfg(feature = "enterprise")]
1495    use common_error::ext::ErrorExt;
1496    #[cfg(feature = "enterprise")]
1497    use common_error::status_code::StatusCode;
1498    #[cfg(feature = "enterprise")]
1499    use common_event_recorder::{PersistentEventContext, ProcedureEventInput, TriggerReason};
1500    use common_procedure::local::LocalManager;
1501    use common_procedure::test_util::InMemoryPoisonStore;
1502    use common_procedure::{BoxedProcedure, ProcedureContext, ProcedureManagerRef};
1503    use store_api::storage::TableId;
1504    use table::table_name::TableName;
1505
1506    use super::DdlManager;
1507    use crate::cache_invalidator::DummyCacheInvalidator;
1508    use crate::ddl::alter_table::AlterTableProcedure;
1509    use crate::ddl::create_database::{
1510        AtomicCreateOutcome, CreateDatabaseMetadataCommitter, CreateDatabaseProcedure,
1511    };
1512    use crate::ddl::create_table::CreateTableProcedure;
1513    use crate::ddl::drop_table::DropTableProcedure;
1514    use crate::ddl::flow_meta::FlowMetadataAllocator;
1515    use crate::ddl::table_meta::TableMetadataAllocator;
1516    use crate::ddl::truncate_table::TruncateTableProcedure;
1517    use crate::ddl::{DdlContext, NoopRegionFailureDetectorControl};
1518    use crate::ddl_manager::{RepartitionProcedureFactory, RepartitionSource};
1519    use crate::key::TableMetadataManager;
1520    use crate::key::flow::FlowMetadataManager;
1521    use crate::kv_backend::memory::MemoryKvBackend;
1522    use crate::node_manager::{DatanodeManager, DatanodeRef, FlownodeManager, FlownodeRef};
1523    use crate::peer::Peer;
1524    use crate::procedure_executor::ExecutorContext;
1525    use crate::region_keeper::MemoryRegionKeeper;
1526    use crate::region_registry::LeaderRegionRegistry;
1527    #[cfg(feature = "enterprise")]
1528    use crate::rpc::ddl::trigger::{CreateTriggerTask, DropTriggerTask};
1529    use crate::rpc::ddl::{CreatorGrantIntent, UndropTableTask};
1530    #[cfg(not(feature = "enterprise"))]
1531    use crate::rpc::ddl::{DdlTask, PurgeDroppedTableTask, QueryContext, SubmitDdlTaskRequest};
1532    #[cfg(feature = "enterprise")]
1533    use crate::rpc::ddl::{DdlTask, QueryContext, SubmitDdlTaskRequest};
1534    use crate::sequence::SequenceBuilder;
1535    use crate::state_store::KvStateStore;
1536    use crate::test_util::{MockDatanodeManager, new_ddl_context};
1537    use crate::wal_provider::WalProvider;
1538
1539    /// A dummy implemented [NodeManager].
1540    pub struct DummyDatanodeManager;
1541
1542    #[async_trait::async_trait]
1543    impl DatanodeManager for DummyDatanodeManager {
1544        async fn datanode(&self, _datanode: &Peer) -> DatanodeRef {
1545            unimplemented!()
1546        }
1547    }
1548
1549    #[async_trait::async_trait]
1550    impl FlownodeManager for DummyDatanodeManager {
1551        async fn flownode(&self, _node: &Peer) -> FlownodeRef {
1552            unimplemented!()
1553        }
1554    }
1555
1556    struct DummyRepartitionProcedureFactory;
1557
1558    #[async_trait::async_trait]
1559    impl RepartitionProcedureFactory for DummyRepartitionProcedureFactory {
1560        fn create(
1561            &self,
1562            _ddl_ctx: &DdlContext,
1563            _table_name: TableName,
1564            _table_id: TableId,
1565            _source: RepartitionSource,
1566            _to_exprs: Vec<String>,
1567            _timeout: Option<Duration>,
1568        ) -> std::result::Result<BoxedProcedure, BoxedError> {
1569            unimplemented!()
1570        }
1571
1572        fn register_loaders(
1573            &self,
1574            _ddl_ctx: &DdlContext,
1575            _procedure_manager: &ProcedureManagerRef,
1576        ) -> std::result::Result<(), BoxedError> {
1577            Ok(())
1578        }
1579
1580        async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError> {
1581            Ok(())
1582        }
1583    }
1584
1585    struct LoaderCommitter;
1586
1587    #[async_trait::async_trait]
1588    impl CreateDatabaseMetadataCommitter for LoaderCommitter {
1589        async fn commit(
1590            &self,
1591            _catalog: &str,
1592            _schema: &str,
1593            _value: &crate::key::schema_name::SchemaNameValue,
1594            _creator: &CreatorGrantIntent,
1595        ) -> std::result::Result<AtomicCreateOutcome, BoxedError> {
1596            unreachable!()
1597        }
1598    }
1599
1600    #[cfg(feature = "enterprise")]
1601    #[derive(Default)]
1602    struct RecordingTriggerDdlManager {
1603        procedure_contexts: Mutex<Vec<ProcedureContext>>,
1604    }
1605
1606    #[cfg(feature = "enterprise")]
1607    #[async_trait::async_trait]
1608    impl super::TriggerDdlManager for RecordingTriggerDdlManager {
1609        async fn create_trigger(
1610            &self,
1611            _create_trigger_task: CreateTriggerTask,
1612            _procedure_manager: ProcedureManagerRef,
1613            _ddl_context: DdlContext,
1614            _query_context: QueryContext,
1615            procedure_context: ProcedureContext,
1616        ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1617            self.procedure_contexts
1618                .lock()
1619                .unwrap()
1620                .push(procedure_context);
1621            Ok(Default::default())
1622        }
1623
1624        async fn drop_trigger(
1625            &self,
1626            _drop_trigger_task: DropTriggerTask,
1627            _procedure_manager: ProcedureManagerRef,
1628            _ddl_context: DdlContext,
1629            _query_context: QueryContext,
1630            procedure_context: ProcedureContext,
1631        ) -> crate::error::Result<crate::rpc::ddl::SubmitDdlTaskResponse> {
1632            self.procedure_contexts
1633                .lock()
1634                .unwrap()
1635                .push(procedure_context);
1636            Ok(Default::default())
1637        }
1638
1639        fn as_any(&self) -> &dyn std::any::Any {
1640            self
1641        }
1642    }
1643
1644    #[test]
1645    fn test_generic_loader_captures_configured_committer() {
1646        let mut context = new_ddl_context(Arc::new(MockDatanodeManager::new(())));
1647        let committer = Arc::new(LoaderCommitter);
1648        context.create_database_metadata_committer = Some(committer.clone());
1649        let loader = procedure_loader!(CreateDatabaseProcedure)[0].1(context);
1650        assert_eq!(Arc::strong_count(&committer), 2);
1651
1652        let loaded = loader(
1653            r#"{
1654                "state":"CreateMetadata",
1655                "catalog":"greptime",
1656                "schema":"metrics",
1657                "create_if_not_exists":false,
1658                "options":{},
1659                "creator":{"username":"alice","created_at_ns":42}
1660            }"#,
1661        )
1662        .unwrap();
1663
1664        assert_eq!(loaded.type_name(), "metasrv-procedure::CreateDatabase");
1665        assert_eq!(Arc::strong_count(&committer), 3);
1666        drop(loaded);
1667        assert_eq!(Arc::strong_count(&committer), 2);
1668    }
1669
1670    #[test]
1671    fn test_register_loaders() {
1672        let kv_backend = Arc::new(MemoryKvBackend::new());
1673        let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1674        let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1675            Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1676            Arc::new(WalProvider::default()),
1677        ));
1678        let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1679        let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1680            Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1681        ));
1682
1683        let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1684        let poison_manager = Arc::new(InMemoryPoisonStore::default());
1685        let procedure_manager = Arc::new(LocalManager::new(
1686            Default::default(),
1687            state_store,
1688            poison_manager,
1689            None,
1690            None,
1691        ));
1692
1693        let ddl_manager = DdlManager::new(
1694            DdlContext {
1695                node_manager: Arc::new(DummyDatanodeManager),
1696                cache_invalidator: Arc::new(DummyCacheInvalidator),
1697                table_metadata_manager,
1698                table_metadata_allocator,
1699                flow_metadata_manager,
1700                flow_metadata_allocator,
1701                memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1702                leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1703                region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1704                soft_drop_enabled: false,
1705                soft_drop_retention: None,
1706                create_database_metadata_committer: None,
1707            },
1708            procedure_manager.clone(),
1709            Arc::new(DummyRepartitionProcedureFactory),
1710        );
1711        ddl_manager.register_loaders().unwrap();
1712
1713        let expected_loaders = vec![
1714            CreateTableProcedure::TYPE_NAME,
1715            AlterTableProcedure::TYPE_NAME,
1716            DropTableProcedure::TYPE_NAME,
1717            TruncateTableProcedure::TYPE_NAME,
1718        ];
1719
1720        for loader in expected_loaders {
1721            assert!(procedure_manager.contains_loader(loader));
1722        }
1723
1724        let soft_drop_loaders = [
1725            "metasrv-procedure::UndropTable",
1726            "metasrv-procedure::PurgeDroppedTable",
1727            "metasrv-procedure::PurgeExpiredDroppedTable",
1728        ];
1729        for loader in soft_drop_loaders {
1730            assert_eq!(
1731                cfg!(feature = "enterprise"),
1732                procedure_manager.contains_loader(loader)
1733            );
1734        }
1735    }
1736
1737    fn build_soft_drop_test_ddl_manager() -> DdlManager {
1738        let kv_backend = Arc::new(MemoryKvBackend::new());
1739        let table_metadata_manager = Arc::new(TableMetadataManager::new(kv_backend.clone()));
1740        let table_metadata_allocator = Arc::new(TableMetadataAllocator::new(
1741            Arc::new(SequenceBuilder::new("test", kv_backend.clone()).build()),
1742            Arc::new(WalProvider::default()),
1743        ));
1744        let flow_metadata_manager = Arc::new(FlowMetadataManager::new(kv_backend.clone()));
1745        let flow_metadata_allocator = Arc::new(FlowMetadataAllocator::with_noop_peer_allocator(
1746            Arc::new(SequenceBuilder::new("flow-test", kv_backend.clone()).build()),
1747        ));
1748
1749        let state_store = Arc::new(KvStateStore::new(kv_backend.clone()));
1750        let poison_manager = Arc::new(InMemoryPoisonStore::default());
1751        let procedure_manager = Arc::new(LocalManager::new(
1752            Default::default(),
1753            state_store,
1754            poison_manager,
1755            None,
1756            None,
1757        ));
1758
1759        DdlManager::new(
1760            DdlContext {
1761                node_manager: Arc::new(DummyDatanodeManager),
1762                cache_invalidator: Arc::new(DummyCacheInvalidator),
1763                table_metadata_manager,
1764                table_metadata_allocator,
1765                flow_metadata_manager,
1766                flow_metadata_allocator,
1767                memory_region_keeper: Arc::new(MemoryRegionKeeper::default()),
1768                leader_region_registry: Arc::new(LeaderRegionRegistry::default()),
1769                region_failure_detector_controller: Arc::new(NoopRegionFailureDetectorControl),
1770                soft_drop_enabled: true,
1771                soft_drop_retention: Some(std::time::Duration::from_secs(1)),
1772                create_database_metadata_committer: None,
1773            },
1774            procedure_manager,
1775            Arc::new(DummyRepartitionProcedureFactory),
1776        )
1777    }
1778
1779    #[cfg(feature = "enterprise")]
1780    #[tokio::test]
1781    async fn test_trigger_ddl_forwards_procedure_context() {
1782        let trigger_ddl_manager = Arc::new(RecordingTriggerDdlManager::default());
1783        let ddl_manager = build_soft_drop_test_ddl_manager()
1784            .with_trigger_ddl_manager(trigger_ddl_manager.clone());
1785        let procedure_context = ProcedureContext {
1786            actor: Some("test-user".to_string()),
1787            event_context: Some(
1788                PersistentEventContext::new(TriggerReason::Manual).with_protocol("mysql"),
1789            ),
1790        };
1791        let executor_context = || ExecutorContext {
1792            query_context: Some(QueryContext {
1793                channel: Channel::Mysql as u8,
1794                ..Default::default()
1795            }),
1796            actor: Some("test-user".to_string()),
1797            event_input: Some(ProcedureEventInput::new(TriggerReason::Manual)),
1798            ..Default::default()
1799        };
1800
1801        ddl_manager
1802            .submit_ddl_task(
1803                executor_context(),
1804                SubmitDdlTaskRequest::new(DdlTask::CreateTrigger(CreateTriggerTask {
1805                    catalog_name: "greptime".to_string(),
1806                    trigger_name: "test_trigger".to_string(),
1807                    if_not_exists: false,
1808                    sql: "SELECT 1".to_string(),
1809                    channels: vec![],
1810                    labels: Default::default(),
1811                    annotations: Default::default(),
1812                    interval: Duration::from_secs(1),
1813                    raw_interval_expr: None,
1814                    r#for: None,
1815                    for_raw_expr: None,
1816                    keep_firing_for: None,
1817                    keep_firing_for_raw_expr: None,
1818                })),
1819            )
1820            .await
1821            .unwrap();
1822        ddl_manager
1823            .submit_ddl_task(
1824                executor_context(),
1825                SubmitDdlTaskRequest::new(DdlTask::DropTrigger(DropTriggerTask {
1826                    catalog_name: "greptime".to_string(),
1827                    trigger_name: "test_trigger".to_string(),
1828                    drop_if_exists: false,
1829                })),
1830            )
1831            .await
1832            .unwrap();
1833
1834        assert_eq!(
1835            *trigger_ddl_manager.procedure_contexts.lock().unwrap(),
1836            vec![procedure_context.clone(), procedure_context]
1837        );
1838    }
1839
1840    #[cfg(feature = "enterprise")]
1841    #[tokio::test]
1842    async fn test_submit_undrop_missing_tombstone_returns_table_not_found_directly() {
1843        let ddl_manager = build_soft_drop_test_ddl_manager();
1844
1845        let err = ddl_manager
1846            .submit_undrop_table_task(
1847                UndropTableTask { table_id: 1024 },
1848                ProcedureContext::default(),
1849            )
1850            .await
1851            .unwrap_err();
1852
1853        assert_eq!(err.status_code(), StatusCode::TableNotFound);
1854        assert!(matches!(err, crate::error::Error::TableNotFound { .. }));
1855    }
1856
1857    #[cfg(not(feature = "enterprise"))]
1858    #[tokio::test]
1859    async fn test_submit_undrop_and_purge_rejected_in_non_enterprise_build() {
1860        let ddl_manager = build_soft_drop_test_ddl_manager();
1861
1862        let err = ddl_manager
1863            .submit_undrop_table_task(
1864                UndropTableTask { table_id: 1024 },
1865                ProcedureContext::default(),
1866            )
1867            .await
1868            .unwrap_err();
1869        assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1870
1871        let err = ddl_manager
1872            .submit_purge_dropped_table_task(
1873                PurgeDroppedTableTask { table_id: 1024 },
1874                ProcedureContext::default(),
1875            )
1876            .await
1877            .unwrap_err();
1878        assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1879
1880        let err = ddl_manager
1881            .submit_expired_purge_dropped_table_task(PurgeDroppedTableTask { table_id: 1024 })
1882            .await
1883            .unwrap_err();
1884        assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1885
1886        for task in [
1887            DdlTask::UndropTable(UndropTableTask { table_id: 1024 }),
1888            DdlTask::PurgeDroppedTable(PurgeDroppedTableTask { table_id: 1024 }),
1889        ] {
1890            let err = ddl_manager
1891                .submit_ddl_task(
1892                    ExecutorContext {
1893                        query_context: Some(QueryContext::default()),
1894                        ..Default::default()
1895                    },
1896                    SubmitDdlTaskRequest::new(task),
1897                )
1898                .await
1899                .unwrap_err();
1900            assert!(matches!(err, crate::error::Error::Unsupported { .. }));
1901        }
1902    }
1903}