servers/grpc/
prom_query_gateway.rs1use 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 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}