Skip to main content

frontend/instance/
grpc.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::pin::Pin;
16use std::sync::Arc;
17use std::time::Instant;
18
19use api::helper::from_pb_time_ranges;
20use api::v1::ddl_request::{Expr as DdlExpr, Expr};
21use api::v1::greptime_request::Request;
22use api::v1::query_request::Query;
23use api::v1::{
24    DeleteRequests, DropFlowExpr, InsertIntoPlan, InsertRequests, RowDeleteRequests,
25    RowInsertRequests,
26};
27use async_stream::try_stream;
28use async_trait::async_trait;
29use auth::{
30    PermissionChecker, PermissionCheckerRef, PermissionReq, PermissionResp, PermissionTableTargets,
31};
32use common_error::ext::BoxedError;
33use common_grpc::flight::do_put::DoPutResponse;
34use common_query::Output;
35use common_query::logical_plan::add_insert_to_logical_plan;
36use common_telemetry::tracing::{self};
37use datafusion::datasource::DefaultTableSource;
38use futures::Stream;
39use futures::stream::StreamExt;
40use query::parser::PromQuery;
41use servers::error as server_error;
42use servers::http::prom_store::PHYSICAL_TABLE_PARAM;
43use servers::interceptor::{GrpcQueryInterceptor, GrpcQueryInterceptorRef};
44use servers::query_handler::grpc::GrpcQueryHandler;
45use session::context::QueryContextRef;
46use snafu::{OptionExt, ResultExt, ensure};
47use table::TableRef;
48use table::table::adapter::DfTableProviderAdapter;
49use table::table_name::TableName;
50
51use crate::error::{
52    CatalogSnafu, DataFusionSnafu, Error, ExternalSnafu, IncompleteGrpcRequestSnafu,
53    NotSupportedSnafu, PermissionSnafu, PlanStatementSnafu, Result,
54    SubstraitDecodeLogicalPlanSnafu, TableNotFoundSnafu, TableOperationSnafu,
55};
56use crate::instance::{Instance, attach_timer};
57use crate::metrics::{
58    GRPC_HANDLE_PLAN_ELAPSED, GRPC_HANDLE_PROMQL_ELAPSED, GRPC_HANDLE_SQL_ELAPSED,
59};
60
61#[async_trait]
62impl GrpcQueryHandler for Instance {
63    async fn do_query(
64        &self,
65        request: Request,
66        ctx: QueryContextRef,
67    ) -> server_error::Result<Output> {
68        let result: Result<Output> = async {
69            let interceptor_ref = self.plugins.get::<GrpcQueryInterceptorRef<Error>>();
70            let interceptor = interceptor_ref.as_ref();
71            interceptor.pre_execute(&request, ctx.clone())?;
72
73            self.plugins
74                .get::<PermissionCheckerRef>()
75                .as_ref()
76                .check_permission_with_context(
77                    ctx.current_user(),
78                    PermissionReq::GrpcRequest(&request),
79                    Some(&ctx.current_schema()),
80                )
81                .context(PermissionSnafu)?;
82
83            let output = match request {
84                Request::Inserts(requests) => self.handle_inserts(requests, ctx.clone()).await?,
85                Request::RowInserts(requests) => match ctx.extension(PHYSICAL_TABLE_PARAM) {
86                    Some(physical_table) => {
87                        self.handle_metric_row_inserts(
88                            requests,
89                            ctx.clone(),
90                            physical_table.to_string(),
91                        )
92                        .await?
93                    }
94                    None => {
95                        self.handle_row_inserts(requests, ctx.clone(), false, false)
96                            .await?
97                    }
98                },
99                Request::Deletes(requests) => self.handle_deletes(requests, ctx.clone()).await?,
100                Request::RowDeletes(requests) => self.handle_row_deletes(requests, ctx.clone()).await?,
101                Request::Query(query_request) => {
102                    let query = query_request.query.context(IncompleteGrpcRequestSnafu {
103                        err_msg: "Missing field 'QueryRequest.query'",
104                    })?;
105                    match query {
106                        Query::Sql(sql) => {
107                            let timer = GRPC_HANDLE_SQL_ELAPSED.start_timer();
108                            let mut result = self.do_query_inner(&sql, ctx.clone()).await;
109                            ensure!(
110                                result.len() == 1,
111                                NotSupportedSnafu {
112                                    feat: "execute multiple statements in SQL query string through GRPC interface"
113                                }
114                            );
115                            let output = result.remove(0)?;
116                            attach_timer(output, timer)
117                        }
118                        Query::LogicalPlan(plan) => {
119                            // this path is useful internally when flownode needs to execute a logical plan through gRPC interface
120                            let timer = GRPC_HANDLE_PLAN_ELAPSED.start_timer();
121
122                            // use dummy catalog to provide table
123                            let plan_decoder = self
124                                .query_engine()
125                                .engine_context(ctx.clone())
126                                .new_plan_decoder()
127                                .context(PlanStatementSnafu)?;
128
129                            let dummy_catalog_list =
130                                Arc::new(catalog::table_source::dummy_catalog::DummyCatalogList::new_with_query_ctx(
131                                    self.catalog_manager().clone(),
132                                    ctx.clone(),
133                                ));
134
135                            let logical_plan = plan_decoder
136                                .decode(bytes::Bytes::from(plan), dummy_catalog_list, true)
137                                .await
138                                .context(SubstraitDecodeLogicalPlanSnafu)?;
139                            let output =
140                                self.do_exec_plan_inner(logical_plan, None, ctx.clone()).await?;
141
142                            attach_timer(output, timer)
143                        }
144                        Query::InsertIntoPlan(insert) => {
145                            self.handle_insert_plan(insert, ctx.clone()).await?
146                        }
147                        Query::PromRangeQuery(promql) => {
148                            let timer = GRPC_HANDLE_PROMQL_ELAPSED.start_timer();
149                            let prom_query = PromQuery {
150                                query: promql.query,
151                                start: promql.start,
152                                end: promql.end,
153                                step: promql.step,
154                                lookback: promql.lookback,
155                                alias: None,
156                            };
157                            let mut result =
158                                self.do_promql_query_inner(&prom_query, ctx.clone()).await;
159                            ensure!(
160                                result.len() == 1,
161                                NotSupportedSnafu {
162                                    feat: "execute multiple statements in PromQL query string through GRPC interface"
163                                }
164                            );
165                            let output = result.remove(0)?;
166                            attach_timer(output, timer)
167                        }
168                    }
169                }
170                Request::Ddl(request) => {
171                    let mut expr = request.expr.context(IncompleteGrpcRequestSnafu {
172                        err_msg: "'expr' is absent in DDL request",
173                    })?;
174
175                    fill_catalog_and_schema_from_context(&mut expr, &ctx);
176
177                    match expr {
178                        DdlExpr::CreateTable(mut expr) => {
179                            let _ = self
180                                .statement_executor
181                                .create_table_inner(&mut expr, None, ctx.clone())
182                                .await?;
183                            Output::new_with_affected_rows(0)
184                        }
185                        DdlExpr::AlterDatabase(expr) => {
186                            let _ = self
187                                .statement_executor
188                                .alter_database_inner(expr, ctx.clone())
189                                .await?;
190                            Output::new_with_affected_rows(0)
191                        }
192                        DdlExpr::AlterTable(expr) => {
193                            self.statement_executor
194                                .alter_table_inner(expr, ctx.clone())
195                                .await?
196                        }
197                        DdlExpr::CreateDatabase(expr) => {
198                            self.statement_executor
199                                .create_database(
200                                    &expr.schema_name,
201                                    expr.create_if_not_exists,
202                                    expr.options,
203                                    ctx.clone(),
204                                )
205                                .await?
206                        }
207                        DdlExpr::DropTable(expr) => {
208                            let table_name =
209                                TableName::new(&expr.catalog_name, &expr.schema_name, &expr.table_name);
210                            self.statement_executor
211                                .drop_table(table_name, expr.drop_if_exists, ctx.clone())
212                                .await?
213                        }
214                        DdlExpr::TruncateTable(expr) => {
215                            let table_name =
216                                TableName::new(&expr.catalog_name, &expr.schema_name, &expr.table_name);
217                            let time_ranges = from_pb_time_ranges(expr.time_ranges.unwrap_or_default())
218                                .map_err(BoxedError::new)
219                                .context(ExternalSnafu)?;
220                            self.statement_executor
221                                .truncate_table(table_name, time_ranges, ctx.clone())
222                                .await?
223                        }
224                        DdlExpr::CreateFlow(expr) => {
225                            self.statement_executor
226                                .create_flow_inner(expr, ctx.clone())
227                                .await?
228                        }
229                        DdlExpr::DropFlow(DropFlowExpr {
230                            catalog_name,
231                            flow_name,
232                            drop_if_exists,
233                            ..
234                        }) => {
235                            self.statement_executor
236                                .drop_flow(catalog_name, flow_name, drop_if_exists, ctx.clone())
237                                .await?
238                        }
239                        DdlExpr::CreateView(expr) => {
240                            let _ = self
241                                .statement_executor
242                                .create_view_by_expr(expr, ctx.clone())
243                                .await?;
244
245                            Output::new_with_affected_rows(0)
246                        }
247                        DdlExpr::DropView(_) => {
248                            todo!("implemented in the following PR")
249                        }
250                        DdlExpr::CommentOn(expr) => {
251                            self.statement_executor
252                                .comment_by_expr(expr, ctx.clone())
253                                .await?
254                        }
255                    }
256                }
257            };
258
259            let output = interceptor.post_execute(output, ctx)?;
260            Ok(output)
261        }
262        .await;
263
264        result
265            .map_err(BoxedError::new)
266            .context(server_error::ExecuteGrpcQuerySnafu)
267    }
268
269    fn handle_put_record_batch_stream(
270        &self,
271        stream: servers::grpc::flight::PutRecordBatchRequestStream,
272        ctx: QueryContextRef,
273    ) -> Pin<Box<dyn Stream<Item = server_error::Result<DoPutResponse>> + Send>> {
274        Box::pin(
275            self.handle_put_record_batch_stream_inner(stream, ctx)
276                .map(|result| {
277                    result
278                        .map_err(BoxedError::new)
279                        .context(server_error::ExecuteGrpcRequestSnafu)
280                }),
281        )
282    }
283}
284
285fn fill_catalog_and_schema_from_context(ddl_expr: &mut DdlExpr, ctx: &QueryContextRef) {
286    let catalog = ctx.current_catalog();
287    let schema = ctx.current_schema();
288
289    macro_rules! check_and_fill {
290        ($expr:ident) => {
291            if $expr.catalog_name.is_empty() {
292                $expr.catalog_name = catalog.to_string();
293            }
294            if $expr.schema_name.is_empty() {
295                $expr.schema_name = schema.to_string();
296            }
297        };
298    }
299
300    match ddl_expr {
301        Expr::CreateDatabase(_) | Expr::AlterDatabase(_) => { /* do nothing*/ }
302        Expr::CreateTable(expr) => {
303            check_and_fill!(expr);
304        }
305        Expr::AlterTable(expr) => {
306            check_and_fill!(expr);
307        }
308        Expr::DropTable(expr) => {
309            check_and_fill!(expr);
310        }
311        Expr::TruncateTable(expr) => {
312            check_and_fill!(expr);
313        }
314        Expr::CreateFlow(expr) => {
315            if expr.catalog_name.is_empty() {
316                expr.catalog_name = catalog.to_string();
317            }
318        }
319        Expr::DropFlow(expr) => {
320            if expr.catalog_name.is_empty() {
321                expr.catalog_name = catalog.to_string();
322            }
323        }
324        Expr::CreateView(expr) => {
325            check_and_fill!(expr);
326        }
327        Expr::DropView(expr) => {
328            check_and_fill!(expr);
329        }
330        Expr::CommentOn(expr) => {
331            check_and_fill!(expr);
332        }
333    }
334}
335
336impl Instance {
337    pub(crate) fn check_table_permission(
338        &self,
339        ctx: &QueryContextRef,
340        req: PermissionReq<'_>,
341        targets: PermissionTableTargets,
342    ) -> auth::error::Result<PermissionResp> {
343        self.plugins
344            .get::<PermissionCheckerRef>()
345            .as_ref()
346            .check_permission_with_table_targets(ctx.current_user(), req, targets)
347    }
348
349    /// Checks every logical table targeted by normalized row inserts.
350    pub(crate) fn check_row_insert_permission(
351        &self,
352        requests: &RowInsertRequests,
353        ctx: &QueryContextRef,
354        req: PermissionReq<'_>,
355    ) -> auth::error::Result<PermissionResp> {
356        let catalog = ctx.current_catalog();
357        let schema = ctx.current_schema();
358        let targets = PermissionTableTargets::from_row_insert_requests(catalog, &schema, requests);
359
360        self.check_table_permission(ctx, req, targets)
361    }
362
363    fn handle_put_record_batch_stream_inner(
364        &self,
365        mut stream: servers::grpc::flight::PutRecordBatchRequestStream,
366        ctx: QueryContextRef,
367    ) -> Pin<Box<dyn Stream<Item = Result<DoPutResponse>> + Send>> {
368        // Clone all necessary data to make it 'static
369        let catalog_manager = self.catalog_manager().clone();
370        let plugins = self.plugins.clone();
371        let inserter = self.inserter.clone();
372        let ctx = ctx.clone();
373        let mut table_ref: Option<TableRef> = None;
374        let mut table_checked = false;
375
376        Box::pin(try_stream! {
377            // Process each request in the stream
378            while let Some(request_result) = stream.next().await {
379                let request = request_result.map_err(|e| {
380                    let error_msg = format!("Stream error: {:?}", e);
381                    IncompleteGrpcRequestSnafu { err_msg: error_msg }.build()
382                })?;
383
384                // Resolve table and check permissions on first RecordBatch (after schema is received)
385                if !table_checked {
386                    let table_name = &request.table_name;
387
388                    plugins
389                        .get::<PermissionCheckerRef>()
390                        .as_ref()
391                        .check_permission(
392                            ctx.current_user(),
393                            PermissionReq::BulkInsert {
394                                catalog: &table_name.catalog_name,
395                                schema: &table_name.schema_name,
396                                table: &table_name.table_name,
397                            },
398                        )
399                        .context(PermissionSnafu)?;
400
401                    // Resolve table reference
402                    table_ref = Some(
403                        catalog_manager
404                            .table(
405                                &table_name.catalog_name,
406                                &table_name.schema_name,
407                                &table_name.table_name,
408                                None,
409                            )
410                            .await
411                            .context(CatalogSnafu)?
412                            .with_context(|| TableNotFoundSnafu {
413                                table_name: table_name.to_string(),
414                            })?,
415                    );
416
417                    // Check permissions for the table
418                    let interceptor_ref = plugins.get::<GrpcQueryInterceptorRef<Error>>();
419                    let interceptor = interceptor_ref.as_ref();
420                    interceptor.pre_bulk_insert(table_ref.clone().unwrap(), ctx.clone())?;
421
422                    table_checked = true;
423                }
424
425                let request_id = request.request_id;
426                let start = Instant::now();
427                let rows = inserter
428                    .handle_bulk_insert(
429                        table_ref.clone().unwrap(),
430                        request.flight_data,
431                        request.record_batch,
432                        request.schema_bytes,
433                    )
434                    .await
435                    .context(TableOperationSnafu)?;
436                let elapsed_secs = start.elapsed().as_secs_f64();
437                yield DoPutResponse::new(request_id, rows, elapsed_secs);
438            }
439        })
440    }
441
442    async fn handle_insert_plan(
443        &self,
444        insert: InsertIntoPlan,
445        ctx: QueryContextRef,
446    ) -> Result<Output> {
447        let timer = GRPC_HANDLE_PLAN_ELAPSED.start_timer();
448        let table_name = insert.table_name.context(IncompleteGrpcRequestSnafu {
449            err_msg: "'table_name' is absent in InsertIntoPlan",
450        })?;
451
452        // use dummy catalog to provide table
453        let plan_decoder = self
454            .query_engine()
455            .engine_context(ctx.clone())
456            .new_plan_decoder()
457            .context(PlanStatementSnafu)?;
458
459        let dummy_catalog_list = Arc::new(
460            catalog::table_source::dummy_catalog::DummyCatalogList::new_with_query_ctx(
461                self.catalog_manager().clone(),
462                ctx.clone(),
463            ),
464        );
465
466        // no optimize yet since we still need to add stuff
467        let logical_plan = plan_decoder
468            .decode(
469                bytes::Bytes::from(insert.logical_plan),
470                dummy_catalog_list,
471                false,
472            )
473            .await
474            .context(SubstraitDecodeLogicalPlanSnafu)?;
475
476        let table = self
477            .catalog_manager()
478            .table(
479                &table_name.catalog_name,
480                &table_name.schema_name,
481                &table_name.table_name,
482                None,
483            )
484            .await
485            .context(CatalogSnafu)?
486            .with_context(|| TableNotFoundSnafu {
487                table_name: [
488                    table_name.catalog_name.clone(),
489                    table_name.schema_name.clone(),
490                    table_name.table_name.clone(),
491                ]
492                .join("."),
493            })?;
494        let table_provider = Arc::new(DfTableProviderAdapter::new(table));
495        let table_source = Arc::new(DefaultTableSource::new(table_provider));
496
497        let insert_into = add_insert_to_logical_plan(table_name, table_source, logical_plan)
498            .context(SubstraitDecodeLogicalPlanSnafu)?;
499
500        let engine_ctx = self.query_engine().engine_context(ctx.clone());
501        let state = engine_ctx.state();
502        // Analyze the plan
503        let analyzed_plan = state
504            .analyzer()
505            .execute_and_check(insert_into, state.config_options(), |_, _| {})
506            .context(DataFusionSnafu)?;
507
508        // Optimize the plan
509        let optimized_plan = state.optimize(&analyzed_plan).context(DataFusionSnafu)?;
510
511        let output = self
512            .do_exec_plan_inner(optimized_plan, None, ctx.clone())
513            .await?;
514
515        Ok(attach_timer(output, timer))
516    }
517    #[tracing::instrument(skip_all)]
518    pub async fn handle_inserts(
519        &self,
520        requests: InsertRequests,
521        ctx: QueryContextRef,
522    ) -> Result<Output> {
523        self.inserter
524            .handle_column_inserts(requests, ctx, self.statement_executor.as_ref())
525            .await
526            .context(TableOperationSnafu)
527    }
528
529    #[tracing::instrument(skip_all)]
530    pub async fn handle_row_inserts(
531        &self,
532        requests: RowInsertRequests,
533        ctx: QueryContextRef,
534        accommodate_existing_schema: bool,
535        is_single_value: bool,
536    ) -> Result<Output> {
537        self.inserter
538            .handle_row_inserts(
539                requests,
540                ctx,
541                self.statement_executor.as_ref(),
542                accommodate_existing_schema,
543                is_single_value,
544            )
545            .await
546            .context(TableOperationSnafu)
547    }
548
549    #[tracing::instrument(skip_all)]
550    pub async fn handle_influx_row_inserts(
551        &self,
552        requests: RowInsertRequests,
553        ctx: QueryContextRef,
554    ) -> Result<Output> {
555        self.inserter
556            .handle_last_non_null_inserts(
557                requests,
558                ctx,
559                self.statement_executor.as_ref(),
560                true,
561                // Influx protocol may writes multiple fields (values).
562                false,
563            )
564            .await
565            .context(TableOperationSnafu)
566    }
567
568    #[tracing::instrument(skip_all)]
569    pub async fn handle_metric_row_inserts(
570        &self,
571        requests: RowInsertRequests,
572        ctx: QueryContextRef,
573        physical_table: String,
574    ) -> Result<Output> {
575        self.inserter
576            .handle_metric_row_inserts(requests, ctx, &self.statement_executor, physical_table)
577            .await
578            .context(TableOperationSnafu)
579    }
580
581    #[tracing::instrument(skip_all)]
582    pub async fn handle_deletes(
583        &self,
584        requests: DeleteRequests,
585        ctx: QueryContextRef,
586    ) -> Result<Output> {
587        self.deleter
588            .handle_column_deletes(requests, ctx)
589            .await
590            .context(TableOperationSnafu)
591    }
592
593    #[tracing::instrument(skip_all)]
594    pub async fn handle_row_deletes(
595        &self,
596        requests: RowDeleteRequests,
597        ctx: QueryContextRef,
598    ) -> Result<Output> {
599        self.deleter
600            .handle_row_deletes(requests, ctx)
601            .await
602            .context(TableOperationSnafu)
603    }
604}