Skip to main content

common_event_recorder/
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 api::v1::column_data_type_extension::TypeExt;
16use api::v1::value::ValueData;
17use api::v1::{
18    ColumnDataType, ColumnDataTypeExtension, ColumnSchema, JsonTypeExtension, SemanticType, Value,
19};
20
21/// A canonical column in the shared event table.
22#[derive(Debug, Clone, Copy, PartialEq, Eq)]
23pub struct EventTableColumn {
24    name: &'static str,
25    datatype: ColumnDataType,
26    semantic_type: SemanticType,
27    json_binary: bool,
28}
29
30impl EventTableColumn {
31    const fn new(
32        name: &'static str,
33        datatype: ColumnDataType,
34        semantic_type: SemanticType,
35    ) -> Self {
36        Self {
37            name,
38            datatype,
39            semantic_type,
40            json_binary: false,
41        }
42    }
43
44    const fn json_binary(
45        name: &'static str,
46        datatype: ColumnDataType,
47        semantic_type: SemanticType,
48    ) -> Self {
49        Self {
50            name,
51            datatype,
52            semantic_type,
53            json_binary: true,
54        }
55    }
56
57    /// Returns the canonical column name.
58    pub const fn name(&self) -> &'static str {
59        self.name
60    }
61
62    /// Builds the canonical API schema for this column.
63    pub fn column_schema(&self) -> ColumnSchema {
64        ColumnSchema {
65            column_name: self.name.to_string(),
66            datatype: self.datatype.into(),
67            semantic_type: self.semantic_type.into(),
68            datatype_extension: self.json_binary.then(|| ColumnDataTypeExtension {
69                type_ext: Some(TypeExt::JsonType(JsonTypeExtension::JsonBinary.into())),
70            }),
71            ..Default::default()
72        }
73    }
74}
75
76/// The canonical event type column.
77pub const TYPE_COLUMN: EventTableColumn =
78    EventTableColumn::new("type", ColumnDataType::String, SemanticType::Tag);
79/// The canonical event payload column.
80pub const PAYLOAD_COLUMN: EventTableColumn =
81    EventTableColumn::json_binary("payload", ColumnDataType::Binary, SemanticType::Field);
82/// The optional context describing why an event was triggered.
83pub const EVENT_CONTEXT_COLUMN: EventTableColumn =
84    EventTableColumn::json_binary("event_context", ColumnDataType::Binary, SemanticType::Field);
85/// The canonical event timestamp column.
86pub const TIMESTAMP_COLUMN: EventTableColumn = EventTableColumn::new(
87    "timestamp",
88    ColumnDataType::TimestampNanosecond,
89    SemanticType::Timestamp,
90);
91/// The effective database user that initiated an event.
92pub const ACTOR_COLUMN: EventTableColumn =
93    EventTableColumn::new("actor", ColumnDataType::String, SemanticType::Field);
94/// The canonical procedure identifier envelope column.
95pub const PROCEDURE_ID_COLUMN: EventTableColumn =
96    EventTableColumn::new("procedure_id", ColumnDataType::String, SemanticType::Field);
97/// The canonical procedure state envelope column.
98pub const PROCEDURE_STATE_COLUMN: EventTableColumn = EventTableColumn::new(
99    "procedure_state",
100    ColumnDataType::String,
101    SemanticType::Field,
102);
103/// The canonical procedure error envelope column.
104pub const PROCEDURE_ERROR_COLUMN: EventTableColumn = EventTableColumn::new(
105    "procedure_error",
106    ColumnDataType::String,
107    SemanticType::Field,
108);
109/// The canonical procedure trigger envelope column.
110pub const PROCEDURE_TRIGGER_COLUMN: EventTableColumn = EventTableColumn::json_binary(
111    "procedure_trigger",
112    ColumnDataType::Binary,
113    SemanticType::Field,
114);
115/// The canonical catalog name dimension.
116pub const CATALOG_NAME_COLUMN: EventTableColumn =
117    EventTableColumn::new("catalog_name", ColumnDataType::String, SemanticType::Field);
118/// The canonical schema name dimension.
119pub const SCHEMA_NAME_COLUMN: EventTableColumn =
120    EventTableColumn::new("schema_name", ColumnDataType::String, SemanticType::Field);
121/// The canonical Flow name dimension.
122pub const FLOW_NAME_COLUMN: EventTableColumn =
123    EventTableColumn::new("flow_name", ColumnDataType::String, SemanticType::Field);
124/// The canonical Flow identifier dimension.
125pub const FLOW_ID_COLUMN: EventTableColumn =
126    EventTableColumn::new("flow_id", ColumnDataType::Uint32, SemanticType::Field);
127/// The canonical View name dimension.
128pub const VIEW_NAME_COLUMN: EventTableColumn =
129    EventTableColumn::new("view_name", ColumnDataType::String, SemanticType::Field);
130/// The canonical View identifier dimension.
131pub const VIEW_ID_COLUMN: EventTableColumn =
132    EventTableColumn::new("view_id", ColumnDataType::Uint32, SemanticType::Field);
133/// The canonical Kafka topic name dimension.
134pub const TOPIC_NAME_COLUMN: EventTableColumn =
135    EventTableColumn::new("topic_name", ColumnDataType::String, SemanticType::Field);
136/// The requested WAL prune boundary. It is only an attempted boundary on non-`Succeeded`
137/// procedure events.
138pub const PRUNABLE_ENTRY_ID_COLUMN: EventTableColumn = EventTableColumn::new(
139    "prunable_entry_id",
140    ColumnDataType::Uint64,
141    SemanticType::Field,
142);
143/// The canonical Kafka latest offset, which is an exclusive upper bound.
144pub const LATEST_OFFSET_COLUMN: EventTableColumn =
145    EventTableColumn::new("latest_offset", ColumnDataType::Uint64, SemanticType::Field);
146/// The canonical per-region GC report field.
147pub const GC_REPORT_COLUMN: EventTableColumn =
148    EventTableColumn::json_binary("gc_report", ColumnDataType::Binary, SemanticType::Field);
149
150/// The canonical ADMIN function name field.
151pub const ADMIN_FUNCTION_NAME_COLUMN: EventTableColumn = EventTableColumn::new(
152    "admin_function_name",
153    ColumnDataType::String,
154    SemanticType::Field,
155);
156/// The canonical ADMIN function execution status field.
157pub const ADMIN_FUNCTION_STATUS_COLUMN: EventTableColumn = EventTableColumn::new(
158    "admin_function_status",
159    ColumnDataType::String,
160    SemanticType::Field,
161);
162/// The canonical ADMIN function output field.
163pub const ADMIN_FUNCTION_OUTPUT_COLUMN: EventTableColumn = EventTableColumn::json_binary(
164    "admin_function_output",
165    ColumnDataType::Binary,
166    SemanticType::Field,
167);
168
169/// The canonical table name field for table DDL events.
170pub const TABLE_NAME_COLUMN: EventTableColumn =
171    EventTableColumn::new("table_name", ColumnDataType::String, SemanticType::Field);
172/// The canonical table identifier field for table DDL events.
173pub const TABLE_ID_COLUMN: EventTableColumn =
174    EventTableColumn::new("table_id", ColumnDataType::Uint32, SemanticType::Field);
175/// The canonical physical table identifier dimension.
176pub const PHYSICAL_TABLE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
177    "physical_table_id",
178    ColumnDataType::Uint32,
179    SemanticType::Field,
180);
181/// The canonical region identifier field for region events.
182pub const REGION_ID_COLUMN: EventTableColumn =
183    EventTableColumn::new("region_id", ColumnDataType::Uint64, SemanticType::Field);
184/// The canonical region number field for region events.
185pub const REGION_NUMBER_COLUMN: EventTableColumn =
186    EventTableColumn::new("region_number", ColumnDataType::Uint32, SemanticType::Field);
187/// The canonical region migration trigger reason field.
188pub const REGION_MIGRATION_TRIGGER_REASON_COLUMN: EventTableColumn = EventTableColumn::new(
189    "region_migration_trigger_reason",
190    ColumnDataType::String,
191    SemanticType::Field,
192);
193/// The canonical region migration source node identifier field.
194pub const REGION_MIGRATION_SRC_NODE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
195    "region_migration_src_node_id",
196    ColumnDataType::Uint64,
197    SemanticType::Field,
198);
199/// The canonical region migration source peer address field.
200pub const REGION_MIGRATION_SRC_PEER_ADDR_COLUMN: EventTableColumn = EventTableColumn::new(
201    "region_migration_src_peer_addr",
202    ColumnDataType::String,
203    SemanticType::Field,
204);
205/// The canonical region migration destination node identifier field.
206pub const REGION_MIGRATION_DST_NODE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
207    "region_migration_dst_node_id",
208    ColumnDataType::Uint64,
209    SemanticType::Field,
210);
211/// The canonical region migration destination peer address field.
212pub const REGION_MIGRATION_DST_PEER_ADDR_COLUMN: EventTableColumn = EventTableColumn::new(
213    "region_migration_dst_peer_addr",
214    ColumnDataType::String,
215    SemanticType::Field,
216);
217/// The canonical parent procedure identifier field for child procedure events.
218pub const PARENT_PROCEDURE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
219    "parent_procedure_id",
220    ColumnDataType::String,
221    SemanticType::Field,
222);
223/// The canonical repartition group identifier field.
224pub const REPARTITION_GROUP_ID_COLUMN: EventTableColumn = EventTableColumn::new(
225    "repartition_group_id",
226    ColumnDataType::String,
227    SemanticType::Field,
228);
229/// The canonical repartition source region identifier field.
230pub const SOURCE_REGION_ID_COLUMN: EventTableColumn = EventTableColumn::new(
231    "source_region_id",
232    ColumnDataType::Uint64,
233    SemanticType::Field,
234);
235/// The canonical repartition source region number field.
236pub const SOURCE_REGION_NUMBER_COLUMN: EventTableColumn = EventTableColumn::new(
237    "source_region_number",
238    ColumnDataType::Uint32,
239    SemanticType::Field,
240);
241/// The canonical repartition source partition expression field.
242pub const SOURCE_PARTITION_EXPR_COLUMN: EventTableColumn = EventTableColumn::new(
243    "source_partition_expr",
244    ColumnDataType::String,
245    SemanticType::Field,
246);
247/// The canonical repartition target region identifier field.
248pub const TARGET_REGION_ID_COLUMN: EventTableColumn = EventTableColumn::new(
249    "target_region_id",
250    ColumnDataType::Uint64,
251    SemanticType::Field,
252);
253/// The canonical repartition target region number field.
254pub const TARGET_REGION_NUMBER_COLUMN: EventTableColumn = EventTableColumn::new(
255    "target_region_number",
256    ColumnDataType::Uint32,
257    SemanticType::Field,
258);
259/// The canonical repartition target partition expression field.
260pub const TARGET_PARTITION_EXPR_COLUMN: EventTableColumn = EventTableColumn::new(
261    "target_partition_expr",
262    ColumnDataType::String,
263    SemanticType::Field,
264);
265
266/// Builds API schemas from canonical event-table columns while preserving their order.
267pub fn column_schemas<'a>(
268    columns: impl IntoIterator<Item = &'a EventTableColumn>,
269) -> Vec<ColumnSchema> {
270    columns
271        .into_iter()
272        .map(EventTableColumn::column_schema)
273        .collect()
274}
275
276/// Builds the canonical base schema for every recorded event.
277pub fn base_column_schemas() -> Vec<ColumnSchema> {
278    column_schemas([&TYPE_COLUMN, &PAYLOAD_COLUMN, &TIMESTAMP_COLUMN])
279}
280
281/// Builds the canonical procedure event envelope schema.
282pub fn procedure_event_column_schemas() -> Vec<ColumnSchema> {
283    column_schemas([
284        &PROCEDURE_ID_COLUMN,
285        &PROCEDURE_STATE_COLUMN,
286        &PROCEDURE_ERROR_COLUMN,
287        &PROCEDURE_TRIGGER_COLUMN,
288    ])
289}
290
291/// Builds an API value from an optional typed value.
292pub fn nullable_value(value: Option<ValueData>) -> Value {
293    Value { value_data: value }
294}
295
296/// Builds a nullable API string value.
297pub fn nullable_string<T>(value: Option<T>) -> Value
298where
299    T: AsRef<str>,
300{
301    nullable_value(value.map(|value| ValueData::StringValue(value.as_ref().to_string())))
302}
303
304/// Builds a nullable API JSONB value.
305pub fn nullable_json(value: Option<&serde_json::Value>) -> Value {
306    nullable_value(value.map(|value| ValueData::BinaryValue(jsonb::Value::from(value).to_vec())))
307}
308
309/// Builds a JSONB API value.
310pub fn jsonb_value(value: &serde_json::Value) -> Value {
311    ValueData::BinaryValue(jsonb::Value::from(value).to_vec()).into()
312}