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;
+}