1use 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#[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 pub const fn name(&self) -> &'static str {
59 self.name
60 }
61
62 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
76pub const TYPE_COLUMN: EventTableColumn =
78 EventTableColumn::new("type", ColumnDataType::String, SemanticType::Tag);
79pub const PAYLOAD_COLUMN: EventTableColumn =
81 EventTableColumn::json_binary("payload", ColumnDataType::Binary, SemanticType::Field);
82pub const EVENT_CONTEXT_COLUMN: EventTableColumn =
84 EventTableColumn::json_binary("event_context", ColumnDataType::Binary, SemanticType::Field);
85pub const TIMESTAMP_COLUMN: EventTableColumn = EventTableColumn::new(
87 "timestamp",
88 ColumnDataType::TimestampNanosecond,
89 SemanticType::Timestamp,
90);
91pub const ACTOR_COLUMN: EventTableColumn =
93 EventTableColumn::new("actor", ColumnDataType::String, SemanticType::Field);
94pub const PROCEDURE_ID_COLUMN: EventTableColumn =
96 EventTableColumn::new("procedure_id", ColumnDataType::String, SemanticType::Field);
97pub const PROCEDURE_STATE_COLUMN: EventTableColumn = EventTableColumn::new(
99 "procedure_state",
100 ColumnDataType::String,
101 SemanticType::Field,
102);
103pub const PROCEDURE_ERROR_COLUMN: EventTableColumn = EventTableColumn::new(
105 "procedure_error",
106 ColumnDataType::String,
107 SemanticType::Field,
108);
109pub const PROCEDURE_TRIGGER_COLUMN: EventTableColumn = EventTableColumn::json_binary(
111 "procedure_trigger",
112 ColumnDataType::Binary,
113 SemanticType::Field,
114);
115pub const CATALOG_NAME_COLUMN: EventTableColumn =
117 EventTableColumn::new("catalog_name", ColumnDataType::String, SemanticType::Field);
118pub const SCHEMA_NAME_COLUMN: EventTableColumn =
120 EventTableColumn::new("schema_name", ColumnDataType::String, SemanticType::Field);
121pub const FLOW_NAME_COLUMN: EventTableColumn =
123 EventTableColumn::new("flow_name", ColumnDataType::String, SemanticType::Field);
124pub const FLOW_ID_COLUMN: EventTableColumn =
126 EventTableColumn::new("flow_id", ColumnDataType::Uint32, SemanticType::Field);
127pub const VIEW_NAME_COLUMN: EventTableColumn =
129 EventTableColumn::new("view_name", ColumnDataType::String, SemanticType::Field);
130pub const VIEW_ID_COLUMN: EventTableColumn =
132 EventTableColumn::new("view_id", ColumnDataType::Uint32, SemanticType::Field);
133pub const TOPIC_NAME_COLUMN: EventTableColumn =
135 EventTableColumn::new("topic_name", ColumnDataType::String, SemanticType::Field);
136pub const PRUNABLE_ENTRY_ID_COLUMN: EventTableColumn = EventTableColumn::new(
139 "prunable_entry_id",
140 ColumnDataType::Uint64,
141 SemanticType::Field,
142);
143pub const LATEST_OFFSET_COLUMN: EventTableColumn =
145 EventTableColumn::new("latest_offset", ColumnDataType::Uint64, SemanticType::Field);
146pub const GC_REPORT_COLUMN: EventTableColumn =
148 EventTableColumn::json_binary("gc_report", ColumnDataType::Binary, SemanticType::Field);
149
150pub const ADMIN_FUNCTION_NAME_COLUMN: EventTableColumn = EventTableColumn::new(
152 "admin_function_name",
153 ColumnDataType::String,
154 SemanticType::Field,
155);
156pub const ADMIN_FUNCTION_STATUS_COLUMN: EventTableColumn = EventTableColumn::new(
158 "admin_function_status",
159 ColumnDataType::String,
160 SemanticType::Field,
161);
162pub const ADMIN_FUNCTION_OUTPUT_COLUMN: EventTableColumn = EventTableColumn::json_binary(
164 "admin_function_output",
165 ColumnDataType::Binary,
166 SemanticType::Field,
167);
168
169pub const TABLE_NAME_COLUMN: EventTableColumn =
171 EventTableColumn::new("table_name", ColumnDataType::String, SemanticType::Field);
172pub const TABLE_ID_COLUMN: EventTableColumn =
174 EventTableColumn::new("table_id", ColumnDataType::Uint32, SemanticType::Field);
175pub const PHYSICAL_TABLE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
177 "physical_table_id",
178 ColumnDataType::Uint32,
179 SemanticType::Field,
180);
181pub const REGION_ID_COLUMN: EventTableColumn =
183 EventTableColumn::new("region_id", ColumnDataType::Uint64, SemanticType::Field);
184pub const REGION_NUMBER_COLUMN: EventTableColumn =
186 EventTableColumn::new("region_number", ColumnDataType::Uint32, SemanticType::Field);
187pub const REGION_MIGRATION_TRIGGER_REASON_COLUMN: EventTableColumn = EventTableColumn::new(
189 "region_migration_trigger_reason",
190 ColumnDataType::String,
191 SemanticType::Field,
192);
193pub const REGION_MIGRATION_SRC_NODE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
195 "region_migration_src_node_id",
196 ColumnDataType::Uint64,
197 SemanticType::Field,
198);
199pub const REGION_MIGRATION_SRC_PEER_ADDR_COLUMN: EventTableColumn = EventTableColumn::new(
201 "region_migration_src_peer_addr",
202 ColumnDataType::String,
203 SemanticType::Field,
204);
205pub const REGION_MIGRATION_DST_NODE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
207 "region_migration_dst_node_id",
208 ColumnDataType::Uint64,
209 SemanticType::Field,
210);
211pub const REGION_MIGRATION_DST_PEER_ADDR_COLUMN: EventTableColumn = EventTableColumn::new(
213 "region_migration_dst_peer_addr",
214 ColumnDataType::String,
215 SemanticType::Field,
216);
217pub const PARENT_PROCEDURE_ID_COLUMN: EventTableColumn = EventTableColumn::new(
219 "parent_procedure_id",
220 ColumnDataType::String,
221 SemanticType::Field,
222);
223pub const REPARTITION_GROUP_ID_COLUMN: EventTableColumn = EventTableColumn::new(
225 "repartition_group_id",
226 ColumnDataType::String,
227 SemanticType::Field,
228);
229pub const SOURCE_REGION_ID_COLUMN: EventTableColumn = EventTableColumn::new(
231 "source_region_id",
232 ColumnDataType::Uint64,
233 SemanticType::Field,
234);
235pub const SOURCE_REGION_NUMBER_COLUMN: EventTableColumn = EventTableColumn::new(
237 "source_region_number",
238 ColumnDataType::Uint32,
239 SemanticType::Field,
240);
241pub const SOURCE_PARTITION_EXPR_COLUMN: EventTableColumn = EventTableColumn::new(
243 "source_partition_expr",
244 ColumnDataType::String,
245 SemanticType::Field,
246);
247pub const TARGET_REGION_ID_COLUMN: EventTableColumn = EventTableColumn::new(
249 "target_region_id",
250 ColumnDataType::Uint64,
251 SemanticType::Field,
252);
253pub const TARGET_REGION_NUMBER_COLUMN: EventTableColumn = EventTableColumn::new(
255 "target_region_number",
256 ColumnDataType::Uint32,
257 SemanticType::Field,
258);
259pub const TARGET_PARTITION_EXPR_COLUMN: EventTableColumn = EventTableColumn::new(
261 "target_partition_expr",
262 ColumnDataType::String,
263 SemanticType::Field,
264);
265
266pub 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
276pub fn base_column_schemas() -> Vec<ColumnSchema> {
278 column_schemas([&TYPE_COLUMN, &PAYLOAD_COLUMN, &TIMESTAMP_COLUMN])
279}
280
281pub 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
291pub fn nullable_value(value: Option<ValueData>) -> Value {
293 Value { value_data: value }
294}
295
296pub 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
304pub 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
309pub fn jsonb_value(value: &serde_json::Value) -> Value {
311 ValueData::BinaryValue(jsonb::Value::from(value).to_vec()).into()
312}