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_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/// A configurator that customizes or enhances a [`DdlManager`].
92#[async_trait::async_trait]
93pub trait DdlManagerConfigurator<C>: Send + Sync {
94    /// Configures the given [`DdlManager`] using the provided [`DdlManagerConfigureContext`].
95    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/// The [DdlManager] provides the ability to execute Ddl.
109#[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/// This trait is responsible for handling DDL tasks about triggers. e.g.,
123/// create trigger, drop trigger, etc.
124#[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/// This trait is responsible for handling custom/extension CREATE FLOW tasks.
152#[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/// Hook for classifying and handling DROP FLOW requests.
169#[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        /// Full target partition columns to overwrite table metadata.
212        ///
213        /// `None` means the repartition keeps using the current table
214        /// partition columns, so the procedure won't update
215        /// `partition_key_indices`.
216        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    /// Persists the cluster-level GC requirement before submitting a repartition.
243    async fn ensure_gc_requirement(&self) -> std::result::Result<(), BoxedError>;
244}
245
246/// The options for DDL tasks.
247///
248/// Note: These options may not be utilized by all procedures.
249/// At present, they are specifically applied in `RepartitionProcedure`.
250#[derive(Debug, Clone, Copy)]
251pub struct DdlOptions {
252    /// The timeout will be passed to the procedure.
253    ///
254    /// Note: Each procedure may implement its own timeout handling mechanism.
255    pub timeout: Duration,
256    /// The flag that controls whether to wait for the procedure to complete.
257    ///
258    /// If wait is `true`, the procedure will wait for completion(success or failure) and the result will be returned.
259    /// Otherwise, the procedure will be submitted and return the [ProcedureId](common_procedure::ProcedureId) immediately.
260    ///
261    /// 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.
262    pub wait: bool,
263}
264
265impl DdlManager {
266    /// Returns a new [DdlManager].
267    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    /// Returns the [TableMetadataManagerRef].
312    pub fn table_metadata_manager(&self) -> &TableMetadataManagerRef {
313        &self.ddl_context.table_metadata_manager
314    }
315
316    /// Returns the [DdlContext]
317    pub fn create_context(&self) -> DdlContext {
318        self.ddl_context.clone()
319    }
320
321    /// Registers all Ddl loaders.
322    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    /// Submits a repartition procedure for the specified table.
379    ///
380    /// This creates a repartition procedure using the provided `table_id`,
381    /// `table_name`, and `Repartition` configuration, and then either executes it
382    /// to completion or just submits it for asynchronous execution.
383    ///
384    /// The `Repartition` argument contains the original (`from_partition_exprs`)
385    /// and target (`into_partition_exprs`) partition expressions that define how
386    /// the table should be repartitioned.
387    ///
388    /// The `wait` flag controls whether this method waits for the repartition
389    /// procedure to finish:
390    /// - If `wait` is `true`, the procedure is executed and this method awaits
391    ///   its completion, returning both the generated `ProcedureId` and the
392    ///   final `Output` of the procedure.
393    /// - If `wait` is `false`, the procedure is only submitted to the procedure
394    ///   manager for asynchronous execution, and this method returns the
395    ///   `ProcedureId` along with `None` as the output.
396    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                // Reads the deprecated field for backward compatibility with old persisted DDL tasks.
422                #[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    /// Submits and executes an alter table task.
457    #[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        // make alter_table_task mutable so we can call .take() on its field
466        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            // Hold the same physical region locks as migration while validating that
497            // this route snapshot is still the one the procedure will mutate.
498            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    /// Submits and executes a create table task.
546    #[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    /// Submits and executes a `[CreateViewTask]`.
568    #[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    /// Submits and executes a create multiple logical table tasks.
585    #[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    /// Submits and executes alter multiple table tasks.
604    #[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        // Resolve the logical table ids up front: procedure locks are fixed
614        // at submission, so `lock_key` cannot derive them during `Prepare`.
615        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    /// Submits and executes a drop table task.
641    #[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    /// Submits and executes an undrop table task.
658    #[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    /// Submits and executes a purge dropped table task.
698    #[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    /// Submits and executes a purge task that first rechecks the tombstone deadline.
726    #[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    /// Submits and executes a create database task.
758    #[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    /// Submits and executes a drop table task.
786    #[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    /// Submits and executes a create flow task.
818    #[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    /// Submits and executes a drop flow task.
834    #[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    /// Submits and executes a drop view task.
849    #[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    /// Submits and executes a truncate table task.
864    #[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    /// Submits and executes a comment on task.
881    #[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    /// Executes a procedure and waits for the result.
902    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    /// Submits a procedure and returns the procedure id.
922    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    // Classify before the route guard: a mixed annotation batch must surface
1142    // its own error here, not a misleading "non-physical route" one. Families
1143    // that only rewrite the logical table's metadata may target logical tables.
1144    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    // Use the physical table id in the first logical table, then it will be checked in the procedure.
1531    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
1556/// Handle the `[CreateViewTask]` and returns the DDL response when success.
1557async 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    /// A dummy implemented [NodeManager].
1669    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        // A semantic alter on a logical table passes the route guard, updates
2262        // only the logical table's metadata, and dispatches nothing.
2263        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        // A mixed batch fails with its own error, not the route guard's.
2296        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        // The repartition hint drives physical repartitioning; on a logical
2313        // route it stays rejected by the guard.
2314        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}