1use std::any::Any;
16
17use api::v1::value::ValueData;
18use api::v1::{ColumnDataType, ColumnSchema, Row, SemanticType};
19use common_event_recorder::Event;
20use common_event_recorder::error::Result;
21use serde::Serialize;
22
23pub const SLOW_QUERY_TABLE_NAME: &str = "slow_queries";
24pub const SLOW_QUERY_TABLE_COST_COLUMN_NAME: &str = "cost";
25pub const SLOW_QUERY_TABLE_THRESHOLD_COLUMN_NAME: &str = "threshold";
26pub const SLOW_QUERY_TABLE_QUERY_COLUMN_NAME: &str = "query";
27pub const SLOW_QUERY_TABLE_SCHEMA_NAME_COLUMN_NAME: &str = "schema_name";
28pub const SLOW_QUERY_TABLE_TIMESTAMP_COLUMN_NAME: &str = "timestamp";
29pub const SLOW_QUERY_TABLE_IS_PROMQL_COLUMN_NAME: &str = "is_promql";
30pub const SLOW_QUERY_TABLE_PROMQL_START_COLUMN_NAME: &str = "promql_start";
31pub const SLOW_QUERY_TABLE_PROMQL_END_COLUMN_NAME: &str = "promql_end";
32pub const SLOW_QUERY_TABLE_PROMQL_RANGE_COLUMN_NAME: &str = "promql_range";
33pub const SLOW_QUERY_TABLE_PROMQL_STEP_COLUMN_NAME: &str = "promql_step";
34pub const SLOW_QUERY_EVENT_TYPE: &str = "slow_query";
35
36#[derive(Debug, Serialize)]
38pub struct SlowQueryEvent {
39 pub cost: u64,
40 pub threshold: u64,
41 pub query: String,
42 pub schema_name: String,
43 pub is_promql: bool,
44 pub promql_range: Option<u64>,
45 pub promql_step: Option<u64>,
46 pub promql_start: Option<i64>,
47 pub promql_end: Option<i64>,
48 pub payload: serde_json::Value,
49}
50
51impl Event for SlowQueryEvent {
52 fn table_name(&self) -> &str {
53 SLOW_QUERY_TABLE_NAME
54 }
55
56 fn event_type(&self) -> &str {
57 SLOW_QUERY_EVENT_TYPE
58 }
59
60 fn json_payload(&self) -> Result<serde_json::Value> {
61 Ok(self.payload.clone())
62 }
63
64 fn extra_schema(&self) -> Vec<ColumnSchema> {
65 vec![
66 ColumnSchema {
67 column_name: SLOW_QUERY_TABLE_COST_COLUMN_NAME.to_string(),
68 datatype: ColumnDataType::Uint64.into(),
69 semantic_type: SemanticType::Field.into(),
70 ..Default::default()
71 },
72 ColumnSchema {
73 column_name: SLOW_QUERY_TABLE_THRESHOLD_COLUMN_NAME.to_string(),
74 datatype: ColumnDataType::Uint64.into(),
75 semantic_type: SemanticType::Field.into(),
76 ..Default::default()
77 },
78 ColumnSchema {
79 column_name: SLOW_QUERY_TABLE_QUERY_COLUMN_NAME.to_string(),
80 datatype: ColumnDataType::String.into(),
81 semantic_type: SemanticType::Field.into(),
82 ..Default::default()
83 },
84 ColumnSchema {
85 column_name: SLOW_QUERY_TABLE_IS_PROMQL_COLUMN_NAME.to_string(),
86 datatype: ColumnDataType::Boolean.into(),
87 semantic_type: SemanticType::Field.into(),
88 ..Default::default()
89 },
90 ColumnSchema {
91 column_name: SLOW_QUERY_TABLE_PROMQL_RANGE_COLUMN_NAME.to_string(),
92 datatype: ColumnDataType::Uint64.into(),
93 semantic_type: SemanticType::Field.into(),
94 ..Default::default()
95 },
96 ColumnSchema {
97 column_name: SLOW_QUERY_TABLE_PROMQL_STEP_COLUMN_NAME.to_string(),
98 datatype: ColumnDataType::Uint64.into(),
99 semantic_type: SemanticType::Field.into(),
100 ..Default::default()
101 },
102 ColumnSchema {
103 column_name: SLOW_QUERY_TABLE_PROMQL_START_COLUMN_NAME.to_string(),
104 datatype: ColumnDataType::TimestampMillisecond.into(),
105 semantic_type: SemanticType::Field.into(),
106 ..Default::default()
107 },
108 ColumnSchema {
109 column_name: SLOW_QUERY_TABLE_PROMQL_END_COLUMN_NAME.to_string(),
110 datatype: ColumnDataType::TimestampMillisecond.into(),
111 semantic_type: SemanticType::Field.into(),
112 ..Default::default()
113 },
114 ColumnSchema {
115 column_name: SLOW_QUERY_TABLE_SCHEMA_NAME_COLUMN_NAME.to_string(),
116 datatype: ColumnDataType::String.into(),
117 semantic_type: SemanticType::Field.into(),
118 ..Default::default()
119 },
120 ]
121 }
122
123 fn extra_rows(&self) -> Result<Vec<Row>> {
124 Ok(vec![Row {
125 values: vec![
126 ValueData::U64Value(self.cost).into(),
127 ValueData::U64Value(self.threshold).into(),
128 ValueData::StringValue(self.query.to_string()).into(),
129 ValueData::BoolValue(self.is_promql).into(),
130 ValueData::U64Value(self.promql_range.unwrap_or(0)).into(),
131 ValueData::U64Value(self.promql_step.unwrap_or(0)).into(),
132 ValueData::TimestampMillisecondValue(self.promql_start.unwrap_or(0)).into(),
133 ValueData::TimestampMillisecondValue(self.promql_end.unwrap_or(0)).into(),
134 ValueData::StringValue(self.schema_name.clone()).into(),
135 ],
136 }])
137 }
138
139 fn as_any(&self) -> &dyn Any {
140 self
141 }
142}
143
144#[cfg(test)]
145mod tests {
146 use api::v1::value::ValueData;
147 use common_event_recorder::Event;
148
149 use super::*;
150
151 #[test]
152 fn slow_query_event_includes_schema() {
153 let event = SlowQueryEvent {
154 cost: 100,
155 threshold: 10,
156 query: "SELECT * FROM numbers".to_string(),
157 schema_name: "public".to_string(),
158 is_promql: false,
159 promql_range: None,
160 promql_step: None,
161 promql_start: None,
162 promql_end: None,
163 payload: serde_json::Value::Null,
164 };
165
166 let schema = event.extra_schema();
167 let column_names = schema
168 .iter()
169 .map(|column| column.column_name.as_str())
170 .collect::<Vec<_>>();
171 assert_eq!(
172 column_names,
173 vec![
174 SLOW_QUERY_TABLE_COST_COLUMN_NAME,
175 SLOW_QUERY_TABLE_THRESHOLD_COLUMN_NAME,
176 SLOW_QUERY_TABLE_QUERY_COLUMN_NAME,
177 SLOW_QUERY_TABLE_IS_PROMQL_COLUMN_NAME,
178 SLOW_QUERY_TABLE_PROMQL_RANGE_COLUMN_NAME,
179 SLOW_QUERY_TABLE_PROMQL_STEP_COLUMN_NAME,
180 SLOW_QUERY_TABLE_PROMQL_START_COLUMN_NAME,
181 SLOW_QUERY_TABLE_PROMQL_END_COLUMN_NAME,
182 SLOW_QUERY_TABLE_SCHEMA_NAME_COLUMN_NAME,
183 ]
184 );
185 assert_eq!(schema[8].semantic_type, SemanticType::Field as i32);
186 assert_eq!(event.json_payload().unwrap(), serde_json::Value::Null);
187
188 let rows = event.extra_rows().unwrap();
189 assert_eq!(rows.len(), 1);
190 assert_eq!(
191 rows[0].values[8].value_data,
192 Some(ValueData::StringValue("public".to_string()))
193 );
194 }
195
196 #[test]
197 fn slow_query_event_includes_timeout_payload() {
198 let payload = serde_json::json!({
199 "timed_out": true,
200 "metrics": [{"stage": 0}],
201 });
202 let event = SlowQueryEvent {
203 cost: 100,
204 threshold: 10,
205 query: "EXPLAIN ANALYZE VERBOSE SELECT 1".to_string(),
206 schema_name: "public".to_string(),
207 is_promql: false,
208 promql_range: None,
209 promql_step: None,
210 promql_start: None,
211 promql_end: None,
212 payload: payload.clone(),
213 };
214
215 assert_eq!(event.json_payload().unwrap(), payload);
216 }
217}