From 47ca5c362ec79f890e683d0fec389efb485d1dfe Mon Sep 17 00:00:00 2001 From: Weny Xu Date: Thu, 30 Jul 2026 11:24:47 +0800 Subject: [PATCH] feat: add table DDL procedure events (#8627) * feat(meta): emit table DDL procedure events Signed-off-by: WenyXu * fix(meta): honor table DDL event filters Signed-off-by: WenyXu * test(meta): cover table DDL event filters Signed-off-by: WenyXu * refactor(meta): align table DDL event conventions Signed-off-by: WenyXu * refactor(meta): bound table DDL event payloads Signed-off-by: WenyXu * test(meta): consolidate table DDL event tests Signed-off-by: WenyXu * style(meta): use crate visibility in event tests Signed-off-by: WenyXu * test: stabilize table DDL event assertions Signed-off-by: WenyXu * fix(meta): exclude repartition from alter table events Signed-off-by: WenyXu * fix(meta): resolve table event rebase conflicts Signed-off-by: WenyXu --------- Signed-off-by: WenyXu --- config/config.md | 4 +- config/metasrv.example.toml | 4 +- config/standalone.example.toml | 2 + src/common/event-recorder/src/event_table.rs | 28 + .../meta/src/ddl/alter_logical_tables.rs | 36 +- src/common/meta/src/ddl/alter_table.rs | 35 +- .../meta/src/ddl/create_logical_tables.rs | 60 +- src/common/meta/src/ddl/create_table.rs | 41 +- src/common/meta/src/ddl/drop_table.rs | 26 +- src/common/meta/src/ddl/event.rs | 1 + src/common/meta/src/ddl/event/table.rs | 426 ++++++++++++++ .../meta/src/ddl/purge_dropped_table.rs | 22 +- .../src/ddl/tests/alter_logical_tables.rs | 2 +- src/common/meta/src/ddl/tests/alter_table.rs | 2 +- src/common/meta/src/ddl/tests/event.rs | 1 + .../meta/src/ddl/tests/event/database.rs | 45 +- src/common/meta/src/ddl/tests/event/flow.rs | 42 +- src/common/meta/src/ddl/tests/event/table.rs | 552 ++++++++++++++++++ .../meta/src/ddl/tests/event/test_util.rs | 2 +- src/common/meta/src/ddl/truncate_table.rs | 24 +- src/common/meta/src/ddl/undrop_table.rs | 34 +- .../tests/event_recorder_test_util.rs | 28 + tests-integration/tests/main.rs | 1 + tests-integration/tests/table_ddl_event.rs | 452 ++++++++++++++ 24 files changed, 1779 insertions(+), 91 deletions(-) create mode 100644 src/common/meta/src/ddl/event/table.rs create mode 100644 src/common/meta/src/ddl/tests/event/table.rs create mode 100644 tests-integration/tests/table_ddl_event.rs diff --git a/config/config.md b/config/config.md index 6f75398436..2e93bd37a9 100644 --- a/config/config.md +++ b/config/config.md @@ -231,7 +231,7 @@ | `tracing.tokio_console_addr` | String | Unset | The tokio console address. | | `event_recorder` | -- | -- | Configuration options for the event recorder. | | `event_recorder.ttl` | String | `90d` | TTL for the events table that will be used to store the events. Default is `90d`. | -| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `create_database`,
`alter_database`, `drop_database`, `create_flow`, `drop_flow`,
`create_view`, `drop_view`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | +| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `create_database`,
`alter_database`, `drop_database`, `create_flow`, `drop_flow`,
`create_table`, `create_logical_tables`, `alter_table`, `alter_logical_tables`,
`drop_table`, `undrop_table`, `purge_dropped_table`, `truncate_table`,
`create_view`, `drop_view`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | | `memory` | -- | -- | The memory options. | | `memory.enable_heap_profiling` | Bool | `true` | Whether to enable heap profiling activation during startup.
When enabled, heap profiling will be activated if the `MALLOC_CONF` environment variable
is set to "prof:true,prof_active:false". The official image adds this env variable.
Default is true. | @@ -444,7 +444,7 @@ | `wal.create_topic_timeout` | String | `30s` | The timeout for creating a Kafka topic.
**It's only used when the provider is `kafka`**. | | `event_recorder` | -- | -- | Configuration options for the event recorder. | | `event_recorder.ttl` | String | `90d` | TTL for the events table that will be used to store the events. Default is `90d`. | -| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,
`create_database`, `alter_database`, `drop_database`, `create_flow`,
`drop_flow`, `create_view`, `drop_view`, `repartition`,
`repartition_group`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | +| `event_recorder.event_types` | Array | -- | Event types to record. Current available event types: `region_migration`,
`create_database`, `alter_database`, `drop_database`, `create_flow`,
`drop_flow`, `create_table`, `create_logical_tables`, `alter_table`,
`alter_logical_tables`, `drop_table`, `undrop_table`, `purge_dropped_table`,
`truncate_table`, `create_view`, `drop_view`, `repartition`,
`repartition_group`.
When omitted, all current and future event types are recorded.
Set to an empty array to disable event recording. | | `stats_persistence` | -- | -- | Configuration options for the stats persistence. | | `stats_persistence.ttl` | String | `0s` | TTL for the stats table that will be used to store the stats.
Set to `0s` to disable stats persistence.
Default is `0s`.
If you want to enable stats persistence, set the TTL to a value greater than 0.
It is recommended to set a small value, e.g., `3h`. | | `stats_persistence.interval` | String | `10m` | The interval to persist the stats. Default is `10m`.
The minimum value is `10m`, if the value is less than `10m`, it will be overridden to `10m`. | diff --git a/config/metasrv.example.toml b/config/metasrv.example.toml index c9c26e0010..37995c22ea 100644 --- a/config/metasrv.example.toml +++ b/config/metasrv.example.toml @@ -314,7 +314,9 @@ create_topic_timeout = "30s" ttl = "90d" ## Event types to record. Current available event types: `region_migration`, ## `create_database`, `alter_database`, `drop_database`, `create_flow`, -## `drop_flow`, `create_view`, `drop_view`, `repartition`, +## `drop_flow`, `create_table`, `create_logical_tables`, `alter_table`, +## `alter_logical_tables`, `drop_table`, `undrop_table`, `purge_dropped_table`, +## `truncate_table`, `create_view`, `drop_view`, `repartition`, ## `repartition_group`. ## When omitted, all current and future event types are recorded. ## Set to an empty array to disable event recording. diff --git a/config/standalone.example.toml b/config/standalone.example.toml index 5b353df5a7..371b64864c 100644 --- a/config/standalone.example.toml +++ b/config/standalone.example.toml @@ -905,6 +905,8 @@ default_ratio = 1.0 ttl = "90d" ## Event types to record. Current available event types: `create_database`, ## `alter_database`, `drop_database`, `create_flow`, `drop_flow`, +## `create_table`, `create_logical_tables`, `alter_table`, `alter_logical_tables`, +## `drop_table`, `undrop_table`, `purge_dropped_table`, `truncate_table`, ## `create_view`, `drop_view`. ## When omitted, all current and future event types are recorded. ## Set to an empty array to disable event recording. diff --git a/src/common/event-recorder/src/event_table.rs b/src/common/event-recorder/src/event_table.rs index 54ca232367..6e524d933b 100644 --- a/src/common/event-recorder/src/event_table.rs +++ b/src/common/event-recorder/src/event_table.rs @@ -131,6 +131,12 @@ pub const TABLE_NAME_COLUMN: EventTableColumn = /// The canonical table identifier field for table DDL events. pub const TABLE_ID_COLUMN: EventTableColumn = EventTableColumn::new("table_id", ColumnDataType::Uint32, SemanticType::Field); +/// The canonical physical table identifier dimension. +pub const PHYSICAL_TABLE_ID_COLUMN: EventTableColumn = EventTableColumn::new( + "physical_table_id", + ColumnDataType::Uint32, + SemanticType::Field, +); /// The canonical region identifier field for region events. pub const REGION_ID_COLUMN: EventTableColumn = EventTableColumn::new("region_id", ColumnDataType::Uint64, SemanticType::Field); @@ -337,6 +343,28 @@ mod tests { ); } + #[test] + fn table_dimension_schema_preserves_names_types_semantics_and_order() { + assert_eq!( + column_schemas([ + &TABLE_NAME_COLUMN, + &TABLE_ID_COLUMN, + &PHYSICAL_TABLE_ID_COLUMN, + ]), + [ + ("table_name", ColumnDataType::String), + ("table_id", ColumnDataType::Uint32), + ("physical_table_id", ColumnDataType::Uint32), + ] + .map(|(column_name, datatype)| ColumnSchema { + column_name: column_name.to_string(), + datatype: datatype.into(), + semantic_type: SemanticType::Field.into(), + ..Default::default() + }) + ); + } + #[test] fn view_dimension_schema_preserves_names_types_semantics_and_order() { assert_eq!( diff --git a/src/common/meta/src/ddl/alter_logical_tables.rs b/src/common/meta/src/ddl/alter_logical_tables.rs index ac41358a51..683248965f 100644 --- a/src/common/meta/src/ddl/alter_logical_tables.rs +++ b/src/common/meta/src/ddl/alter_logical_tables.rs @@ -20,7 +20,7 @@ use api::region::RegionResponse; use async_trait::async_trait; use common_catalog::format_full_table_name; use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu}; -use common_procedure::{Context, LockKey, Procedure, Status}; +use common_procedure::{Context, EventContext, EventTrigger, LockKey, Procedure, Status}; use common_telemetry::{debug, error, info, warn}; pub use executor::make_alter_region_request; use serde::{Deserialize, Serialize}; @@ -36,6 +36,9 @@ use crate::ddl::alter_logical_tables::executor::AlterLogicalTablesExecutor; use crate::ddl::alter_logical_tables::validator::{ AlterLogicalTableValidator, ValidatorResult, retain_unskipped, }; +use crate::ddl::event::table::{ + TableDdlEvent, TableDdlEventType, TableDdlLocator, alter_table_kind_name, +}; use crate::ddl::utils::{extract_column_metadatas, map_to_procedure_error, sync_follower_regions}; use crate::error::Result; use crate::instruction::CacheIdent; @@ -316,6 +319,37 @@ impl Procedure for AlterLogicalTablesProcedure { LockKey::new(lock_key) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::AlterLogicalTables.as_str()) + { + return None; + } + if ctx.trigger != EventTrigger::Submitted { + return Some(Box::new(TableDdlEvent::lifecycle( + TableDdlEventType::AlterLogicalTables, + ))); + } + + let locators = self.data.tasks.iter().map(|task| { + let table_ref = task.table_ref(); + TableDdlLocator::new(table_ref.catalog, table_ref.schema, table_ref.table) + .with_physical_table_id(self.data.physical_table_id) + }); + let kinds = self + .data + .tasks + .iter() + .filter_map(|task| task.alter_table.kind.as_ref()) + .filter_map(alter_table_kind_name); + Some(Box::new(TableDdlEvent::alter_logical_tables_submitted( + locators, + self.data.tasks.len(), + kinds, + ))) + } } #[derive(Debug, Serialize, Deserialize)] diff --git a/src/common/meta/src/ddl/alter_table.rs b/src/common/meta/src/ddl/alter_table.rs index cf2dfc7fbc..1516772e51 100644 --- a/src/common/meta/src/ddl/alter_table.rs +++ b/src/common/meta/src/ddl/alter_table.rs @@ -25,8 +25,8 @@ use async_trait::async_trait; use common_error::ext::BoxedError; use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu}; use common_procedure::{ - Context as ProcedureContext, ContextProvider, Error as ProcedureError, LockKey, PoisonKey, - PoisonKeys, Procedure, ProcedureId, Status, StringKey, + Context as ProcedureContext, ContextProvider, Error as ProcedureError, EventContext, + EventTrigger, LockKey, PoisonKey, PoisonKeys, Procedure, ProcedureId, Status, StringKey, }; use common_telemetry::{error, info, warn}; use serde::{Deserialize, Serialize}; @@ -40,6 +40,9 @@ use table::table_reference::TableReference; use crate::ddl::DdlContext; use crate::ddl::alter_table::executor::AlterTableExecutor; +use crate::ddl::event::table::{ + TableDdlEvent, TableDdlEventType, TableDdlLocator, alter_table_kind_name, +}; use crate::ddl::utils::{ MultipleResults, extract_column_metadatas, handle_multiple_results, map_to_procedure_error, sync_follower_regions, @@ -394,6 +397,34 @@ impl Procedure for AlterTableProcedure { fn poison_keys(&self) -> PoisonKeys { PoisonKeys::new(vec![self.table_poison_key()]) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::AlterTable.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let table_ref = self.data.table_ref(); + let locator = + TableDdlLocator::new(table_ref.catalog, table_ref.schema, table_ref.table) + .with_table_id(self.data.table_id()); + let kind = self + .data + .task + .alter_table + .kind + .as_ref() + .and_then(alter_table_kind_name); + TableDdlEvent::alter_table_submitted(locator, kind) + } + _ => TableDdlEvent::lifecycle(TableDdlEventType::AlterTable), + }; + + Some(Box::new(event)) + } } #[derive(Debug, Serialize, Deserialize, AsRefStr)] diff --git a/src/common/meta/src/ddl/create_logical_tables.rs b/src/common/meta/src/ddl/create_logical_tables.rs index 78f6e225b7..08f6cb855c 100644 --- a/src/common/meta/src/ddl/create_logical_tables.rs +++ b/src/common/meta/src/ddl/create_logical_tables.rs @@ -21,8 +21,12 @@ use api::region::RegionResponse; use api::v1::CreateTableExpr; use async_trait::async_trait; use common_catalog::consts::METRIC_ENGINE; +use common_event_recorder::Event; use common_procedure::error::{FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu}; -use common_procedure::{Context as ProcedureContext, LockKey, Procedure, Status}; +use common_procedure::{ + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, ProcedureState, + Status, +}; use common_telemetry::{debug, error, warn}; use futures::future; pub use region_request::create_region_request_builder; @@ -35,6 +39,7 @@ use strum::AsRefStr; use table::metadata::{TableId, TableInfo}; use crate::ddl::DdlContext; +use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator}; use crate::ddl::utils::{ add_peer_context_if_needed, extract_column_metadatas, map_to_procedure_error, sync_follower_regions, @@ -252,6 +257,59 @@ impl Procedure for CreateLogicalTablesProcedure { } LockKey::new(lock_key) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::CreateLogicalTables.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let locators = self.data.tasks.iter().map(|task| { + TableDdlLocator::new( + &task.create_table.catalog_name, + &task.create_table.schema_name, + &task.create_table.table_name, + ) + .with_physical_table_id(self.data.physical_table_id) + }); + TableDdlEvent::create_logical_tables_submitted(locators, self.data.tasks.len()) + } + EventTrigger::Succeeded => match ctx.lifecycle_state { + ProcedureState::Done { + output: Some(output), + } => output + .downcast_ref::>() + .map(|table_ids| { + debug_assert_eq!(self.data.tasks.len(), table_ids.len()); + let locators = + self.data + .tasks + .iter() + .zip(table_ids) + .map(|(task, table_id)| { + TableDdlLocator::new( + &task.create_table.catalog_name, + &task.create_table.schema_name, + &task.create_table.table_name, + ) + .with_table_id(*table_id) + .with_physical_table_id(self.data.physical_table_id) + }); + TableDdlEvent::create_logical_tables_succeeded(locators) + }) + .unwrap_or_else(|| { + TableDdlEvent::lifecycle(TableDdlEventType::CreateLogicalTables) + }), + _ => TableDdlEvent::lifecycle(TableDdlEventType::CreateLogicalTables), + }, + _ => TableDdlEvent::lifecycle(TableDdlEventType::CreateLogicalTables), + }; + + Some(Box::new(event)) + } } #[derive(Debug, Serialize, Deserialize)] diff --git a/src/common/meta/src/ddl/create_table.rs b/src/common/meta/src/ddl/create_table.rs index acf52221c2..298fa8e56a 100644 --- a/src/common/meta/src/ddl/create_table.rs +++ b/src/common/meta/src/ddl/create_table.rs @@ -22,7 +22,10 @@ use common_procedure::error::{ ExternalSnafu, FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu, }; use common_procedure::local::DynamicKeyLockGuard; -use common_procedure::{Context as ProcedureContext, LockKey, Procedure, ProcedureId, Status}; +use common_procedure::{ + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, ProcedureId, + ProcedureState, Status, +}; use common_telemetry::info; use serde::{Deserialize, Serialize}; use snafu::{OptionExt, ResultExt}; @@ -35,6 +38,7 @@ pub(crate) use template::{CreateRequestBuilder, build_template_from_raw_table_in use crate::ddl::create_table::executor::CreateTableExecutor; use crate::ddl::create_table::template::build_template; +use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator}; use crate::ddl::utils::map_to_procedure_error; use crate::ddl::{DdlContext, TableMetadata}; use crate::error::{self, Result}; @@ -394,6 +398,41 @@ impl Procedure for CreateTableProcedure { TableNameLock::new(table_ref.catalog, table_ref.schema, table_ref.table).into(), ]) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::CreateTable.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let table_ref = self.data.table_ref(); + let locator = + TableDdlLocator::new(table_ref.catalog, table_ref.schema, table_ref.table); + let create_table = &self.data.task.create_table; + TableDdlEvent::create_table_submitted( + locator, + create_table.create_if_not_exists, + &create_table.engine, + ) + } + EventTrigger::Succeeded => match ctx.lifecycle_state { + ProcedureState::Done { + output: Some(output), + } => output + .downcast_ref::() + .copied() + .map(TableDdlEvent::create_table_succeeded) + .unwrap_or_else(|| TableDdlEvent::lifecycle(TableDdlEventType::CreateTable)), + _ => TableDdlEvent::lifecycle(TableDdlEventType::CreateTable), + }, + _ => TableDdlEvent::lifecycle(TableDdlEventType::CreateTable), + }; + + Some(Box::new(event)) + } } #[derive(Debug, Clone, Serialize, Deserialize, AsRefStr, PartialEq)] diff --git a/src/common/meta/src/ddl/drop_table.rs b/src/common/meta/src/ddl/drop_table.rs index 177edd2553..90e2d50f31 100644 --- a/src/common/meta/src/ddl/drop_table.rs +++ b/src/common/meta/src/ddl/drop_table.rs @@ -19,10 +19,11 @@ use std::collections::HashMap; use async_trait::async_trait; use common_error::ext::BoxedError; +use common_event_recorder::Event; use common_procedure::error::{ExternalSnafu, FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ - Context as ProcedureContext, Error as ProcedureError, LockKey, Procedure, - Result as ProcedureResult, Status, + Context as ProcedureContext, Error as ProcedureError, EventContext, EventTrigger, LockKey, + Procedure, Result as ProcedureResult, Status, }; use common_telemetry::info; use common_telemetry::tracing::warn; @@ -38,6 +39,7 @@ use uuid::Uuid; use self::executor::DropTableExecutor; use crate::ddl::DdlContext; +use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator}; use crate::ddl::utils::{convert_region_routes_to_detecting_regions, map_to_procedure_error}; use crate::error::{self, Result}; use crate::key::table_route::TableRouteValue; @@ -365,6 +367,26 @@ impl Procedure for DropTableProcedure { LockKey::new(lock_key) } + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::DropTable.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let task = &self.data.task; + let locator = TableDdlLocator::new(&task.catalog, &task.schema, &task.table) + .with_table_id(task.table_id); + TableDdlEvent::drop_table_submitted(locator, task.drop_if_exists) + } + _ => TableDdlEvent::lifecycle(TableDdlEventType::DropTable), + }; + + Some(Box::new(event)) + } + fn rollback_supported(&self) -> bool { !matches!(self.data.state, DropTableState::Prepare) && self.data.allow_rollback } diff --git a/src/common/meta/src/ddl/event.rs b/src/common/meta/src/ddl/event.rs index 0e151b2306..bbcff178a2 100644 --- a/src/common/meta/src/ddl/event.rs +++ b/src/common/meta/src/ddl/event.rs @@ -16,4 +16,5 @@ pub(crate) mod database; pub(crate) mod flow; +pub(crate) mod table; pub(crate) mod view; diff --git a/src/common/meta/src/ddl/event/table.rs b/src/common/meta/src/ddl/event/table.rs new file mode 100644 index 0000000000..52508e092a --- /dev/null +++ b/src/common/meta/src/ddl/event/table.rs @@ -0,0 +1,426 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::any::Any; +use std::collections::BTreeSet; + +use api::v1::alter_table_expr::Kind as AlterTableKind; +use api::v1::value::ValueData; +use api::v1::{ColumnSchema, Row}; +use common_event_recorder::Event; +use common_event_recorder::error::{Result, SerializeEventSnafu}; +use common_event_recorder::event_table::{ + CATALOG_NAME_COLUMN, PHYSICAL_TABLE_ID_COLUMN, SCHEMA_NAME_COLUMN, TABLE_ID_COLUMN, + TABLE_NAME_COLUMN, column_schemas, nullable_string, nullable_value, +}; +use serde::Serialize; +use serde_json::Value as JsonValue; +use snafu::ResultExt; +use store_api::storage::TableId; + +/// Current version of table DDL event payloads. +pub(crate) const TABLE_DDL_PAYLOAD_VERSION: u8 = 1; + +/// A table DDL event type and its fixed domain schema. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum TableDdlEventType { + CreateTable, + CreateLogicalTables, + AlterTable, + AlterLogicalTables, + DropTable, + UndropTable, + PurgeDroppedTable, + TruncateTable, +} + +impl TableDdlEventType { + /// Returns the stable event type stored in the events table. + pub(crate) const fn as_str(self) -> &'static str { + match self { + Self::CreateTable => "create_table", + Self::CreateLogicalTables => "create_logical_tables", + Self::AlterTable => "alter_table", + Self::AlterLogicalTables => "alter_logical_tables", + Self::DropTable => "drop_table", + Self::UndropTable => "undrop_table", + Self::PurgeDroppedTable => "purge_dropped_table", + Self::TruncateTable => "truncate_table", + } + } + + const fn has_physical_table_id(self) -> bool { + matches!(self, Self::CreateLogicalTables | Self::AlterLogicalTables) + } +} + +/// Nullable table locator columns stored alongside a table DDL event. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub(crate) struct TableDdlLocator { + /// Catalog containing the table. + pub(crate) catalog_name: Option, + /// Schema containing the table. + pub(crate) schema_name: Option, + /// Table name. + pub(crate) table_name: Option, + /// Table ID when known at this lifecycle point. + pub(crate) table_id: Option, + /// Physical table ID for a logical table event. + pub(crate) physical_table_id: Option, +} + +impl TableDdlLocator { + /// Creates a locator from a fully qualified table name. + pub(crate) fn new( + catalog_name: impl Into, + schema_name: impl Into, + table_name: impl Into, + ) -> Self { + Self { + catalog_name: Some(catalog_name.into()), + schema_name: Some(schema_name.into()), + table_name: Some(table_name.into()), + ..Default::default() + } + } + + /// Creates a locator containing only a table ID. + pub(crate) fn from_table_id(table_id: TableId) -> Self { + Self { + table_id: Some(table_id), + ..Default::default() + } + } + + /// Adds a table ID to the locator. + pub(crate) fn with_table_id(mut self, table_id: TableId) -> Self { + self.table_id = Some(table_id); + self + } + + /// Adds a physical table ID to a logical-table locator. + pub(crate) fn with_physical_table_id(mut self, physical_table_id: TableId) -> Self { + self.physical_table_id = Some(physical_table_id); + self + } +} + +#[derive(Debug, Serialize)] +#[serde(untagged)] +enum TableDdlPayload { + CreateTable(CreateTablePayload), + CreateLogicalTables(CreateLogicalTablesPayload), + AlterTable(AlterTablePayload), + AlterLogicalTables(AlterLogicalTablesPayload), + DropTable(DropTablePayload), + UndropTable(UndropTablePayload), + PurgeDroppedTable(PurgeDroppedTablePayload), + TruncateTable(TruncateTablePayload), +} + +#[derive(Debug, Serialize)] +struct CreateTablePayload { + version: u8, + create_if_not_exists: bool, + engine: String, +} + +#[derive(Debug, Serialize)] +struct CreateLogicalTablesPayload { + version: u8, + table_count: usize, +} + +#[derive(Debug, Serialize)] +struct AlterTablePayload { + version: u8, + kind: Option<&'static str>, +} + +#[derive(Debug, Serialize)] +struct AlterLogicalTablesPayload { + version: u8, + table_count: usize, + kinds: Vec<&'static str>, +} + +#[derive(Debug, Serialize)] +struct DropTablePayload { + version: u8, + drop_if_exists: bool, +} + +#[derive(Debug, Serialize)] +struct UndropTablePayload { + version: u8, +} + +#[derive(Debug, Serialize)] +struct PurgeDroppedTablePayload { + version: u8, +} + +#[derive(Debug, Serialize)] +struct TruncateTablePayload { + version: u8, + time_range_count: usize, +} + +/// Returns the stable kind stored in an Alter Table payload, if supported. +pub(crate) fn alter_table_kind_name(kind: &AlterTableKind) -> Option<&'static str> { + match kind { + AlterTableKind::AddColumns(_) => Some("add_columns"), + AlterTableKind::DropColumns(_) => Some("drop_columns"), + AlterTableKind::RenameTable(_) => Some("rename_table"), + AlterTableKind::ModifyColumnTypes(_) => Some("modify_column_types"), + AlterTableKind::SetTableOptions(_) => Some("set_table_options"), + AlterTableKind::UnsetTableOptions(_) => Some("unset_table_options"), + AlterTableKind::SetIndex(_) => Some("set_index"), + AlterTableKind::UnsetIndex(_) => Some("unset_index"), + AlterTableKind::DropDefaults(_) => Some("drop_defaults"), + AlterTableKind::SetIndexes(_) => Some("set_indexes"), + AlterTableKind::UnsetIndexes(_) => Some("unset_indexes"), + AlterTableKind::SetDefaults(_) => Some("set_defaults"), + // Repartition is handled by RepartitionProcedure. + AlterTableKind::Repartition(_) => None, + } +} + +/// Shared event representation used by table DDL procedures. +#[derive(Debug)] +pub(crate) struct TableDdlEvent { + event_type: TableDdlEventType, + locators: Vec, + payload: Option, +} + +impl TableDdlEvent { + /// Builds the bounded event emitted when creating a table is submitted. + pub(crate) fn create_table_submitted( + locator: TableDdlLocator, + create_if_not_exists: bool, + engine: &str, + ) -> Self { + Self::submitted( + TableDdlEventType::CreateTable, + [locator], + TableDdlPayload::CreateTable(CreateTablePayload { + version: TABLE_DDL_PAYLOAD_VERSION, + create_if_not_exists, + engine: engine.to_string(), + }), + ) + } + + /// Builds the bounded event emitted when creating logical tables is submitted. + pub(crate) fn create_logical_tables_submitted( + locators: impl IntoIterator, + table_count: usize, + ) -> Self { + Self::submitted( + TableDdlEventType::CreateLogicalTables, + locators, + TableDdlPayload::CreateLogicalTables(CreateLogicalTablesPayload { + version: TABLE_DDL_PAYLOAD_VERSION, + table_count, + }), + ) + } + + /// Builds the bounded event emitted when altering a table is submitted. + pub(crate) fn alter_table_submitted( + locator: TableDdlLocator, + kind: Option<&'static str>, + ) -> Self { + Self::submitted( + TableDdlEventType::AlterTable, + [locator], + TableDdlPayload::AlterTable(AlterTablePayload { + version: TABLE_DDL_PAYLOAD_VERSION, + kind, + }), + ) + } + + /// Builds the bounded event emitted when altering logical tables is submitted. + pub(crate) fn alter_logical_tables_submitted( + locators: impl IntoIterator, + table_count: usize, + kinds: impl IntoIterator, + ) -> Self { + let kinds = kinds + .into_iter() + .collect::>() + .into_iter() + .collect(); + Self::submitted( + TableDdlEventType::AlterLogicalTables, + locators, + TableDdlPayload::AlterLogicalTables(AlterLogicalTablesPayload { + version: TABLE_DDL_PAYLOAD_VERSION, + table_count, + kinds, + }), + ) + } + + /// Builds the bounded event emitted when dropping a table is submitted. + pub(crate) fn drop_table_submitted(locator: TableDdlLocator, drop_if_exists: bool) -> Self { + Self::submitted( + TableDdlEventType::DropTable, + [locator], + TableDdlPayload::DropTable(DropTablePayload { + version: TABLE_DDL_PAYLOAD_VERSION, + drop_if_exists, + }), + ) + } + + /// Builds the bounded event emitted when restoring a dropped table is submitted. + pub(crate) fn undrop_table_submitted(locator: TableDdlLocator) -> Self { + Self::submitted( + TableDdlEventType::UndropTable, + [locator], + TableDdlPayload::UndropTable(UndropTablePayload { + version: TABLE_DDL_PAYLOAD_VERSION, + }), + ) + } + + /// Builds the bounded event emitted when purging a dropped table is submitted. + pub(crate) fn purge_dropped_table_submitted(locator: TableDdlLocator) -> Self { + Self::submitted( + TableDdlEventType::PurgeDroppedTable, + [locator], + TableDdlPayload::PurgeDroppedTable(PurgeDroppedTablePayload { + version: TABLE_DDL_PAYLOAD_VERSION, + }), + ) + } + + /// Builds the bounded event emitted when truncating a table is submitted. + pub(crate) fn truncate_table_submitted( + locator: TableDdlLocator, + time_range_count: usize, + ) -> Self { + Self::submitted( + TableDdlEventType::TruncateTable, + [locator], + TableDdlPayload::TruncateTable(TruncateTablePayload { + version: TABLE_DDL_PAYLOAD_VERSION, + time_range_count, + }), + ) + } + + /// Builds a lightweight lifecycle event with null domain columns and payload. + pub(crate) fn lifecycle(event_type: TableDdlEventType) -> Self { + Self { + event_type, + locators: vec![TableDdlLocator::default()], + payload: None, + } + } + + /// Builds a Create Table success event containing only the allocated table ID. + pub(crate) fn create_table_succeeded(table_id: TableId) -> Self { + Self { + event_type: TableDdlEventType::CreateTable, + locators: vec![TableDdlLocator::from_table_id(table_id)], + payload: None, + } + } + + /// Builds Create Logical Tables success rows from their allocated locators. + pub(crate) fn create_logical_tables_succeeded( + locators: impl IntoIterator, + ) -> Self { + Self { + event_type: TableDdlEventType::CreateLogicalTables, + locators: locators.into_iter().collect(), + payload: None, + } + } + + fn submitted( + event_type: TableDdlEventType, + locators: impl IntoIterator, + payload: TableDdlPayload, + ) -> Self { + Self { + event_type, + locators: locators.into_iter().collect(), + payload: Some(payload), + } + } + + fn schema() -> Vec { + column_schemas([ + &CATALOG_NAME_COLUMN, + &SCHEMA_NAME_COLUMN, + &TABLE_NAME_COLUMN, + &TABLE_ID_COLUMN, + ]) + } + + fn locator_row(&self, locator: &TableDdlLocator) -> Row { + let mut values = vec![ + nullable_string(locator.catalog_name.as_deref()), + nullable_string(locator.schema_name.as_deref()), + nullable_string(locator.table_name.as_deref()), + nullable_table_id(locator.table_id), + ]; + if self.event_type.has_physical_table_id() { + values.push(nullable_table_id(locator.physical_table_id)); + } + Row { values } + } +} + +impl Event for TableDdlEvent { + fn event_type(&self) -> &str { + self.event_type.as_str() + } + + fn json_payload(&self) -> Result { + match &self.payload { + Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu), + None => Ok(JsonValue::Null), + } + } + + fn extra_schema(&self) -> Vec { + let mut schema = Self::schema(); + if self.event_type.has_physical_table_id() { + schema.push(PHYSICAL_TABLE_ID_COLUMN.column_schema()); + } + schema + } + + fn extra_rows(&self) -> Result> { + Ok(self + .locators + .iter() + .map(|locator| self.locator_row(locator)) + .collect()) + } + + fn as_any(&self) -> &dyn Any { + self + } +} + +fn nullable_table_id(value: Option) -> api::v1::Value { + nullable_value(value.map(ValueData::U32Value)) +} diff --git a/src/common/meta/src/ddl/purge_dropped_table.rs b/src/common/meta/src/ddl/purge_dropped_table.rs index 1548c72241..bc2f74ee57 100644 --- a/src/common/meta/src/ddl/purge_dropped_table.rs +++ b/src/common/meta/src/ddl/purge_dropped_table.rs @@ -17,7 +17,8 @@ use std::collections::HashMap; use async_trait::async_trait; use common_procedure::error::{FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ - Context as ProcedureContext, LockKey, Procedure, Result as ProcedureResult, Status, + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, + Result as ProcedureResult, Status, }; use common_telemetry::info; use common_time::util::current_time_millis; @@ -31,6 +32,7 @@ use table::table_name::TableName; use crate::ddl::DdlContext; use crate::ddl::drop_table::executor::DropTableExecutor; +use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator}; use crate::ddl::utils::{ convert_region_routes_to_detecting_regions, is_metric_engine_logical_table, map_to_procedure_error, @@ -256,6 +258,24 @@ impl Procedure for PurgeDroppedTableProcedure { fn lock_key(&self) -> LockKey { LockKey::new(vec![TableLock::Write(self.data.task.table_id).into()]) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::PurgeDroppedTable.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let locator = TableDdlLocator::from_table_id(self.data.task.table_id); + TableDdlEvent::purge_dropped_table_submitted(locator) + } + _ => TableDdlEvent::lifecycle(TableDdlEventType::PurgeDroppedTable), + }; + + Some(Box::new(event)) + } } #[derive(Debug, Serialize, Deserialize)] diff --git a/src/common/meta/src/ddl/tests/alter_logical_tables.rs b/src/common/meta/src/ddl/tests/alter_logical_tables.rs index bca56bafc2..bef06cd839 100644 --- a/src/common/meta/src/ddl/tests/alter_logical_tables.rs +++ b/src/common/meta/src/ddl/tests/alter_logical_tables.rs @@ -47,7 +47,7 @@ use crate::rpc::ddl::AlterTableTask; use crate::rpc::router::{Region, RegionRoute}; use crate::test_util::{MockDatanodeManager, new_ddl_context}; -fn make_alter_logical_table_add_column_task( +pub(crate) fn make_alter_logical_table_add_column_task( schema: Option<&str>, table: &str, add_columns: Vec, diff --git a/src/common/meta/src/ddl/tests/alter_table.rs b/src/common/meta/src/ddl/tests/alter_table.rs index 14ee71b3e6..173ca870b8 100644 --- a/src/common/meta/src/ddl/tests/alter_table.rs +++ b/src/common/meta/src/ddl/tests/alter_table.rs @@ -133,7 +133,7 @@ async fn test_on_prepare_table_not_exists_err() { assert_matches!(err.status_code(), StatusCode::TableNotFound); } -fn test_alter_table_task(table_name: &str) -> AlterTableTask { +pub(crate) fn test_alter_table_task(table_name: &str) -> AlterTableTask { AlterTableTask { alter_table: AlterTableExpr { catalog_name: DEFAULT_CATALOG_NAME.to_string(), diff --git a/src/common/meta/src/ddl/tests/event.rs b/src/common/meta/src/ddl/tests/event.rs index c45638fc7a..7bb5af53fc 100644 --- a/src/common/meta/src/ddl/tests/event.rs +++ b/src/common/meta/src/ddl/tests/event.rs @@ -14,5 +14,6 @@ mod database; mod flow; +mod table; mod test_util; mod view; diff --git a/src/common/meta/src/ddl/tests/event/database.rs b/src/common/meta/src/ddl/tests/event/database.rs index 453e1a0f16..c617a12bf0 100644 --- a/src/common/meta/src/ddl/tests/event/database.rs +++ b/src/common/meta/src/ddl/tests/event/database.rs @@ -12,12 +12,13 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::{HashMap, HashSet}; +use std::collections::HashMap; use std::sync::Arc; use api::v1::alter_database_expr::Kind as PbAlterDatabaseKind; use api::v1::value::ValueData; use api::v1::{AlterDatabaseExpr, Row, SetDatabaseOptions as PbSetDatabaseOptions, Value}; +use common_event_recorder::Event; use common_event_recorder::event_table::{ CATALOG_NAME_COLUMN as EVENT_TABLE_CATALOG_NAME_COLUMN, PROCEDURE_ERROR_COLUMN as EVENT_TABLE_PROCEDURE_ERROR_COLUMN, @@ -27,11 +28,9 @@ use common_event_recorder::event_table::{ SCHEMA_NAME_COLUMN as EVENT_TABLE_SCHEMA_NAME_COLUMN, }; use common_event_recorder::testing::assert_event_contract; -use common_event_recorder::{Event, EventTypeFilter}; -use common_procedure::{ - EventContext, EventTrigger, Procedure, ProcedureEvent, ProcedureId, ProcedureState, -}; +use common_procedure::{EventTrigger, ProcedureEvent, ProcedureId, ProcedureState}; +use super::test_util::assert_event_filter; use crate::ddl::alter_database::AlterDatabaseProcedure; use crate::ddl::create_database::CreateDatabaseProcedure; use crate::ddl::drop_database::DropDatabaseProcedure; @@ -221,7 +220,7 @@ fn test_create_database_event_filter() { new_ddl_context(Arc::new(MockDatanodeManager::new(()))), ); - assert_database_event_filter(&procedure, CREATE_DATABASE_EVENT_TYPE); + assert_event_filter(&procedure, CREATE_DATABASE_EVENT_TYPE); } #[test] @@ -242,7 +241,7 @@ fn test_alter_database_event_filter() { ) .unwrap(); - assert_database_event_filter(&procedure, ALTER_DATABASE_EVENT_TYPE); + assert_event_filter(&procedure, ALTER_DATABASE_EVENT_TYPE); } #[test] @@ -254,37 +253,7 @@ fn test_drop_database_event_filter() { new_ddl_context(Arc::new(MockDatanodeManager::new(()))), ); - assert_database_event_filter(&procedure, DROP_DATABASE_EVENT_TYPE); -} - -fn assert_database_event_filter(procedure: &dyn Procedure, event_type: &str) { - let state = ProcedureState::Running; - let event_context = |event_type_filter| EventContext { - procedure_id: ProcedureId::random(), - lifecycle_state: &state, - trigger: EventTrigger::Submitted, - event_type_filter: Arc::new(event_type_filter), - }; - - let allowed = procedure - .event(&event_context(EventTypeFilter::Only(HashSet::from([ - event_type.to_string(), - ])))) - .unwrap(); - assert_eq!(allowed.event_type(), event_type); - - assert!( - procedure - .event(&event_context(EventTypeFilter::Only(HashSet::from([ - String::from("other_event",) - ])))) - .is_none() - ); - assert!( - procedure - .event(&event_context(EventTypeFilter::Only(HashSet::new()))) - .is_none() - ); + assert_event_filter(&procedure, DROP_DATABASE_EVENT_TYPE); } fn assert_event_locator( diff --git a/src/common/meta/src/ddl/tests/event/flow.rs b/src/common/meta/src/ddl/tests/event/flow.rs index 47a69ab921..4c4fd8a222 100644 --- a/src/common/meta/src/ddl/tests/event/flow.rs +++ b/src/common/meta/src/ddl/tests/event/flow.rs @@ -12,23 +12,21 @@ // See the License for the specific language governing permissions and // limitations under the License. -use std::collections::HashSet; use std::sync::Arc; use api::v1::value::ValueData; use api::v1::{ColumnSchema, Row, Value}; use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; +use common_event_recorder::Event; use common_event_recorder::event_table::{ CATALOG_NAME_COLUMN, FLOW_ID_COLUMN, FLOW_NAME_COLUMN, PROCEDURE_ERROR_COLUMN, PROCEDURE_ID_COLUMN, PROCEDURE_STATE_COLUMN, PROCEDURE_TRIGGER_COLUMN, }; use common_event_recorder::testing::assert_event_contract; -use common_event_recorder::{Event, EventTypeFilter}; -use common_procedure::{ - EventContext, EventTrigger, Procedure, ProcedureEvent, ProcedureId, ProcedureState, -}; +use common_procedure::{EventTrigger, ProcedureEvent, ProcedureId, ProcedureState}; use table::table_name::TableName; +use super::test_util::assert_event_filter; use crate::ddl::create_flow::CreateFlowProcedure; use crate::ddl::drop_flow::DropFlowProcedure; use crate::ddl::event::flow::{ @@ -191,7 +189,7 @@ fn test_create_flow_event_filter() { test_query_context(), new_ddl_context(Arc::new(MockFlownodeManager::new(NaiveFlownodeHandler))), ); - assert_flow_event_filter(&procedure, CREATE_FLOW_EVENT_TYPE); + assert_event_filter(&procedure, CREATE_FLOW_EVENT_TYPE); } #[test] @@ -200,7 +198,7 @@ fn test_drop_flow_event_filter() { test_drop_flow_task("flow", 42, false), new_ddl_context(Arc::new(MockFlownodeManager::new(NaiveFlownodeHandler))), ); - assert_flow_event_filter(&procedure, DROP_FLOW_EVENT_TYPE); + assert_event_filter(&procedure, DROP_FLOW_EVENT_TYPE); } fn flow_schema() -> Vec { @@ -253,36 +251,6 @@ fn assert_procedure_event_contract( ); } -fn assert_flow_event_filter(procedure: &dyn Procedure, event_type: &str) { - let state = ProcedureState::Running; - let event_context = |event_type_filter| EventContext { - procedure_id: ProcedureId::random(), - lifecycle_state: &state, - trigger: EventTrigger::Submitted, - event_type_filter: Arc::new(event_type_filter), - }; - - assert!( - procedure - .event(&event_context(EventTypeFilter::Only(HashSet::from([ - event_type.to_string() - ])))) - .is_some() - ); - assert!( - procedure - .event(&event_context(EventTypeFilter::Only(HashSet::new()))) - .is_none() - ); - assert!( - procedure - .event(&event_context(EventTypeFilter::Only(HashSet::from([ - "other_event".to_string(), - ])))) - .is_none() - ); -} - fn optional_string(value: Option<&str>) -> Value { value .map(|value| ValueData::StringValue(value.to_string()).into()) diff --git a/src/common/meta/src/ddl/tests/event/table.rs b/src/common/meta/src/ddl/tests/event/table.rs new file mode 100644 index 0000000000..eb617abaa8 --- /dev/null +++ b/src/common/meta/src/ddl/tests/event/table.rs @@ -0,0 +1,552 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; + +use api::v1::alter_table_expr::Kind as AlterTableKind; +use api::v1::value::ValueData; +use api::v1::{ColumnDataType, Repartition, SemanticType, Value}; +use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME}; +use common_event_recorder::event_table::{ + CATALOG_NAME_COLUMN, PHYSICAL_TABLE_ID_COLUMN, SCHEMA_NAME_COLUMN, TABLE_ID_COLUMN, + TABLE_NAME_COLUMN, +}; +use common_event_recorder::{Event, EventTypeFilter}; +use common_procedure::{ + ChildSubmissionOutcome, EventContext, EventTrigger, Procedure, ProcedureId, ProcedureState, + RetryPhase, +}; +use common_time::Timestamp; +use serde_json::{Value as JsonValue, json}; + +use super::test_util::assert_event_filter; +use crate::ddl::alter_logical_tables::AlterLogicalTablesProcedure; +use crate::ddl::alter_table::AlterTableProcedure; +use crate::ddl::create_logical_tables::CreateLogicalTablesProcedure; +use crate::ddl::create_table::CreateTableProcedure; +use crate::ddl::drop_table::DropTableProcedure; +use crate::ddl::event::table::{ + TABLE_DDL_PAYLOAD_VERSION, TableDdlEvent, TableDdlEventType, TableDdlLocator, + alter_table_kind_name, +}; +use crate::ddl::purge_dropped_table::PurgeDroppedTableProcedure; +use crate::ddl::test_util::create_table::test_create_table_task as test_create_table_task_with_id; +use crate::ddl::test_util::test_create_logical_table_task; +use crate::ddl::tests::alter_logical_tables::make_alter_logical_table_add_column_task; +use crate::ddl::tests::alter_table::test_alter_table_task; +use crate::ddl::tests::create_table::test_create_table_task; +use crate::ddl::truncate_table::TruncateTableProcedure; +use crate::ddl::undrop_table::UndropTableProcedure; +use crate::key::DeserializedValueWithBytes; +use crate::key::table_info::TableInfoValue; +use crate::rpc::ddl::{DropTableTask, PurgeDroppedTableTask, TruncateTableTask, UndropTableTask}; +use crate::test_util::{MockDatanodeManager, new_ddl_context}; + +struct EventCase { + event_type: TableDdlEventType, + event: TableDdlEvent, + payload: JsonValue, + rows: Vec>, +} + +struct ProcedureCase { + procedure: Box, + event_type: &'static str, + payload: JsonValue, + rows: Vec>, +} + +#[test] +fn repartition_kind_is_not_supported_by_alter_table_events() { + assert_eq!( + None, + alter_table_kind_name(&AlterTableKind::Repartition(Repartition::default())) + ); +} + +#[test] +fn submitted_event_contracts_are_bounded_and_fixed() { + for case in event_cases() { + assert_eq!(case.event.event_type(), case.event_type.as_str()); + assert_eq!(case.event.json_payload().unwrap(), case.payload); + assert_eq!( + case.event + .extra_rows() + .unwrap() + .into_iter() + .map(|row| row.values) + .collect::>(), + case.rows + ); + + let schema = case.event.extra_schema(); + assert_eq!( + schema + .iter() + .map(|column| (column.column_name.as_str(), column.datatype)) + .collect::>(), + expected_schema(case.event_type) + ); + assert!( + schema + .iter() + .all(|column| column.semantic_type == SemanticType::Field as i32) + ); + + let lifecycle = TableDdlEvent::lifecycle(case.event_type); + assert_eq!(lifecycle.extra_schema(), schema); + assert_eq!(lifecycle.json_payload().unwrap(), JsonValue::Null); + assert_eq!( + lifecycle.extra_rows().unwrap()[0].values, + vec![Value::default(); schema.len()] + ); + } +} + +#[test] +fn procedures_map_tasks_to_submitted_events() { + for case in procedure_cases() { + let event = event_for(case.procedure.as_ref(), EventTrigger::Submitted); + + assert_eq!(event.event_type(), case.event_type); + assert_eq!(event.json_payload().unwrap(), case.payload); + assert_eq!( + event + .extra_rows() + .unwrap() + .into_iter() + .map(|row| row.values) + .collect::>(), + case.rows + ); + } +} + +#[test] +fn later_lifecycle_events_are_uniform() { + let triggers = [ + EventTrigger::Recovered, + EventTrigger::ChildSubmitted { + procedure_id: ProcedureId::random(), + outcome: ChildSubmissionOutcome::Accepted, + }, + EventTrigger::Retrying { + phase: RetryPhase::Execute, + attempt: 1, + }, + EventTrigger::RollingBack, + EventTrigger::Succeeded, + EventTrigger::Failed, + EventTrigger::Poisoned, + ]; + + for case in procedure_cases() { + let submitted = event_for(case.procedure.as_ref(), EventTrigger::Submitted); + let schema = submitted.extra_schema(); + + for trigger in &triggers { + let event = event_for(case.procedure.as_ref(), trigger.clone()); + + assert_eq!(event.event_type(), case.event_type); + assert_eq!(event.extra_schema(), schema); + assert_eq!(event.json_payload().unwrap(), JsonValue::Null); + assert_eq!( + event.extra_rows().unwrap()[0].values, + vec![Value::default(); schema.len()] + ); + } + } +} + +#[test] +fn create_success_events_keep_allocated_ids() { + let mut task = test_create_table_task("create_success"); + task.table_info.ident.table_id = 7; + let create_table = CreateTableProcedure::new(task, test_context()).unwrap(); + let state = ProcedureState::Done { + output: Some(Arc::new(42_u32)), + }; + let event = event_for_state(&create_table, EventTrigger::Succeeded, &state); + + assert_eq!(event.event_type(), "create_table"); + assert_eq!(event.json_payload().unwrap(), JsonValue::Null); + assert_eq!( + event.extra_rows().unwrap()[0].values, + table_locator_values(None, Some(42)) + ); + + let logical_tables = CreateLogicalTablesProcedure::new( + vec![ + test_create_logical_table_task("foo"), + test_create_logical_table_task("bar"), + ], + 1024, + test_context(), + ); + let state = ProcedureState::Done { + output: Some(Arc::new(vec![1025_u32, 1026_u32])), + }; + let event = event_for_state(&logical_tables, EventTrigger::Succeeded, &state); + + assert_eq!(event.event_type(), "create_logical_tables"); + assert_eq!(event.json_payload().unwrap(), JsonValue::Null); + assert_eq!( + event + .extra_rows() + .unwrap() + .into_iter() + .map(|row| row.values) + .collect::>(), + vec![ + logical_locator_values("foo", Some(1025), 1024), + logical_locator_values("bar", Some(1026), 1024), + ] + ); +} + +#[test] +fn table_procedures_honor_event_type_filter() { + for case in procedure_cases() { + assert_event_filter(case.procedure.as_ref(), case.event_type); + } +} + +fn event_cases() -> Vec { + vec![ + EventCase { + event_type: TableDdlEventType::CreateTable, + event: TableDdlEvent::create_table_submitted( + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "create"), + true, + "mito2", + ), + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "create_if_not_exists": true, + "engine": "mito2", + }), + rows: vec![table_locator_values(Some("create"), None)], + }, + EventCase { + event_type: TableDdlEventType::CreateLogicalTables, + event: TableDdlEvent::create_logical_tables_submitted( + [ + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "logical1") + .with_physical_table_id(10), + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "logical2") + .with_physical_table_id(10), + ], + 2, + ), + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "table_count": 2, + }), + rows: vec![ + logical_locator_values("logical1", None, 10), + logical_locator_values("logical2", None, 10), + ], + }, + EventCase { + event_type: TableDdlEventType::AlterTable, + event: TableDdlEvent::alter_table_submitted( + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "alter") + .with_table_id(11), + Some("drop_columns"), + ), + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "kind": "drop_columns", + }), + rows: vec![table_locator_values(Some("alter"), Some(11))], + }, + EventCase { + event_type: TableDdlEventType::AlterLogicalTables, + event: TableDdlEvent::alter_logical_tables_submitted( + [ + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "logical1") + .with_physical_table_id(10), + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "logical2") + .with_physical_table_id(10), + ], + 2, + ["rename_table", "add_columns", "add_columns"], + ), + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "table_count": 2, + "kinds": ["add_columns", "rename_table"], + }), + rows: vec![ + logical_locator_values("logical1", None, 10), + logical_locator_values("logical2", None, 10), + ], + }, + EventCase { + event_type: TableDdlEventType::DropTable, + event: TableDdlEvent::drop_table_submitted( + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "drop") + .with_table_id(12), + true, + ), + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "drop_if_exists": true, + }), + rows: vec![table_locator_values(Some("drop"), Some(12))], + }, + EventCase { + event_type: TableDdlEventType::UndropTable, + event: TableDdlEvent::undrop_table_submitted(TableDdlLocator::from_table_id(13)), + payload: json!({"version": TABLE_DDL_PAYLOAD_VERSION}), + rows: vec![table_locator_values(None, Some(13))], + }, + EventCase { + event_type: TableDdlEventType::PurgeDroppedTable, + event: TableDdlEvent::purge_dropped_table_submitted(TableDdlLocator::from_table_id(14)), + payload: json!({"version": TABLE_DDL_PAYLOAD_VERSION}), + rows: vec![table_locator_values(None, Some(14))], + }, + EventCase { + event_type: TableDdlEventType::TruncateTable, + event: TableDdlEvent::truncate_table_submitted( + TableDdlLocator::new(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, "truncate") + .with_table_id(15), + 4, + ), + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "time_range_count": 4, + }), + rows: vec![table_locator_values(Some("truncate"), Some(15))], + }, + ] +} + +fn procedure_cases() -> Vec { + let create_table = + CreateTableProcedure::new(test_create_table_task("create"), test_context()).unwrap(); + let create_logical_tables = CreateLogicalTablesProcedure::new( + vec![ + test_create_logical_table_task("logical1"), + test_create_logical_table_task("logical2"), + ], + 41, + test_context(), + ); + let alter_table = + AlterTableProcedure::new(42, test_alter_table_task("alter"), test_context()).unwrap(); + let alter_logical_tables = AlterLogicalTablesProcedure::new( + vec![ + make_alter_logical_table_add_column_task( + Some(DEFAULT_SCHEMA_NAME), + "logical1", + vec!["tag1".to_string()], + ), + make_alter_logical_table_add_column_task( + Some(DEFAULT_SCHEMA_NAME), + "logical2", + vec!["tag2".to_string()], + ), + ], + 43, + test_context(), + ); + let drop_table = DropTableProcedure::new( + DropTableTask { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table: "drop".to_string(), + table_id: 44, + drop_if_exists: true, + }, + test_context(), + ); + let undrop_table = UndropTableProcedure::new(UndropTableTask { table_id: 45 }, test_context()); + let purge_dropped_table = + PurgeDroppedTableProcedure::new(PurgeDroppedTableTask { table_id: 46 }, test_context()); + let truncate_table = truncate_procedure(TruncateTableTask { + catalog: DEFAULT_CATALOG_NAME.to_string(), + schema: DEFAULT_SCHEMA_NAME.to_string(), + table: "truncate".to_string(), + table_id: 47, + time_ranges: vec![( + Timestamp::new_millisecond(1_000), + Timestamp::new_millisecond(2_000), + )], + }); + + vec![ + ProcedureCase { + procedure: Box::new(create_table), + event_type: "create_table", + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "create_if_not_exists": false, + "engine": "mito2", + }), + rows: vec![table_locator_values(Some("create"), None)], + }, + ProcedureCase { + procedure: Box::new(create_logical_tables), + event_type: "create_logical_tables", + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "table_count": 2, + }), + rows: vec![ + logical_locator_values("logical1", None, 41), + logical_locator_values("logical2", None, 41), + ], + }, + ProcedureCase { + procedure: Box::new(alter_table), + event_type: "alter_table", + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "kind": "drop_columns", + }), + rows: vec![table_locator_values(Some("alter"), Some(42))], + }, + ProcedureCase { + procedure: Box::new(alter_logical_tables), + event_type: "alter_logical_tables", + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "table_count": 2, + "kinds": ["add_columns"], + }), + rows: vec![ + logical_locator_values("logical1", None, 43), + logical_locator_values("logical2", None, 43), + ], + }, + ProcedureCase { + procedure: Box::new(drop_table), + event_type: "drop_table", + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "drop_if_exists": true, + }), + rows: vec![table_locator_values(Some("drop"), Some(44))], + }, + ProcedureCase { + procedure: Box::new(undrop_table), + event_type: "undrop_table", + payload: json!({"version": TABLE_DDL_PAYLOAD_VERSION}), + rows: vec![table_locator_values(None, Some(45))], + }, + ProcedureCase { + procedure: Box::new(purge_dropped_table), + event_type: "purge_dropped_table", + payload: json!({"version": TABLE_DDL_PAYLOAD_VERSION}), + rows: vec![table_locator_values(None, Some(46))], + }, + ProcedureCase { + procedure: Box::new(truncate_table), + event_type: "truncate_table", + payload: json!({ + "version": TABLE_DDL_PAYLOAD_VERSION, + "time_range_count": 1, + }), + rows: vec![table_locator_values(Some("truncate"), Some(47))], + }, + ] +} + +fn expected_schema(event_type: TableDdlEventType) -> Vec<(&'static str, i32)> { + let mut schema = vec![ + (CATALOG_NAME_COLUMN.name(), ColumnDataType::String as i32), + (SCHEMA_NAME_COLUMN.name(), ColumnDataType::String as i32), + (TABLE_NAME_COLUMN.name(), ColumnDataType::String as i32), + (TABLE_ID_COLUMN.name(), ColumnDataType::Uint32 as i32), + ]; + if matches!( + event_type, + TableDdlEventType::CreateLogicalTables | TableDdlEventType::AlterLogicalTables + ) { + schema.push(( + PHYSICAL_TABLE_ID_COLUMN.name(), + ColumnDataType::Uint32 as i32, + )); + } + schema +} + +fn table_locator_values(table_name: Option<&str>, table_id: Option) -> Vec { + let (catalog_name, schema_name) = if table_name.is_some() { + ( + string_value(DEFAULT_CATALOG_NAME), + string_value(DEFAULT_SCHEMA_NAME), + ) + } else { + (Value::default(), Value::default()) + }; + vec![ + catalog_name, + schema_name, + table_name.map(string_value).unwrap_or_default(), + table_id.map(table_id_value).unwrap_or_default(), + ] +} + +fn logical_locator_values( + table_name: &str, + table_id: Option, + physical_table_id: u32, +) -> Vec { + let mut values = table_locator_values(Some(table_name), table_id); + values.push(table_id_value(physical_table_id)); + values +} + +fn string_value(value: &str) -> Value { + ValueData::StringValue(value.to_string()).into() +} + +fn table_id_value(value: u32) -> Value { + ValueData::U32Value(value).into() +} + +fn truncate_procedure(task: TruncateTableTask) -> TruncateTableProcedure { + let table_info = test_create_table_task_with_id("metrics", task.table_id).table_info; + TruncateTableProcedure::new( + task, + DeserializedValueWithBytes::from_inner(TableInfoValue::new(table_info)), + test_context(), + ) +} + +fn test_context() -> crate::ddl::DdlContext { + new_ddl_context(Arc::new(MockDatanodeManager::new(()))) +} + +fn event_for(procedure: &dyn Procedure, trigger: EventTrigger) -> Box { + event_for_state(procedure, trigger, &ProcedureState::Running) +} + +fn event_for_state( + procedure: &dyn Procedure, + trigger: EventTrigger, + lifecycle_state: &ProcedureState, +) -> Box { + procedure + .event(&EventContext { + procedure_id: ProcedureId::random(), + lifecycle_state, + trigger, + event_type_filter: Arc::new(EventTypeFilter::All), + }) + .unwrap() +} diff --git a/src/common/meta/src/ddl/tests/event/test_util.rs b/src/common/meta/src/ddl/tests/event/test_util.rs index f861340dea..884fa3ec4d 100644 --- a/src/common/meta/src/ddl/tests/event/test_util.rs +++ b/src/common/meta/src/ddl/tests/event/test_util.rs @@ -18,7 +18,7 @@ use std::sync::Arc; use common_event_recorder::EventTypeFilter; use common_procedure::{EventContext, EventTrigger, Procedure, ProcedureId, ProcedureState}; -pub(super) fn assert_event_filter(procedure: &dyn Procedure, event_type: &str) { +pub(crate) fn assert_event_filter(procedure: &dyn Procedure, event_type: &str) { let state = ProcedureState::Running; let event_context = |event_type_filter| EventContext { procedure_id: ProcedureId::random(), diff --git a/src/common/meta/src/ddl/truncate_table.rs b/src/common/meta/src/ddl/truncate_table.rs index 21491b6230..454b7e2de4 100644 --- a/src/common/meta/src/ddl/truncate_table.rs +++ b/src/common/meta/src/ddl/truncate_table.rs @@ -20,7 +20,8 @@ use api::v1::region::{ use async_trait::async_trait; use common_procedure::error::{FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ - Context as ProcedureContext, LockKey, Procedure, Result as ProcedureResult, Status, + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, + Result as ProcedureResult, Status, }; use common_telemetry::debug; use common_telemetry::tracing_context::TracingContext; @@ -34,6 +35,7 @@ use table::table_name::TableName; use table::table_reference::TableReference; use crate::ddl::DdlContext; +use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator}; use crate::ddl::utils::{add_peer_context_if_needed, map_to_procedure_error}; use crate::error::{ConvertTimeRangesSnafu, Result, TableNotFoundSnafu}; use crate::key::DeserializedValueWithBytes; @@ -86,6 +88,26 @@ impl Procedure for TruncateTableProcedure { LockKey::new(lock_key) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::TruncateTable.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let task = &self.data.task; + let locator = TableDdlLocator::new(&task.catalog, &task.schema, &task.table) + .with_table_id(task.table_id); + TableDdlEvent::truncate_table_submitted(locator, task.time_ranges.len()) + } + _ => TableDdlEvent::lifecycle(TableDdlEventType::TruncateTable), + }; + + Some(Box::new(event)) + } } impl TruncateTableProcedure { diff --git a/src/common/meta/src/ddl/undrop_table.rs b/src/common/meta/src/ddl/undrop_table.rs index 57fe8efa89..8b21c72f29 100644 --- a/src/common/meta/src/ddl/undrop_table.rs +++ b/src/common/meta/src/ddl/undrop_table.rs @@ -20,7 +20,8 @@ use api::v1::region::{ use async_trait::async_trait; use common_procedure::error::{FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ - Context as ProcedureContext, LockKey, Procedure, Result as ProcedureResult, Status, + Context as ProcedureContext, EventContext, EventTrigger, LockKey, Procedure, + Result as ProcedureResult, Status, }; use common_telemetry::tracing_context::TracingContext; use common_telemetry::warn; @@ -34,6 +35,7 @@ use table::metadata::TableId; use table::table_name::TableName; use crate::ddl::drop_table::executor::DropTableExecutor; +use crate::ddl::event::table::{TableDdlEvent, TableDdlEventType, TableDdlLocator}; use crate::ddl::utils::{ add_peer_context_if_needed, convert_region_routes_to_detecting_regions, is_metric_engine_logical_table, map_to_procedure_error, region_storage_path, @@ -426,6 +428,36 @@ impl Procedure for UndropTableProcedure { lock_key.push(TableLock::Write(self.data.task.table_id).into()); LockKey::new(lock_key) } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + if !ctx + .event_type_filter + .allows(TableDdlEventType::UndropTable.as_str()) + { + return None; + } + let event = match &ctx.trigger { + EventTrigger::Submitted => { + let locator = self + .data + .table_name + .as_ref() + .map(|table_name| { + TableDdlLocator::new( + &table_name.catalog_name, + &table_name.schema_name, + &table_name.table_name, + ) + }) + .unwrap_or_default() + .with_table_id(self.data.task.table_id); + TableDdlEvent::undrop_table_submitted(locator) + } + _ => TableDdlEvent::lifecycle(TableDdlEventType::UndropTable), + }; + + Some(Box::new(event)) + } } pub(crate) async fn open_regions( diff --git a/tests-integration/tests/event_recorder_test_util.rs b/tests-integration/tests/event_recorder_test_util.rs index 43d923f810..5b21ef45b2 100644 --- a/tests-integration/tests/event_recorder_test_util.rs +++ b/tests-integration/tests/event_recorder_test_util.rs @@ -18,6 +18,7 @@ use std::time::Duration; use client::OutputData; use common_recordbatch::RecordBatches; use datatypes::arrow::array::AsArray; +use datatypes::arrow::datatypes::UInt32Type; use frontend::instance::Instance; use servers::query_handler::sql::SqlQueryHandler; use session::context::QueryContext; @@ -53,6 +54,22 @@ pub(crate) async fn find_eventually_string( panic!("timed out waiting for event query: {query}"); } +/// Returns the first non-null u32 from `column` once the query produces one. +pub(crate) async fn find_eventually_u32( + instance: &Arc, + query: &str, + column: &str, +) -> u32 { + for _ in 0..MAX_ATTEMPTS { + if let Some(value) = query_first_u32(instance, query, column).await { + return value; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + + panic!("timed out waiting for event query: {query}"); +} + /// Asserts that the query eventually returns the expected pretty-printed rows. pub(crate) async fn assert_eventually_eq(instance: &Arc, query: &str, expected: &str) { let mut last_actual = None; @@ -94,6 +111,17 @@ async fn query_first_string(instance: &Arc, query: &str, column: &str) .map(ToString::to_string) } +async fn query_first_u32(instance: &Arc, query: &str, column: &str) -> Option { + let batches = query_record_batches(instance, query).await?; + let batch = batches.take().into_iter().next()?; + batch + .column_by_name(column)? + .as_primitive::() + .iter() + .next() + .flatten() +} + async fn query_pretty_print(instance: &Arc, query: &str) -> Option { query_record_batches(instance, query) .await? diff --git a/tests-integration/tests/main.rs b/tests-integration/tests/main.rs index 1bc4c1a24e..247cb8d54f 100644 --- a/tests-integration/tests/main.rs +++ b/tests-integration/tests/main.rs @@ -27,6 +27,7 @@ mod jsonbench; mod sql; #[macro_use] mod region_migration; +mod table_ddl_event; #[macro_use] mod repartition; mod repartition_event; diff --git a/tests-integration/tests/table_ddl_event.rs b/tests-integration/tests/table_ddl_event.rs new file mode 100644 index 0000000000..26f4afd1a0 --- /dev/null +++ b/tests-integration/tests/table_ddl_event.rs @@ -0,0 +1,452 @@ +// Copyright 2023 Greptime Team +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +use std::sync::Arc; +use std::time::Duration; + +use client::OutputData; +use common_event_recorder::event_table::{ + CATALOG_NAME_COLUMN, PAYLOAD_COLUMN, PHYSICAL_TABLE_ID_COLUMN, PROCEDURE_ID_COLUMN, + PROCEDURE_TRIGGER_COLUMN, SCHEMA_NAME_COLUMN, TABLE_ID_COLUMN, TABLE_NAME_COLUMN, TYPE_COLUMN, +}; +use common_test_util::temp_dir::create_temp_dir; +use frontend::instance::Instance; +use meta_srv::gc::GcSchedulerOptions; +use mito2::gc::GcConfig; +use servers::query_handler::sql::SqlQueryHandler; +use session::context::QueryContext; +use tests_integration::cluster::GreptimeDbClusterBuilder; + +use crate::event_recorder_test_util::{ + assert_eventually_eq, assert_single_event, find_eventually_string, find_eventually_u32, +}; + +const EVENTS_TABLE: &str = "greptime_private.events"; +const TABLE: &str = "table_ddl_events"; +const PHYSICAL_TABLE: &str = "table_ddl_events_phy"; +const LOGICAL_TABLE: &str = "table_ddl_events_logical"; + +#[tokio::test(flavor = "multi_thread")] +async fn test_table_ddl_procedure_events() { + common_telemetry::init_default_ut_logging(); + + // Arrange: use a distributed cluster to submit Undrop and Purge directly + // through the metasrv DDL manager and exercise logical-table task grouping. + let home_dir = create_temp_dir("table_ddl_procedure_events"); + let mut gc_options = GcSchedulerOptions { + enable: true, + ..Default::default() + }; + gc_options.experimental_soft_drop.enable = true; + gc_options.experimental_soft_drop.retention = Duration::from_secs(1); + let cluster = GreptimeDbClusterBuilder::new("table_ddl_procedure_events") + .await + .with_datanodes(1) + .with_metasrv_gc_config(gc_options) + .with_datanode_gc_config(GcConfig { + enable: true, + ..Default::default() + }) + .with_shared_home_dir(Arc::new(home_dir)) + .build(true) + .await; + let frontend = cluster.fe_instance().clone(); + + // Act / Assert: Create Table retains its rich submitted row and records the + // created table ID when it completes. + run_sql( + &frontend, + &format!( + "CREATE TABLE {TABLE} (host STRING PRIMARY KEY, ts TIMESTAMP TIME INDEX, val DOUBLE)" + ), + ) + .await; + let table_id = find_table_id(&frontend, TABLE).await; + let create_table_procedure_id = submitted_procedure_id(&frontend, "create_table", TABLE).await; + assert_named_submitted_event( + &frontend, + "create_table", + &create_table_procedure_id, + "\ ++--------------+------------------+ +| type | name | ++--------------+------------------+ +| create_table | table_ddl_events | ++--------------+------------------+", + ) + .await; + assert_id_terminal_event( + &frontend, + "create_table", + &create_table_procedure_id, + table_id, + "Succeeded", + "\ ++--------------+ +| type | ++--------------+ +| create_table | ++--------------+", + ) + .await; + + // Act / Assert: Alter and Truncate have a rich submitted row and lightweight + // terminal lifecycle row. + run_sql( + &frontend, + &format!("ALTER TABLE {TABLE} ADD COLUMN extra STRING"), + ) + .await; + let alter_table_procedure_id = submitted_procedure_id(&frontend, "alter_table", TABLE).await; + assert_named_submitted_event( + &frontend, + "alter_table", + &alter_table_procedure_id, + "\ ++-------------+------------------+ +| type | name | ++-------------+------------------+ +| alter_table | table_ddl_events | ++-------------+------------------+", + ) + .await; + assert_lightweight_terminal_event( + &frontend, + "alter_table", + &alter_table_procedure_id, + "Succeeded", + ) + .await; + run_sql( + &frontend, + &format!("INSERT INTO {TABLE} VALUES ('a', 0, 1, 'b')"), + ) + .await; + run_sql(&frontend, &format!("TRUNCATE TABLE {TABLE}")).await; + let truncate_table_procedure_id = + submitted_procedure_id(&frontend, "truncate_table", TABLE).await; + assert_named_submitted_event( + &frontend, + "truncate_table", + &truncate_table_procedure_id, + "\ ++----------------+------------------+ +| type | name | ++----------------+------------------+ +| truncate_table | table_ddl_events | ++----------------+------------------+", + ) + .await; + assert_lightweight_terminal_event( + &frontend, + "truncate_table", + &truncate_table_procedure_id, + "Succeeded", + ) + .await; + + // Act / Assert: Drop records the table name, while Undrop and Purge use the + // table ID as their submitted locator. + run_sql(&frontend, &format!("DROP TABLE {TABLE}")).await; + let drop_table_procedure_id = submitted_procedure_id(&frontend, "drop_table", TABLE).await; + assert_named_submitted_event( + &frontend, + "drop_table", + &drop_table_procedure_id, + "\ ++------------+------------------+ +| type | name | ++------------+------------------+ +| drop_table | table_ddl_events | ++------------+------------------+", + ) + .await; + assert_lightweight_terminal_event( + &frontend, + "drop_table", + &drop_table_procedure_id, + "Succeeded", + ) + .await; + let (undrop_table_procedure_id, _) = cluster + .metasrv + .ddl_manager() + .submit_undrop_table_task(common_meta::rpc::ddl::UndropTableTask { table_id }) + .await + .unwrap(); + let undrop_table_procedure_id = undrop_table_procedure_id.to_string(); + assert_id_submitted_event( + &frontend, + "undrop_table", + &undrop_table_procedure_id, + "\ ++--------------+ +| type | ++--------------+ +| undrop_table | ++--------------+", + ) + .await; + assert_lightweight_terminal_event( + &frontend, + "undrop_table", + &undrop_table_procedure_id, + "Succeeded", + ) + .await; + run_sql(&frontend, &format!("DROP TABLE {TABLE}")).await; + let (purge_table_procedure_id, _) = cluster + .metasrv + .ddl_manager() + .submit_purge_dropped_table_task(common_meta::rpc::ddl::PurgeDroppedTableTask { table_id }) + .await + .unwrap(); + let purge_table_procedure_id = purge_table_procedure_id.to_string(); + assert_id_submitted_event( + &frontend, + "purge_dropped_table", + &purge_table_procedure_id, + "\ ++---------------------+ +| type | ++---------------------+ +| purge_dropped_table | ++---------------------+", + ) + .await; + assert_lightweight_terminal_event( + &frontend, + "purge_dropped_table", + &purge_table_procedure_id, + "Succeeded", + ) + .await; + + // Arrange / Act: metric logical-table DDL is executed through the distributed + // frontend, which groups logical DDL tasks for the meta procedure manager. + run_sql( + &frontend, + &format!( + "CREATE TABLE {PHYSICAL_TABLE} (ts TIMESTAMP TIME INDEX, val DOUBLE) ENGINE=metric WITH (\"physical_metric_table\" = \"\")" + ), + ) + .await; + let physical_table_id = find_table_id(&frontend, PHYSICAL_TABLE).await; + run_sql( + &frontend, + &format!( + "CREATE TABLE {LOGICAL_TABLE} (ts TIMESTAMP TIME INDEX, val DOUBLE, host STRING PRIMARY KEY) ENGINE=metric WITH (\"on_physical_table\" = \"{PHYSICAL_TABLE}\")" + ), + ) + .await; + let logical_table_id = find_table_id(&frontend, LOGICAL_TABLE).await; + let create_logical_tables_procedure_id = + submitted_procedure_id(&frontend, "create_logical_tables", LOGICAL_TABLE).await; + + // Assert: logical Create emits one submitted row per logical table and + // preserves the logical and physical table IDs on its terminal row. + assert_named_submitted_event( + &frontend, + "create_logical_tables", + &create_logical_tables_procedure_id, + "\ ++-----------------------+--------------------------+ +| type | name | ++-----------------------+--------------------------+ +| create_logical_tables | table_ddl_events_logical | ++-----------------------+--------------------------+", + ) + .await; + assert_logical_terminal_event( + &frontend, + "create_logical_tables", + &create_logical_tables_procedure_id, + logical_table_id, + physical_table_id, + "Succeeded", + "\ ++-----------------------+--------------------------+ +| type | name | ++-----------------------+--------------------------+ +| create_logical_tables | table_ddl_events_logical | ++-----------------------+--------------------------+", + ) + .await; + + // Act / Assert: logical Alter preserves the one-row-per-logical-table submitted + // contract and emits a lightweight terminal row. + run_sql( + &frontend, + &format!("ALTER TABLE {LOGICAL_TABLE} ADD COLUMN rack STRING PRIMARY KEY"), + ) + .await; + let alter_logical_tables_procedure_id = + submitted_procedure_id(&frontend, "alter_logical_tables", LOGICAL_TABLE).await; + assert_named_submitted_event( + &frontend, + "alter_logical_tables", + &alter_logical_tables_procedure_id, + "\ ++----------------------+--------------------------+ +| type | name | ++----------------------+--------------------------+ +| alter_logical_tables | table_ddl_events_logical | ++----------------------+--------------------------+", + ) + .await; + assert_lightweight_terminal_event( + &frontend, + "alter_logical_tables", + &alter_logical_tables_procedure_id, + "Succeeded", + ) + .await; +} + +async fn run_sql(instance: &Arc, sql: &str) { + let output = instance + .do_query(sql, QueryContext::arc()) + .await + .remove(0) + .unwrap(); + assert!(matches!(output.data, OutputData::AffectedRows(_)), "{sql}"); +} + +async fn find_table_id(instance: &Arc, table_name: &str) -> u32 { + find_eventually_u32( + instance, + &format!( + "SELECT table_id FROM information_schema.tables WHERE table_catalog = 'greptime' AND table_schema = 'public' AND table_name = '{table_name}'" + ), + "table_id", + ) + .await +} + +async fn submitted_procedure_id( + instance: &Arc, + event_type: &str, + table_name: &str, +) -> String { + let query = format!( + "SELECT {} FROM {EVENTS_TABLE} WHERE {} = '{event_type}' AND {} = '{table_name}' AND {} = 'Submitted' AND json_path_match({}, '$.version == 1') ORDER BY timestamp DESC LIMIT 1", + PROCEDURE_ID_COLUMN.name(), + TYPE_COLUMN.name(), + TABLE_NAME_COLUMN.name(), + PROCEDURE_TRIGGER_COLUMN.name(), + PAYLOAD_COLUMN.name(), + ); + find_eventually_string(instance, &query, PROCEDURE_ID_COLUMN.name()).await +} + +async fn assert_named_submitted_event( + instance: &Arc, + event_type: &str, + procedure_id: &str, + expected: &str, +) { + let query = format!( + "SELECT {}, {} AS name FROM {EVENTS_TABLE} WHERE {} = '{event_type}' AND {} = '{procedure_id}' AND {} = 'Submitted' AND json_path_match({}, '$.version == 1')", + TYPE_COLUMN.name(), + TABLE_NAME_COLUMN.name(), + TYPE_COLUMN.name(), + PROCEDURE_ID_COLUMN.name(), + PROCEDURE_TRIGGER_COLUMN.name(), + PAYLOAD_COLUMN.name(), + ); + assert_eventually_eq(instance, &query, expected).await; +} + +async fn assert_id_submitted_event( + instance: &Arc, + event_type: &str, + procedure_id: &str, + expected: &str, +) { + let query = format!( + "SELECT {} FROM {EVENTS_TABLE} WHERE {} = '{event_type}' AND {} = '{procedure_id}' AND {} = 'Submitted' AND json_path_match({}, '$.version == 1')", + TYPE_COLUMN.name(), + TYPE_COLUMN.name(), + PROCEDURE_ID_COLUMN.name(), + PROCEDURE_TRIGGER_COLUMN.name(), + PAYLOAD_COLUMN.name(), + ); + assert_eventually_eq(instance, &query, expected).await; +} + +async fn assert_id_terminal_event( + instance: &Arc, + event_type: &str, + procedure_id: &str, + table_id: u32, + terminal_trigger: &str, + expected: &str, +) { + let query = format!( + "SELECT {} FROM {EVENTS_TABLE} WHERE {} = '{event_type}' AND {} = '{procedure_id}' AND {} = '{terminal_trigger}' AND {} = {table_id} AND json_is_null({}) AND {} IS NULL AND {} IS NULL AND {} IS NULL", + TYPE_COLUMN.name(), + TYPE_COLUMN.name(), + PROCEDURE_ID_COLUMN.name(), + PROCEDURE_TRIGGER_COLUMN.name(), + TABLE_ID_COLUMN.name(), + PAYLOAD_COLUMN.name(), + CATALOG_NAME_COLUMN.name(), + SCHEMA_NAME_COLUMN.name(), + TABLE_NAME_COLUMN.name(), + ); + assert_eventually_eq(instance, &query, expected).await; +} + +async fn assert_logical_terminal_event( + instance: &Arc, + event_type: &str, + procedure_id: &str, + table_id: u32, + physical_table_id: u32, + terminal_trigger: &str, + expected: &str, +) { + let query = format!( + "SELECT {}, {} AS name FROM {EVENTS_TABLE} WHERE {} = '{event_type}' AND {} = '{procedure_id}' AND {} = '{terminal_trigger}' AND {} = {table_id} AND {} = {physical_table_id} AND json_is_null({})", + TYPE_COLUMN.name(), + TABLE_NAME_COLUMN.name(), + TYPE_COLUMN.name(), + PROCEDURE_ID_COLUMN.name(), + PROCEDURE_TRIGGER_COLUMN.name(), + TABLE_ID_COLUMN.name(), + PHYSICAL_TABLE_ID_COLUMN.name(), + PAYLOAD_COLUMN.name(), + ); + assert_eventually_eq(instance, &query, expected).await; +} + +async fn assert_lightweight_terminal_event( + instance: &Arc, + event_type: &str, + procedure_id: &str, + terminal_trigger: &str, +) { + let query = format!( + "SELECT count(*) AS event_count FROM {EVENTS_TABLE} WHERE {} = '{event_type}' AND {} = '{procedure_id}' AND {} = '{terminal_trigger}' AND json_is_null({}) AND {} IS NULL AND {} IS NULL AND {} IS NULL AND {} IS NULL", + TYPE_COLUMN.name(), + PROCEDURE_ID_COLUMN.name(), + PROCEDURE_TRIGGER_COLUMN.name(), + PAYLOAD_COLUMN.name(), + CATALOG_NAME_COLUMN.name(), + SCHEMA_NAME_COLUMN.name(), + TABLE_NAME_COLUMN.name(), + TABLE_ID_COLUMN.name(), + ); + assert_single_event(instance, &query).await; +}