1use std::any::Any;
16
17use api::v1::value::ValueData;
18use api::v1::{ColumnSchema, Row};
19use common_event_recorder::Event;
20use common_event_recorder::error::Result as EventResult;
21use common_event_recorder::event_table::{
22 ACTOR_COLUMN, ADMIN_FUNCTION_NAME_COLUMN, ADMIN_FUNCTION_OUTPUT_COLUMN,
23 ADMIN_FUNCTION_STATUS_COLUMN, column_schemas, jsonb_value,
24};
25use datatypes::value::Value;
26use serde_json::{Value as JsonValue, json};
27use sql::ast::{Expr, FunctionArg, FunctionArgExpr, FunctionArguments, Value as SqlValue};
28use sql::statements::admin::Admin;
29
30use crate::error::Error;
31use crate::statement::admin::AdminFunctionRequest;
32
33pub(crate) const ADMIN_FUNCTION_EVENT_TYPE: &str = "admin_function";
35
36const PAYLOAD_VERSION: u8 = 1;
37const VERSION_KEY: &str = "version";
38const ARGUMENTS_KEY: &str = "arguments";
39const ARGUMENT_TYPE_KEY: &str = "type";
40const ARGUMENT_VALUE_KEY: &str = "value";
41const LIST_ARGUMENT_TYPE: &str = "list";
42const SUBQUERY_ARGUMENT_TYPE: &str = "subquery";
43const RESULT_KEY: &str = "result";
44const ERROR_KEY: &str = "error";
45const SUCCEEDED_STATUS: &str = "Succeeded";
46const FAILED_STATUS: &str = "Failed";
47const UNSUPPORTED: &str = "<unsupported>";
48
49#[derive(Debug, Clone)]
51pub(crate) struct AdminFunctionEventInput {
52 actor: String,
53 function: String,
54 arguments: Option<JsonValue>,
55}
56
57impl AdminFunctionEventInput {
58 pub(crate) fn from_request(request: &AdminFunctionRequest) -> Self {
60 let Admin::Func(function) = &request.statement;
61 let function_name = function.name.to_string().to_lowercase();
62
63 let arguments = match &function.args {
64 FunctionArguments::List(arguments) => Some(json!({
65 (ARGUMENT_TYPE_KEY): LIST_ARGUMENT_TYPE,
66 (ARGUMENT_VALUE_KEY): arguments.args.iter().map(argument_to_json).collect::<Vec<_>>(),
67 })),
68 FunctionArguments::Subquery(query) => Some(json!({
69 (ARGUMENT_TYPE_KEY): SUBQUERY_ARGUMENT_TYPE,
70 (ARGUMENT_VALUE_KEY): query.to_string(),
71 })),
72 FunctionArguments::None => None,
73 };
74
75 Self {
76 actor: request.query_ctx.current_user().username().to_string(),
77 function: function_name,
78 arguments,
79 }
80 }
81}
82
83#[derive(Debug)]
85pub(crate) struct AdminFunctionEvent {
86 actor: String,
87 function_name: String,
88 status: &'static str,
89 output: JsonValue,
90 payload: JsonValue,
91}
92
93impl AdminFunctionEvent {
94 pub(crate) fn success(input: AdminFunctionEventInput, result: Option<&Value>) -> Self {
96 let AdminFunctionEventInput {
97 actor,
98 function,
99 arguments,
100 } = input;
101 Self {
102 actor,
103 function_name: function,
104 status: SUCCEEDED_STATUS,
105 output: result.map_or_else(
106 || json!({}),
107 |result| {
108 json!({
109 (RESULT_KEY): value_to_json(result),
110 })
111 },
112 ),
113 payload: input_payload(arguments),
114 }
115 }
116
117 pub(crate) fn failure(input: AdminFunctionEventInput, error: &Error) -> Self {
119 let AdminFunctionEventInput {
120 actor,
121 function,
122 arguments,
123 } = input;
124 Self {
125 actor,
126 function_name: function,
127 status: FAILED_STATUS,
128 output: json!({
129 (ERROR_KEY): format!("{error:?}"),
130 }),
131 payload: input_payload(arguments),
132 }
133 }
134
135 pub(crate) fn cancelled(input: AdminFunctionEventInput) -> Self {
137 Self::failure(input, &Error::AdminFunctionCancelled)
138 }
139}
140
141impl Event for AdminFunctionEvent {
142 fn event_type(&self) -> &str {
143 ADMIN_FUNCTION_EVENT_TYPE
144 }
145
146 fn json_payload(&self) -> EventResult<JsonValue> {
147 Ok(self.payload.clone())
148 }
149
150 fn extra_schema(&self) -> Vec<ColumnSchema> {
151 column_schemas([
152 &ACTOR_COLUMN,
153 &ADMIN_FUNCTION_NAME_COLUMN,
154 &ADMIN_FUNCTION_STATUS_COLUMN,
155 &ADMIN_FUNCTION_OUTPUT_COLUMN,
156 ])
157 }
158
159 fn extra_rows(&self) -> EventResult<Vec<Row>> {
160 Ok(vec![Row {
161 values: vec![
162 ValueData::StringValue(self.actor.clone()).into(),
163 ValueData::StringValue(self.function_name.clone()).into(),
164 ValueData::StringValue(self.status.to_string()).into(),
165 jsonb_value(&self.output),
166 ],
167 }])
168 }
169
170 fn as_any(&self) -> &dyn Any {
171 self
172 }
173}
174
175fn input_payload(arguments: Option<JsonValue>) -> JsonValue {
176 let mut payload = json!({ VERSION_KEY: PAYLOAD_VERSION });
177 if let Some(arguments) = arguments {
178 payload[ARGUMENTS_KEY] = arguments;
179 }
180 payload
181}
182
183fn argument_to_json(argument: &FunctionArg) -> JsonValue {
184 match argument {
185 FunctionArg::Unnamed(FunctionArgExpr::Expr(Expr::Value(value))) => {
186 sql_value_to_json(&value.value)
187 }
188 _ => JsonValue::String(argument.to_string()),
189 }
190}
191
192fn sql_value_to_json(value: &SqlValue) -> JsonValue {
193 match value {
194 SqlValue::Number(value, _) => {
195 serde_json::from_str(value).unwrap_or_else(|_| JsonValue::String(value.clone()))
196 }
197 SqlValue::Boolean(value) => JsonValue::Bool(*value),
198 SqlValue::Null => JsonValue::Null,
199 SqlValue::SingleQuotedString(value) | SqlValue::DoubleQuotedString(value) => {
200 JsonValue::String(value.clone())
201 }
202 SqlValue::Placeholder(value) => JsonValue::String(value.clone()),
203 _ => value
204 .clone()
205 .into_string()
206 .map(JsonValue::String)
207 .unwrap_or_else(|| JsonValue::String(value.to_string())),
208 }
209}
210
211fn value_to_json(value: &Value) -> JsonValue {
212 match value {
213 Value::Float32(value) if !value.0.is_finite() => JsonValue::String(value.to_string()),
214 Value::Float64(value) if !value.0.is_finite() => JsonValue::String(value.to_string()),
215 _ => JsonValue::try_from(value.clone())
216 .unwrap_or_else(|_| JsonValue::String(UNSUPPORTED.to_string())),
217 }
218}
219
220#[cfg(test)]
221mod tests {
222 use api::v1::Row;
223 use api::v1::value::ValueData;
224 use common_event_recorder::Event;
225 use common_event_recorder::event_table::{
226 ACTOR_COLUMN, ADMIN_FUNCTION_NAME_COLUMN, ADMIN_FUNCTION_OUTPUT_COLUMN,
227 ADMIN_FUNCTION_STATUS_COLUMN, column_schemas, jsonb_value,
228 };
229 use datatypes::value::Value;
230 use serde_json::json;
231 use session::context::QueryContext;
232 use sql::ast::{
233 Expr, FunctionArg, FunctionArgExpr, FunctionArguments, Ident, Value as SqlValue,
234 };
235 use sql::dialect::GreptimeDbDialect;
236 use sql::parser::{ParseOptions, ParserContext};
237 use sql::statements::admin::Admin;
238 use sql::statements::statement::Statement;
239 use sqlparser::ast::{DollarQuotedString, FunctionArgOperator, FunctionArgumentList};
240
241 use crate::error::Error;
242 use crate::statement::admin::AdminFunctionRequest;
243 use crate::statement::admin::event::{
244 ADMIN_FUNCTION_EVENT_TYPE, AdminFunctionEvent, AdminFunctionEventInput, sql_value_to_json,
245 value_to_json,
246 };
247
248 fn request(sql: &str) -> AdminFunctionRequest {
249 let Statement::Admin(statement) =
250 ParserContext::create_with_dialect(sql, &GreptimeDbDialect {}, ParseOptions::default())
251 .unwrap()
252 .remove(0)
253 else {
254 panic!("expected ADMIN statement")
255 };
256 AdminFunctionRequest {
257 statement,
258 query_ctx: QueryContext::arc(),
259 }
260 }
261
262 fn request_with_arguments(arguments: FunctionArguments) -> AdminFunctionRequest {
263 let mut request = request("ADMIN plugin_function()");
264 let Admin::Func(function) = &mut request.statement;
265 function.args = arguments;
266 request
267 }
268
269 fn assert_admin_columns(
270 event: &AdminFunctionEvent,
271 function_name: &str,
272 status: &str,
273 output: serde_json::Value,
274 ) {
275 assert_eq!(event.event_type(), ADMIN_FUNCTION_EVENT_TYPE);
276 assert_eq!(
277 event.extra_schema(),
278 column_schemas([
279 &ACTOR_COLUMN,
280 &ADMIN_FUNCTION_NAME_COLUMN,
281 &ADMIN_FUNCTION_STATUS_COLUMN,
282 &ADMIN_FUNCTION_OUTPUT_COLUMN,
283 ])
284 );
285 assert_eq!(
286 event.extra_rows().unwrap(),
287 vec![Row {
288 values: vec![
289 ValueData::StringValue("greptime".to_string()).into(),
290 ValueData::StringValue(function_name.to_string()).into(),
291 ValueData::StringValue(status.to_string()).into(),
292 jsonb_value(&output),
293 ],
294 }]
295 );
296 }
297
298 #[test]
299 fn success_uses_separate_output_columns() {
300 let input = AdminFunctionEventInput::from_request(&request(
301 "ADMIN flush_table('greptime.public.demo')",
302 ));
303 let event = AdminFunctionEvent::success(input, Some(&Value::UInt64(3)));
304
305 assert_admin_columns(&event, "flush_table", "Succeeded", json!({"result": 3}));
306 assert_eq!(
307 event.json_payload().unwrap(),
308 json!({
309 "version": 1,
310 "arguments": {"type": "list", "value": ["greptime.public.demo"]},
311 })
312 );
313 }
314
315 #[test]
316 fn successful_null_is_explicit() {
317 let input = AdminFunctionEventInput::from_request(&request(
318 "ADMIN migrate_region(NULL, NULL, NULL)",
319 ));
320 let event = AdminFunctionEvent::success(input, Some(&Value::Null));
321
322 assert_admin_columns(
323 &event,
324 "migrate_region",
325 "Succeeded",
326 json!({"result": null}),
327 );
328 assert_eq!(
329 event.json_payload().unwrap(),
330 json!({
331 "version": 1,
332 "arguments": {"type": "list", "value": [null, null, null]},
333 })
334 );
335 }
336
337 #[test]
338 fn unknown_function_values_are_recorded() {
339 let input = AdminFunctionEventInput::from_request(&request(
340 "ADMIN plugin_function('plugin-value', NULL)",
341 ));
342 let event = AdminFunctionEvent::success(input, Some(&Value::from("plugin-result")));
343
344 assert_admin_columns(
345 &event,
346 "plugin_function",
347 "Succeeded",
348 json!({"result": "plugin-result"}),
349 );
350 assert_eq!(
351 event.json_payload().unwrap(),
352 json!({
353 "version": 1,
354 "arguments": {"type": "list", "value": ["plugin-value", null]},
355 })
356 );
357 }
358
359 #[test]
360 fn failure_payload_has_error_only() {
361 let input = AdminFunctionEventInput::from_request(&request(
362 "ADMIN flush_table('greptime.public.demo')",
363 ));
364 let error = Error::BuildAdminFunctionArgs {
365 msg: "invalid input".to_string(),
366 };
367 let expected_error = format!("{error:?}");
368 let event = AdminFunctionEvent::failure(input, &error);
369
370 assert_admin_columns(
371 &event,
372 "flush_table",
373 "Failed",
374 json!({"error": expected_error}),
375 );
376 assert_eq!(
377 event.json_payload().unwrap(),
378 json!({
379 "version": 1,
380 "arguments": {"type": "list", "value": ["greptime.public.demo"]},
381 })
382 );
383 }
384
385 #[test]
386 fn cancellation_uses_debug_error_format() {
387 let input = AdminFunctionEventInput::from_request(&request("ADMIN flush_table('demo')"));
388 let event = AdminFunctionEvent::cancelled(input);
389
390 assert_admin_columns(
391 &event,
392 "flush_table",
393 "Failed",
394 json!({"error": format!("{:?}", Error::AdminFunctionCancelled)}),
395 );
396 }
397
398 #[test]
399 fn subquery_arguments_use_display() {
400 let Statement::Query(query) = ParserContext::create_with_dialect(
401 "SELECT 1",
402 &GreptimeDbDialect {},
403 ParseOptions::default(),
404 )
405 .unwrap()
406 .remove(0) else {
407 panic!("expected query statement")
408 };
409 let input = AdminFunctionEventInput::from_request(&request_with_arguments(
410 FunctionArguments::Subquery(Box::new(query.inner)),
411 ));
412 let event = AdminFunctionEvent::failure(
413 input,
414 &Error::BuildAdminFunctionArgs {
415 msg: "subquery is not executable".to_string(),
416 },
417 );
418
419 assert_eq!(
420 event.json_payload().unwrap(),
421 json!({
422 "version": 1,
423 "arguments": {"type": "subquery", "value": "SELECT 1"},
424 })
425 );
426 }
427
428 #[test]
429 fn no_arguments_omits_arguments_field_and_empty_result() {
430 let input =
431 AdminFunctionEventInput::from_request(&request_with_arguments(FunctionArguments::None));
432 let event = AdminFunctionEvent::success(input, None);
433
434 assert_admin_columns(&event, "plugin_function", "Succeeded", json!({}));
435 assert_eq!(event.json_payload().unwrap(), json!({"version": 1}));
436 }
437
438 #[test]
439 fn empty_argument_list_is_recorded() {
440 let input = AdminFunctionEventInput::from_request(&request("ADMIN plugin_function()"));
441 let event = AdminFunctionEvent::success(input, Some(&Value::UInt64(0)));
442
443 assert_eq!(
444 event.json_payload().unwrap(),
445 json!({
446 "version": 1,
447 "arguments": {"type": "list", "value": []},
448 })
449 );
450 }
451
452 #[test]
453 fn records_extended_sql_values_and_non_literal_arguments() {
454 let input = AdminFunctionEventInput::from_request(&request(
455 "ADMIN plugin_function(1, true, NULL, X'48656c6c6f', table_name)",
456 ));
457 let event = AdminFunctionEvent::success(input, Some(&Value::UInt64(0)));
458
459 assert_eq!(
460 event.json_payload().unwrap(),
461 json!({
462 "version": 1,
463 "arguments": {
464 "type": "list",
465 "value": [1, true, null, "48656c6c6f", "table_name"],
466 },
467 })
468 );
469 }
470
471 #[test]
472 fn string_literals_preserve_logical_content() {
473 let values = [
474 SqlValue::SingleQuotedString("single".to_string()),
475 SqlValue::DoubleQuotedString("double".to_string()),
476 SqlValue::DollarQuotedString(DollarQuotedString {
477 value: "dollar".to_string(),
478 tag: None,
479 }),
480 SqlValue::SingleQuotedRawStringLiteral("raw".to_string()),
481 SqlValue::SingleQuotedByteStringLiteral("bytes".to_string()),
482 SqlValue::NationalStringLiteral("national".to_string()),
483 SqlValue::UnicodeStringLiteral("unicode".to_string()),
484 SqlValue::HexStringLiteral("48656c6c6f".to_string()),
485 ];
486
487 let expected = [
488 "single",
489 "double",
490 "dollar",
491 "raw",
492 "bytes",
493 "national",
494 "unicode",
495 "48656c6c6f",
496 ];
497 for (value, expected) in values.iter().zip(expected) {
498 assert_eq!(sql_value_to_json(value), json!(expected));
499 }
500 }
501
502 #[test]
503 fn non_literal_arguments_use_sql_display() {
504 let arguments = FunctionArgumentList {
505 duplicate_treatment: None,
506 args: vec![
507 FunctionArg::Named {
508 name: Ident::new("named"),
509 arg: FunctionArgExpr::Expr(Expr::Identifier(Ident::new("value"))),
510 operator: FunctionArgOperator::RightArrow,
511 },
512 FunctionArg::Unnamed(FunctionArgExpr::Wildcard),
513 ],
514 clauses: vec![],
515 };
516 let input = AdminFunctionEventInput::from_request(&request_with_arguments(
517 FunctionArguments::List(arguments),
518 ));
519 let event = AdminFunctionEvent::success(input, Some(&Value::UInt64(0)));
520
521 assert_eq!(
522 event.json_payload().unwrap(),
523 json!({
524 "version": 1,
525 "arguments": {"type": "list", "value": ["named => value", "*"]},
526 })
527 );
528 assert_eq!(
529 sql_value_to_json(&SqlValue::Placeholder("$1".to_string())),
530 json!("$1")
531 );
532 }
533
534 #[test]
535 fn serializes_binary_result() {
536 assert_eq!(
537 value_to_json(&Value::from(vec![1_u8, 2, 3])),
538 json!([1, 2, 3])
539 );
540 }
541
542 #[test]
543 fn serializes_non_finite_float_results_as_strings() {
544 assert_eq!(
545 value_to_json(&Value::Float32(f32::NAN.into())),
546 json!("NaN")
547 );
548 assert_eq!(
549 value_to_json(&Value::Float32(f32::INFINITY.into())),
550 json!("inf")
551 );
552 assert_eq!(
553 value_to_json(&Value::Float32(f32::NEG_INFINITY.into())),
554 json!("-inf")
555 );
556 assert_eq!(
557 value_to_json(&Value::Float64(f64::NAN.into())),
558 json!("NaN")
559 );
560 assert_eq!(
561 value_to_json(&Value::Float64(f64::INFINITY.into())),
562 json!("inf")
563 );
564 assert_eq!(
565 value_to_json(&Value::Float64(f64::NEG_INFINITY.into())),
566 json!("-inf")
567 );
568 }
569}