1use 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 let timer = GRPC_HANDLE_PLAN_ELAPSED.start_timer();
121
122 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(_) => { }
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 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 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 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 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 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 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 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 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 let analyzed_plan = state
504 .analyzer()
505 .execute_and_check(insert_into, state.config_options(), |_, _| {})
506 .context(DataFusionSnafu)?;
507
508 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 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}