Skip to main content

operator/statement/admin/
event.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 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
33/// Event type emitted for an ADMIN function execution.
34pub(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/// The ADMIN function input captured before execution.
50#[derive(Debug, Clone)]
51pub(crate) struct AdminFunctionEventInput {
52    actor: String,
53    function: String,
54    arguments: Option<JsonValue>,
55}
56
57impl AdminFunctionEventInput {
58    /// Creates an event input from an ADMIN function request.
59    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/// An event describing the outcome of an ADMIN function execution.
84#[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    /// Creates a successful ADMIN function event.
95    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    /// Creates a failed ADMIN function event with the debug representation of the error.
118    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    /// Creates a failed ADMIN function event for a cancelled execution.
136    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}