Skip to main content

common_meta/ddl/event/
table.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::any::Any;
16use std::collections::BTreeSet;
17
18use api::v1::alter_table_expr::Kind as AlterTableKind;
19use api::v1::value::ValueData;
20use api::v1::{ColumnSchema, Row};
21use common_event_recorder::Event;
22use common_event_recorder::error::{Result, SerializeEventSnafu};
23use common_event_recorder::event_table::{
24    CATALOG_NAME_COLUMN, PHYSICAL_TABLE_ID_COLUMN, SCHEMA_NAME_COLUMN, TABLE_ID_COLUMN,
25    TABLE_NAME_COLUMN, column_schemas, nullable_string, nullable_value,
26};
27use serde::Serialize;
28use serde_json::Value as JsonValue;
29use snafu::ResultExt;
30use store_api::storage::TableId;
31
32/// Current version of table DDL event payloads.
33pub(crate) const TABLE_DDL_PAYLOAD_VERSION: u8 = 1;
34
35/// A table DDL event type and its fixed domain schema.
36#[derive(Debug, Clone, Copy, PartialEq, Eq)]
37pub(crate) enum TableDdlEventType {
38    CreateTable,
39    CreateLogicalTables,
40    AlterTable,
41    AlterLogicalTables,
42    DropTable,
43    #[cfg(feature = "enterprise")]
44    UndropTable,
45    #[cfg(feature = "enterprise")]
46    PurgeDroppedTable,
47    TruncateTable,
48}
49
50impl TableDdlEventType {
51    /// Returns the stable event type stored in the events table.
52    pub(crate) const fn as_str(self) -> &'static str {
53        match self {
54            Self::CreateTable => "create_table",
55            Self::CreateLogicalTables => "create_logical_tables",
56            Self::AlterTable => "alter_table",
57            Self::AlterLogicalTables => "alter_logical_tables",
58            Self::DropTable => "drop_table",
59            #[cfg(feature = "enterprise")]
60            Self::UndropTable => "undrop_table",
61            #[cfg(feature = "enterprise")]
62            Self::PurgeDroppedTable => "purge_dropped_table",
63            Self::TruncateTable => "truncate_table",
64        }
65    }
66
67    const fn has_physical_table_id(self) -> bool {
68        matches!(self, Self::CreateLogicalTables | Self::AlterLogicalTables)
69    }
70}
71
72/// Nullable table locator columns stored alongside a table DDL event.
73#[derive(Debug, Clone, Default, PartialEq, Eq)]
74pub(crate) struct TableDdlLocator {
75    /// Catalog containing the table.
76    pub(crate) catalog_name: Option<String>,
77    /// Schema containing the table.
78    pub(crate) schema_name: Option<String>,
79    /// Table name.
80    pub(crate) table_name: Option<String>,
81    /// Table ID when known at this lifecycle point.
82    pub(crate) table_id: Option<TableId>,
83    /// Physical table ID for a logical table event.
84    pub(crate) physical_table_id: Option<TableId>,
85}
86
87impl TableDdlLocator {
88    /// Creates a locator from a fully qualified table name.
89    pub(crate) fn new(
90        catalog_name: impl Into<String>,
91        schema_name: impl Into<String>,
92        table_name: impl Into<String>,
93    ) -> Self {
94        Self {
95            catalog_name: Some(catalog_name.into()),
96            schema_name: Some(schema_name.into()),
97            table_name: Some(table_name.into()),
98            ..Default::default()
99        }
100    }
101
102    /// Creates a locator containing only a table ID.
103    #[cfg(feature = "enterprise")]
104    pub(crate) fn from_table_id(table_id: TableId) -> Self {
105        Self {
106            table_id: Some(table_id),
107            ..Default::default()
108        }
109    }
110
111    /// Adds a table ID to the locator.
112    pub(crate) fn with_table_id(mut self, table_id: TableId) -> Self {
113        self.table_id = Some(table_id);
114        self
115    }
116
117    /// Adds a physical table ID to a logical-table locator.
118    pub(crate) fn with_physical_table_id(mut self, physical_table_id: TableId) -> Self {
119        self.physical_table_id = Some(physical_table_id);
120        self
121    }
122}
123
124#[derive(Debug, Serialize)]
125#[serde(untagged)]
126enum TableDdlPayload {
127    CreateTable(CreateTablePayload),
128    CreateLogicalTables(CreateLogicalTablesPayload),
129    AlterTable(AlterTablePayload),
130    AlterLogicalTables(AlterLogicalTablesPayload),
131    DropTable(DropTablePayload),
132    #[cfg(feature = "enterprise")]
133    UndropTable(UndropTablePayload),
134    #[cfg(feature = "enterprise")]
135    PurgeDroppedTable(PurgeDroppedTablePayload),
136    TruncateTable(TruncateTablePayload),
137}
138
139#[derive(Debug, Serialize)]
140struct CreateTablePayload {
141    version: u8,
142    create_if_not_exists: bool,
143    engine: String,
144}
145
146#[derive(Debug, Serialize)]
147struct CreateLogicalTablesPayload {
148    version: u8,
149    table_count: usize,
150}
151
152#[derive(Debug, Serialize)]
153struct AlterTablePayload {
154    version: u8,
155    kind: Option<&'static str>,
156}
157
158#[derive(Debug, Serialize)]
159struct AlterLogicalTablesPayload {
160    version: u8,
161    table_count: usize,
162    kinds: Vec<&'static str>,
163}
164
165#[derive(Debug, Serialize)]
166struct DropTablePayload {
167    version: u8,
168    drop_if_exists: bool,
169}
170
171#[cfg(feature = "enterprise")]
172#[derive(Debug, Serialize)]
173struct UndropTablePayload {
174    version: u8,
175}
176
177#[cfg(feature = "enterprise")]
178#[derive(Debug, Serialize)]
179struct PurgeDroppedTablePayload {
180    version: u8,
181}
182
183#[derive(Debug, Serialize)]
184struct TruncateTablePayload {
185    version: u8,
186    time_range_count: usize,
187}
188
189/// Returns the stable kind stored in an Alter Table payload, if supported.
190pub(crate) fn alter_table_kind_name(kind: &AlterTableKind) -> Option<&'static str> {
191    match kind {
192        AlterTableKind::AddColumns(_) => Some("add_columns"),
193        AlterTableKind::DropColumns(_) => Some("drop_columns"),
194        AlterTableKind::RenameTable(_) => Some("rename_table"),
195        AlterTableKind::ModifyColumnTypes(_) => Some("modify_column_types"),
196        AlterTableKind::SetJsonSettings(_) => Some("set_json_settings"),
197        AlterTableKind::SetTableOptions(_) => Some("set_table_options"),
198        AlterTableKind::UnsetTableOptions(_) => Some("unset_table_options"),
199        AlterTableKind::SetIndex(_) => Some("set_index"),
200        AlterTableKind::UnsetIndex(_) => Some("unset_index"),
201        AlterTableKind::DropDefaults(_) => Some("drop_defaults"),
202        AlterTableKind::SetIndexes(_) => Some("set_indexes"),
203        AlterTableKind::UnsetIndexes(_) => Some("unset_indexes"),
204        AlterTableKind::SetDefaults(_) => Some("set_defaults"),
205        // Repartition is handled by RepartitionProcedure.
206        AlterTableKind::Repartition(_) => None,
207    }
208}
209
210/// Shared event representation used by table DDL procedures.
211#[derive(Debug)]
212pub(crate) struct TableDdlEvent {
213    event_type: TableDdlEventType,
214    locators: Vec<TableDdlLocator>,
215    payload: Option<TableDdlPayload>,
216}
217
218impl TableDdlEvent {
219    /// Builds the bounded event emitted when creating a table is submitted.
220    pub(crate) fn create_table_submitted(
221        locator: TableDdlLocator,
222        create_if_not_exists: bool,
223        engine: &str,
224    ) -> Self {
225        Self::submitted(
226            TableDdlEventType::CreateTable,
227            [locator],
228            TableDdlPayload::CreateTable(CreateTablePayload {
229                version: TABLE_DDL_PAYLOAD_VERSION,
230                create_if_not_exists,
231                engine: engine.to_string(),
232            }),
233        )
234    }
235
236    /// Builds the bounded event emitted when creating logical tables is submitted.
237    pub(crate) fn create_logical_tables_submitted(
238        locators: impl IntoIterator<Item = TableDdlLocator>,
239        table_count: usize,
240    ) -> Self {
241        Self::submitted(
242            TableDdlEventType::CreateLogicalTables,
243            locators,
244            TableDdlPayload::CreateLogicalTables(CreateLogicalTablesPayload {
245                version: TABLE_DDL_PAYLOAD_VERSION,
246                table_count,
247            }),
248        )
249    }
250
251    /// Builds the bounded event emitted when altering a table is submitted.
252    pub(crate) fn alter_table_submitted(
253        locator: TableDdlLocator,
254        kind: Option<&'static str>,
255    ) -> Self {
256        Self::submitted(
257            TableDdlEventType::AlterTable,
258            [locator],
259            TableDdlPayload::AlterTable(AlterTablePayload {
260                version: TABLE_DDL_PAYLOAD_VERSION,
261                kind,
262            }),
263        )
264    }
265
266    /// Builds the bounded event emitted when altering logical tables is submitted.
267    pub(crate) fn alter_logical_tables_submitted(
268        locators: impl IntoIterator<Item = TableDdlLocator>,
269        table_count: usize,
270        kinds: impl IntoIterator<Item = &'static str>,
271    ) -> Self {
272        let kinds = kinds
273            .into_iter()
274            .collect::<BTreeSet<_>>()
275            .into_iter()
276            .collect();
277        Self::submitted(
278            TableDdlEventType::AlterLogicalTables,
279            locators,
280            TableDdlPayload::AlterLogicalTables(AlterLogicalTablesPayload {
281                version: TABLE_DDL_PAYLOAD_VERSION,
282                table_count,
283                kinds,
284            }),
285        )
286    }
287
288    /// Builds the bounded event emitted when dropping a table is submitted.
289    pub(crate) fn drop_table_submitted(locator: TableDdlLocator, drop_if_exists: bool) -> Self {
290        Self::submitted(
291            TableDdlEventType::DropTable,
292            [locator],
293            TableDdlPayload::DropTable(DropTablePayload {
294                version: TABLE_DDL_PAYLOAD_VERSION,
295                drop_if_exists,
296            }),
297        )
298    }
299
300    /// Builds the bounded event emitted when restoring a dropped table is submitted.
301    #[cfg(feature = "enterprise")]
302    pub(crate) fn undrop_table_submitted(locator: TableDdlLocator) -> Self {
303        Self::submitted(
304            TableDdlEventType::UndropTable,
305            [locator],
306            TableDdlPayload::UndropTable(UndropTablePayload {
307                version: TABLE_DDL_PAYLOAD_VERSION,
308            }),
309        )
310    }
311
312    /// Builds the bounded event emitted when purging a dropped table is submitted.
313    #[cfg(feature = "enterprise")]
314    pub(crate) fn purge_dropped_table_submitted(locator: TableDdlLocator) -> Self {
315        Self::submitted(
316            TableDdlEventType::PurgeDroppedTable,
317            [locator],
318            TableDdlPayload::PurgeDroppedTable(PurgeDroppedTablePayload {
319                version: TABLE_DDL_PAYLOAD_VERSION,
320            }),
321        )
322    }
323
324    /// Builds the bounded event emitted when truncating a table is submitted.
325    pub(crate) fn truncate_table_submitted(
326        locator: TableDdlLocator,
327        time_range_count: usize,
328    ) -> Self {
329        Self::submitted(
330            TableDdlEventType::TruncateTable,
331            [locator],
332            TableDdlPayload::TruncateTable(TruncateTablePayload {
333                version: TABLE_DDL_PAYLOAD_VERSION,
334                time_range_count,
335            }),
336        )
337    }
338
339    /// Builds a lifecycle event with stable object locators and no intent payload.
340    pub(crate) fn lifecycle(
341        event_type: TableDdlEventType,
342        locators: impl IntoIterator<Item = TableDdlLocator>,
343    ) -> Self {
344        Self {
345            event_type,
346            locators: locators.into_iter().collect(),
347            payload: None,
348        }
349    }
350
351    /// Builds a Create Table success event containing the submitted locator and allocated ID.
352    pub(crate) fn create_table_succeeded(locator: TableDdlLocator, table_id: TableId) -> Self {
353        Self::lifecycle(
354            TableDdlEventType::CreateTable,
355            [locator.with_table_id(table_id)],
356        )
357    }
358
359    /// Builds Create Logical Tables success rows from their allocated locators.
360    pub(crate) fn create_logical_tables_succeeded(
361        locators: impl IntoIterator<Item = TableDdlLocator>,
362    ) -> Self {
363        Self::lifecycle(TableDdlEventType::CreateLogicalTables, locators)
364    }
365
366    fn submitted(
367        event_type: TableDdlEventType,
368        locators: impl IntoIterator<Item = TableDdlLocator>,
369        payload: TableDdlPayload,
370    ) -> Self {
371        Self {
372            event_type,
373            locators: locators.into_iter().collect(),
374            payload: Some(payload),
375        }
376    }
377
378    fn schema() -> Vec<ColumnSchema> {
379        column_schemas([
380            &CATALOG_NAME_COLUMN,
381            &SCHEMA_NAME_COLUMN,
382            &TABLE_NAME_COLUMN,
383            &TABLE_ID_COLUMN,
384        ])
385    }
386
387    fn locator_row(&self, locator: &TableDdlLocator) -> Row {
388        let mut values = vec![
389            nullable_string(locator.catalog_name.as_deref()),
390            nullable_string(locator.schema_name.as_deref()),
391            nullable_string(locator.table_name.as_deref()),
392            nullable_table_id(locator.table_id),
393        ];
394        if self.event_type.has_physical_table_id() {
395            values.push(nullable_table_id(locator.physical_table_id));
396        }
397        Row { values }
398    }
399}
400
401impl Event for TableDdlEvent {
402    fn event_type(&self) -> &str {
403        self.event_type.as_str()
404    }
405
406    fn json_payload(&self) -> Result<JsonValue> {
407        match &self.payload {
408            Some(payload) => serde_json::to_value(payload).context(SerializeEventSnafu),
409            None => Ok(JsonValue::Null),
410        }
411    }
412
413    fn extra_schema(&self) -> Vec<ColumnSchema> {
414        let mut schema = Self::schema();
415        if self.event_type.has_physical_table_id() {
416            schema.push(PHYSICAL_TABLE_ID_COLUMN.column_schema());
417        }
418        schema
419    }
420
421    fn extra_rows(&self) -> Result<Vec<Row>> {
422        Ok(self
423            .locators
424            .iter()
425            .map(|locator| self.locator_row(locator))
426            .collect())
427    }
428
429    fn as_any(&self) -> &dyn Any {
430        self
431    }
432}
433
434fn nullable_table_id(value: Option<TableId>) -> api::v1::Value {
435    nullable_value(value.map(ValueData::U32Value))
436}