Skip to main content

common_meta/
ddl_manager.rs

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