1use 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
32pub(crate) const TABLE_DDL_PAYLOAD_VERSION: u8 = 1;
34
35#[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 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#[derive(Debug, Clone, Default, PartialEq, Eq)]
74pub(crate) struct TableDdlLocator {
75 pub(crate) catalog_name: Option<String>,
77 pub(crate) schema_name: Option<String>,
79 pub(crate) table_name: Option<String>,
81 pub(crate) table_id: Option<TableId>,
83 pub(crate) physical_table_id: Option<TableId>,
85}
86
87impl TableDdlLocator {
88 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 #[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 pub(crate) fn with_table_id(mut self, table_id: TableId) -> Self {
113 self.table_id = Some(table_id);
114 self
115 }
116
117 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
189pub(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 AlterTableKind::Repartition(_) => None,
207 }
208}
209
210#[derive(Debug)]
212pub(crate) struct TableDdlEvent {
213 event_type: TableDdlEventType,
214 locators: Vec<TableDdlLocator>,
215 payload: Option<TableDdlPayload>,
216}
217
218impl TableDdlEvent {
219 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 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 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 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 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 #[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 #[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 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 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 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 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}