Skip to main content

servers/grpc/
prom_query_gateway.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
15//! PrometheusGateway provides a gRPC interface to query Prometheus metrics
16//! by PromQL. The behavior is similar to the Prometheus HTTP API.
17
18use std::sync::Arc;
19
20use api::v1::prometheus_gateway_server::PrometheusGateway;
21use api::v1::promql_request::Promql;
22use api::v1::{PromqlRequest, PromqlResponse, ResponseHeader};
23use async_trait::async_trait;
24use auth::UserProviderRef;
25use common_error::ext::ErrorExt;
26use common_error::status_code::StatusCode;
27use common_time::util::current_time_rfc3339;
28use promql_parser::parser::value::ValueType;
29use query::parser::PromQuery;
30use session::context::{Channel, QueryContext};
31use snafu::OptionExt;
32use tonic::{Request, Response};
33
34use crate::error::InvalidQuerySnafu;
35use crate::grpc::TonicResult;
36use crate::grpc::context_auth::auth;
37use crate::grpc::greptime_handler::create_query_context;
38use crate::http::prometheus::{PrometheusJsonResponse, retrieve_metric_name_and_result_type};
39use crate::prometheus_handler::{ParsedPromQuery, PrometheusHandlerRef};
40
41pub struct PrometheusGatewayService {
42    handler: PrometheusHandlerRef,
43    user_provider: Option<UserProviderRef>,
44}
45
46#[async_trait]
47impl PrometheusGateway for PrometheusGatewayService {
48    async fn handle(&self, req: Request<PromqlRequest>) -> TonicResult<Response<PromqlResponse>> {
49        let mut is_range_query = false;
50        let channel = req
51            .extensions()
52            .get::<Channel>()
53            .copied()
54            .unwrap_or(Channel::Promql);
55        let inner = req.into_inner();
56        let prom_query = match inner.promql.context(InvalidQuerySnafu {
57            reason: "Expecting non-empty PromqlRequest.",
58        })? {
59            Promql::RangeQuery(range_query) => {
60                is_range_query = true;
61                PromQuery {
62                    query: range_query.query,
63                    start: range_query.start,
64                    end: range_query.end,
65                    step: range_query.step,
66                    lookback: range_query.lookback,
67                    alias: None,
68                }
69            }
70            Promql::InstantQuery(instant_query) => {
71                let time = if instant_query.time.is_empty() {
72                    current_time_rfc3339()
73                } else {
74                    instant_query.time
75                };
76                PromQuery {
77                    query: instant_query.query,
78                    start: time.clone(),
79                    end: time,
80                    step: String::from("1s"),
81                    lookback: instant_query.lookback,
82                    alias: None,
83                }
84            }
85        };
86
87        let header = inner.header.as_ref();
88        let query_ctx =
89            create_query_context(channel, header, Default::default(), Default::default())?;
90
91        let user_info = auth(self.user_provider.clone(), header, &query_ctx).await?;
92        query_ctx.set_current_user(user_info);
93
94        let json_response = self
95            .handle_inner(prom_query, query_ctx, is_range_query)
96            .await;
97        let json_bytes = serde_json::to_string(&json_response).unwrap().into_bytes();
98
99        let response = Response::new(PromqlResponse {
100            header: Some(ResponseHeader {
101                status: Some(api::v1::Status {
102                    status_code: StatusCode::Success as _,
103                    ..Default::default()
104                }),
105            }),
106            body: json_bytes,
107        });
108        Ok(response)
109    }
110}
111
112impl PrometheusGatewayService {
113    pub fn new(handler: PrometheusHandlerRef, user_provider: Option<UserProviderRef>) -> Self {
114        Self {
115            handler,
116            user_provider,
117        }
118    }
119
120    async fn handle_inner(
121        &self,
122        query: PromQuery,
123        ctx: Arc<QueryContext>,
124        is_range_query: bool,
125    ) -> PrometheusJsonResponse {
126        let db = ctx.get_db_string();
127        let _timer = crate::metrics::METRIC_SERVER_GRPC_PROM_REQUEST_TIMER
128            .with_label_values(&[db.as_str()])
129            .start_timer();
130
131        let query = match ParsedPromQuery::parse(query, &ctx) {
132            Ok(query) => query,
133            Err(err) => {
134                return PrometheusJsonResponse::error(err.status_code(), err.output_msg());
135            }
136        };
137        let (metric_name, mut result_type) = retrieve_metric_name_and_result_type(query.expr());
138        let query_id = ctx.remote_query_id().map(str::to_string);
139        // A range query only returns a matrix, and matrix serialization sorts
140        // samples and series, so execution order never reaches the response.
141        let query = if is_range_query {
142            result_type = ValueType::Matrix;
143            query.with_unordered_output()
144        } else {
145            query
146        };
147        let result = self.handler.do_query_parsed(query, ctx).await;
148
149        PrometheusJsonResponse::from_query_result(
150            result,
151            metric_name,
152            result_type,
153            query_id.as_deref(),
154        )
155        .await
156    }
157}