Skip to main content

servers/http/
prometheus.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//! prom supply the prometheus HTTP API Server compliance
16
17use std::borrow::Borrow;
18use std::collections::{BTreeMap, HashMap, HashSet};
19use std::hash::{Hash, Hasher};
20use std::sync::Arc;
21
22use arrow::array::{Array, AsArray};
23use arrow::datatypes::{
24    Date32Type, Date64Type, Decimal128Type, Float32Type, Float64Type, Int8Type, Int16Type,
25    Int32Type, Int64Type, IntervalDayTimeType, IntervalMonthDayNanoType, IntervalYearMonthType,
26    UInt8Type, UInt16Type, UInt32Type, UInt64Type,
27};
28use arrow_schema::{DataType, IntervalUnit};
29use auth::{PermissionTableTarget, PermissionTableTargets};
30use axum::extract::{Path, Query, State};
31use axum::{Extension, Form};
32use catalog::CatalogManagerRef;
33use common_catalog::parse_catalog_and_schema_from_db_string;
34use common_decimal::Decimal128;
35use common_error::ext::ErrorExt;
36use common_error::status_code::StatusCode;
37use common_query::native_histogram::is_native_histogram_value_type;
38use common_query::{Output, OutputData};
39use common_recordbatch::{RecordBatch, RecordBatches};
40use common_telemetry::{debug, tracing};
41use common_time::util::{current_time_rfc3339, yesterday_rfc3339};
42use common_time::{Date, IntervalDayTime, IntervalMonthDayNano, IntervalYearMonth};
43use common_version::OwnedBuildInfo;
44use datafusion_common::ScalarValue;
45use datatypes::prelude::ConcreteDataType;
46use datatypes::schema::{ColumnSchema, SchemaRef};
47use datatypes::types::jsonb_to_string;
48use futures::future::join_all;
49use futures::{StreamExt, TryStreamExt};
50use itertools::Itertools;
51use promql_parser::label::{METRIC_NAME, MatchOp, Matcher, Matchers};
52use promql_parser::parser::token::{self};
53use promql_parser::parser::value::ValueType;
54use promql_parser::parser::{
55    AggregateExpr, BinaryExpr, Call, Expr as PromqlExpr, LabelModifier, MatrixSelector, ParenExpr,
56    SubqueryExpr, UnaryExpr, VectorSelector,
57};
58use query::parser::{DEFAULT_LOOKBACK_STRING, PromQuery, QueryStatement};
59use serde::de::{self, MapAccess, Visitor};
60use serde::{Deserialize, Serialize};
61use serde_json::Value;
62use session::context::{QueryContext, QueryContextRef};
63use snafu::{Location, OptionExt, ResultExt};
64use store_api::metric_engine_consts::{
65    DATA_SCHEMA_TABLE_ID_COLUMN_NAME, DATA_SCHEMA_TSID_COLUMN_NAME, LOGICAL_TABLE_METADATA_KEY,
66};
67use table::TableRef;
68use table::metadata::TableInfo;
69use table::requests::{
70    METRIC_TEMPORALITY_DELTA, SEMANTIC_METRIC_TEMPORALITY, SEMANTIC_METRIC_TYPE,
71    SEMANTIC_METRIC_UNIT, SEMANTIC_VALUE_MIXED,
72};
73
74pub use super::result::prometheus_resp::PrometheusJsonResponse;
75use crate::error::{
76    CollectRecordbatchSnafu, ConvertScalarValueSnafu, DataFusionSnafu, Error, InvalidQuerySnafu,
77    NotSupportedSnafu, Result, TableNotFoundSnafu, UnexpectedResultSnafu,
78};
79use crate::http::header::collect_plan_metrics;
80use crate::otlp::metrics::ucum_to_openmetrics_unit;
81use crate::prom_store::{FIELD_NAME_LABEL, METRIC_NAME_LABEL, is_database_selection_label};
82use crate::prometheus_handler::{
83    ParsedPromQuery, PrometheusHandlerRef, resolve_schema_from_matchers,
84};
85
86/// For [ValueType::Vector] result type
87#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
88#[serde(deny_unknown_fields)]
89pub struct PromSeriesVector {
90    pub metric: BTreeMap<String, String>,
91    #[serde(skip_serializing_if = "Option::is_none")]
92    pub value: Option<(f64, String)>,
93    #[serde(skip_serializing_if = "Option::is_none")]
94    pub histogram: Option<(f64, PromNativeHistogram)>,
95}
96
97/// For [ValueType::Matrix] result type
98#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
99#[serde(deny_unknown_fields)]
100pub struct PromSeriesMatrix {
101    pub metric: BTreeMap<String, String>,
102    #[serde(skip_serializing_if = "Vec::is_empty", default)]
103    pub values: Vec<(f64, String)>,
104    #[serde(skip_serializing_if = "Vec::is_empty", default)]
105    pub histograms: Vec<(f64, PromNativeHistogram)>,
106}
107
108/// A native histogram sample in the Prometheus HTTP API JSON format.
109#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
110pub struct PromNativeHistogram {
111    /// Total number of observations, encoded as a string.
112    pub count: String,
113    /// Sum of all observations, encoded as a string.
114    pub sum: String,
115    /// Populated buckets as `(boundary rule, lower bound, upper bound, count)`.
116    ///
117    /// Boundary rules are `0` for lower-open/upper-closed, `1` for
118    /// lower-closed/upper-open, `2` for both bounds open, and `3` for both
119    /// bounds closed. Bounds and counts are encoded as strings.
120    pub buckets: Vec<(u8, String, String, String)>,
121}
122
123/// Variants corresponding to [ValueType]
124#[derive(Debug, Serialize, Deserialize, PartialEq)]
125#[serde(untagged)]
126pub enum PromQueryResult {
127    Matrix(Vec<PromSeriesMatrix>),
128    Vector(Vec<PromSeriesVector>),
129    Scalar(#[serde(skip_serializing_if = "Option::is_none")] Option<(f64, String)>),
130    String(#[serde(skip_serializing_if = "Option::is_none")] Option<(f64, String)>),
131}
132
133impl Default for PromQueryResult {
134    fn default() -> Self {
135        PromQueryResult::Matrix(Default::default())
136    }
137}
138
139#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
140pub struct PromData {
141    #[serde(rename = "resultType")]
142    pub result_type: String,
143    pub result: PromQueryResult,
144}
145
146/// Metadata for a Prometheus metric family.
147#[derive(Debug, Default, Serialize, Deserialize, PartialEq, Eq)]
148pub struct PromMetadata {
149    #[serde(rename = "type")]
150    pub metric_type: String,
151    pub unit: String,
152    pub help: String,
153}
154
155/// A "holder" for the reference([Arc]) to a column name,
156/// to help avoiding cloning [String]s when used as a [HashMap] key.
157#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
158pub struct Column(Arc<String>);
159
160impl From<&str> for Column {
161    fn from(s: &str) -> Self {
162        Self(Arc::new(s.to_string()))
163    }
164}
165
166#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
167#[serde(untagged)]
168pub enum PrometheusResponse {
169    PromData(PromData),
170    Labels(Vec<String>),
171    Series(Vec<HashMap<Column, String>>),
172    Metadata(BTreeMap<String, Vec<PromMetadata>>),
173    LabelValues(Vec<String>),
174    FormatQuery(String),
175    BuildInfo(OwnedBuildInfo),
176    #[serde(skip_deserializing)]
177    ParseResult(promql_parser::parser::Expr),
178    #[default]
179    None,
180}
181
182impl PrometheusResponse {
183    /// Append the other [`PrometheusResponse]`.
184    /// # NOTE
185    ///   Only append matrix and vector results, otherwise just ignore the other response.
186    pub(super) fn append(&mut self, other: PrometheusResponse) {
187        match (self, other) {
188            (
189                PrometheusResponse::PromData(PromData {
190                    result: PromQueryResult::Matrix(lhs),
191                    ..
192                }),
193                PrometheusResponse::PromData(PromData {
194                    result: PromQueryResult::Matrix(rhs),
195                    ..
196                }),
197            ) => {
198                lhs.extend(rhs);
199            }
200
201            (
202                PrometheusResponse::PromData(PromData {
203                    result: PromQueryResult::Vector(lhs),
204                    ..
205                }),
206                PrometheusResponse::PromData(PromData {
207                    result: PromQueryResult::Vector(rhs),
208                    ..
209                }),
210            ) => {
211                lhs.extend(rhs);
212            }
213            _ => {
214                // TODO(dennis): process other cases?
215            }
216        }
217    }
218
219    pub fn is_none(&self) -> bool {
220        matches!(self, PrometheusResponse::None)
221    }
222}
223
224#[derive(Debug, Default, Serialize, Deserialize)]
225pub struct FormatQuery {
226    query: Option<String>,
227}
228
229/// Query parameters for the Prometheus metric metadata endpoint.
230#[derive(Debug, Default, Serialize, Deserialize)]
231pub struct MetadataQuery {
232    db: Option<String>,
233    limit: Option<usize>,
234    metric: Option<String>,
235}
236
237#[axum_macros::debug_handler]
238#[tracing::instrument(
239    skip_all,
240    fields(protocol = "prometheus", request_type = "format_query")
241)]
242pub async fn format_query(
243    State(_handler): State<PrometheusHandlerRef>,
244    Query(params): Query<InstantQuery>,
245    Extension(_query_ctx): Extension<QueryContext>,
246    Form(form_params): Form<InstantQuery>,
247) -> PrometheusJsonResponse {
248    let query = params.query.or(form_params.query).unwrap_or_default();
249    match promql_parser::parser::parse(&query) {
250        Ok(expr) => {
251            let pretty = expr.prettify();
252            PrometheusJsonResponse::success(PrometheusResponse::FormatQuery(pretty))
253        }
254        Err(reason) => {
255            let err = InvalidQuerySnafu { reason }.build();
256            PrometheusJsonResponse::error(err.status_code(), err.output_msg())
257        }
258    }
259}
260
261#[derive(Debug, Default, Serialize, Deserialize)]
262pub struct BuildInfoQuery {}
263
264#[axum_macros::debug_handler]
265#[tracing::instrument(
266    skip_all,
267    fields(protocol = "prometheus", request_type = "build_info_query")
268)]
269pub async fn build_info_query() -> PrometheusJsonResponse {
270    let build_info = common_version::build_info().clone();
271    PrometheusJsonResponse::success(PrometheusResponse::BuildInfo(build_info.into()))
272}
273
274#[axum_macros::debug_handler]
275#[tracing::instrument(
276    skip_all,
277    fields(protocol = "prometheus", request_type = "metadata_query")
278)]
279pub async fn metadata_query(
280    State(handler): State<PrometheusHandlerRef>,
281    Query(params): Query<MetadataQuery>,
282    Extension(mut query_ctx): Extension<QueryContext>,
283) -> PrometheusJsonResponse {
284    let (catalog, schema) = get_catalog_schema(&params.db, &query_ctx);
285    try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
286    let query_ctx = Arc::new(query_ctx);
287
288    if let Err(err) = handler.check_query_permission(&[], &query_ctx).await {
289        return PrometheusJsonResponse::error(err.status_code(), err.to_string());
290    }
291    if let Some(metric) = &params.metric
292        && let Err(err) = handler
293            .check_query_target_permission(
294                current_schema_metric_targets(&query_ctx, std::slice::from_ref(metric)),
295                &query_ctx,
296            )
297            .await
298    {
299        return PrometheusJsonResponse::error(err.status_code(), err.to_string());
300    }
301
302    let mut metadata =
303        match retrieve_metric_metadata(&query_ctx, handler.catalog_manager(), &params).await {
304            Ok(metadata) => metadata,
305            Err(err) => {
306                return PrometheusJsonResponse::error(
307                    StatusCode::InvalidArguments,
308                    err.to_string(),
309                );
310            }
311        };
312    let allowed_metric_names = match handler
313        .filter_metadata_metric_names(
314            metadata.keys().cloned().collect(),
315            query_ctx.current_schema().as_str(),
316            &query_ctx,
317        )
318        .await
319    {
320        Ok(metric_names) => metric_names,
321        Err(err) => {
322            return PrometheusJsonResponse::error(err.status_code(), err.to_string());
323        }
324    }
325    .into_iter()
326    .collect::<HashSet<_>>();
327    metadata.retain(|metric, _| allowed_metric_names.contains(metric));
328    if let Some(limit) = params.limit {
329        metadata = metadata.into_iter().take(limit).collect();
330    }
331
332    PrometheusJsonResponse::success(PrometheusResponse::Metadata(metadata))
333}
334
335#[derive(Debug, Default, Serialize, Deserialize)]
336pub struct InstantQuery {
337    query: Option<String>,
338    lookback: Option<String>,
339    time: Option<String>,
340    timeout: Option<String>,
341    db: Option<String>,
342}
343
344/// Helper macro which try to evaluate the expression and return its results.
345/// If the evaluation fails, return a `PrometheusJsonResponse` early.
346macro_rules! try_call_return_response {
347    (@output_msg $handle: expr, $status_code: expr) => {
348        match $handle {
349            Ok(res) => res,
350            Err(err) => {
351                let msg = err.output_msg();
352                return PrometheusJsonResponse::error($status_code, msg);
353            }
354        }
355    };
356    ($handle: expr, $status_code: expr) => {
357        match $handle {
358            Ok(res) => res,
359            Err(err) => {
360                let msg = err.to_string();
361                return PrometheusJsonResponse::error($status_code, msg);
362            }
363        }
364    };
365    ($handle: expr) => {
366        match $handle {
367            Ok(res) => res,
368            Err(err) => {
369                let status_code = err.status_code();
370                let msg = err.to_string();
371                return PrometheusJsonResponse::error(status_code, msg);
372            }
373        }
374    };
375}
376
377#[axum_macros::debug_handler]
378#[tracing::instrument(
379    skip_all,
380    fields(protocol = "prometheus", request_type = "instant_query")
381)]
382pub async fn instant_query(
383    State(handler): State<PrometheusHandlerRef>,
384    Query(params): Query<InstantQuery>,
385    Extension(mut query_ctx): Extension<QueryContext>,
386    Form(form_params): Form<InstantQuery>,
387) -> PrometheusJsonResponse {
388    // Extract time from query string, or use current server time if not specified.
389    let time = params
390        .time
391        .or(form_params.time)
392        .unwrap_or_else(current_time_rfc3339);
393    let prom_query = PromQuery {
394        query: params.query.or(form_params.query).unwrap_or_default(),
395        start: time.clone(),
396        end: time,
397        step: "1s".to_string(),
398        lookback: params
399            .lookback
400            .or(form_params.lookback)
401            .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string()),
402        alias: None,
403    };
404
405    // update catalog and schema in query context if necessary
406    if let Some(db) = &params.db {
407        let (catalog, schema) = parse_catalog_and_schema_from_db_string(db);
408        try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
409    }
410    let query_ctx = Arc::new(query_ctx);
411    let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
412        .with_label_values(&[query_ctx.get_db_string().as_str(), "instant_query"])
413        .start_timer();
414    let prom_query = try_call_return_response!(
415        @output_msg ParsedPromQuery::parse(prom_query, &query_ctx),
416        StatusCode::InvalidArguments
417    );
418    let promql_expr = prom_query.expr();
419
420    let metric_name_discovery =
421        try_call_return_response!(find_metric_name_not_equal_matchers(promql_expr));
422    if let Some(discovery) = metric_name_discovery {
423        debug!("Find metric name matchers: {:?}", discovery.name_matchers);
424
425        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
426        let (static_targets, unresolved_selectors) =
427            try_call_return_response!(static_promql_targets(promql_expr, &query_ctx));
428        try_call_return_response!(
429            handler
430                .check_query_target_permission(static_targets, &query_ctx,)
431                .await
432        );
433        if unresolved_selectors != 1 {
434            try_call_return_response!(
435                handler
436                    .check_query_target_permission(PermissionTableTargets::Unresolved, &query_ctx,)
437                    .await
438            );
439            return do_instant_query(&handler, prom_query, query_ctx).await;
440        }
441
442        let schema = discovery
443            .schema
444            .unwrap_or_else(|| query_ctx.current_schema());
445        let metric_names = try_call_return_response!(
446            handler
447                .query_metric_names(discovery.name_matchers, &schema, &query_ctx)
448                .await
449        );
450
451        debug!("Find metric names: {:?}", metric_names);
452
453        let prom_queries = expand_metric_name_queries(&prom_query, metric_names);
454        try_call_return_response!(
455            handler
456                .check_query_permission_parsed(&prom_queries, &query_ctx)
457                .await
458        );
459
460        if prom_queries.is_empty() {
461            let result_type = promql_expr.value_type();
462
463            return PrometheusJsonResponse::success(PrometheusResponse::PromData(PromData {
464                result_type: result_type.to_string(),
465                ..Default::default()
466            }));
467        }
468
469        let responses = join_all(prom_queries.into_iter().map(|prom_query| {
470            let query_ctx = query_ctx.clone();
471            let handler = handler.clone();
472
473            async move { do_instant_query(&handler, prom_query, query_ctx).await }
474        }))
475        .await;
476
477        responses
478            .into_iter()
479            .reduce(|mut acc, resp| {
480                acc.append_query_response(resp);
481                acc
482            })
483            .unwrap()
484    } else {
485        do_instant_query(&handler, prom_query, query_ctx).await
486    }
487}
488
489/// Executes a single instant query and returns response
490async fn do_instant_query(
491    handler: &PrometheusHandlerRef,
492    prom_query: ParsedPromQuery,
493    query_ctx: QueryContextRef,
494) -> PrometheusJsonResponse {
495    let (metric_name, result_type) = retrieve_metric_name_and_result_type(prom_query.expr());
496    let query_id = query_ctx.remote_query_id().map(str::to_string);
497    let result = handler.do_query_parsed(prom_query, query_ctx).await;
498    PrometheusJsonResponse::from_query_result(result, metric_name, result_type, query_id.as_deref())
499        .await
500}
501
502#[derive(Debug, Default, Serialize, Deserialize)]
503pub struct RangeQuery {
504    query: Option<String>,
505    start: Option<String>,
506    end: Option<String>,
507    step: Option<String>,
508    lookback: Option<String>,
509    timeout: Option<String>,
510    db: Option<String>,
511}
512
513#[axum_macros::debug_handler]
514#[tracing::instrument(
515    skip_all,
516    fields(protocol = "prometheus", request_type = "range_query")
517)]
518pub async fn range_query(
519    State(handler): State<PrometheusHandlerRef>,
520    Query(params): Query<RangeQuery>,
521    Extension(mut query_ctx): Extension<QueryContext>,
522    Form(form_params): Form<RangeQuery>,
523) -> PrometheusJsonResponse {
524    let prom_query = PromQuery {
525        query: params.query.or(form_params.query).unwrap_or_default(),
526        start: params.start.or(form_params.start).unwrap_or_default(),
527        end: params.end.or(form_params.end).unwrap_or_default(),
528        step: params.step.or(form_params.step).unwrap_or_default(),
529        lookback: params
530            .lookback
531            .or(form_params.lookback)
532            .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string()),
533        alias: None,
534    };
535
536    // update catalog and schema in query context if necessary
537    if let Some(db) = &params.db {
538        let (catalog, schema) = parse_catalog_and_schema_from_db_string(db);
539        try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
540    }
541    let query_ctx = Arc::new(query_ctx);
542    let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
543        .with_label_values(&[query_ctx.get_db_string().as_str(), "range_query"])
544        .start_timer();
545    let prom_query = try_call_return_response!(
546        @output_msg ParsedPromQuery::parse(prom_query, &query_ctx),
547        StatusCode::InvalidArguments
548    );
549    let promql_expr = prom_query.expr();
550
551    let metric_name_discovery =
552        try_call_return_response!(find_metric_name_not_equal_matchers(promql_expr));
553    if let Some(discovery) = metric_name_discovery {
554        debug!("Find metric name matchers: {:?}", discovery.name_matchers);
555
556        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
557        let (static_targets, unresolved_selectors) =
558            try_call_return_response!(static_promql_targets(promql_expr, &query_ctx));
559        try_call_return_response!(
560            handler
561                .check_query_target_permission(static_targets, &query_ctx,)
562                .await
563        );
564        if unresolved_selectors != 1 {
565            try_call_return_response!(
566                handler
567                    .check_query_target_permission(PermissionTableTargets::Unresolved, &query_ctx,)
568                    .await
569            );
570            return do_range_query(&handler, prom_query, query_ctx).await;
571        }
572
573        let schema = discovery
574            .schema
575            .unwrap_or_else(|| query_ctx.current_schema());
576        let metric_names = try_call_return_response!(
577            handler
578                .query_metric_names(discovery.name_matchers, &schema, &query_ctx)
579                .await
580        );
581
582        debug!("Find metric names: {:?}", metric_names);
583
584        let prom_queries = expand_metric_name_queries(&prom_query, metric_names);
585        try_call_return_response!(
586            handler
587                .check_query_permission_parsed(&prom_queries, &query_ctx)
588                .await
589        );
590
591        if prom_queries.is_empty() {
592            return PrometheusJsonResponse::success(PrometheusResponse::PromData(PromData {
593                result_type: ValueType::Matrix.to_string(),
594                ..Default::default()
595            }));
596        }
597
598        let responses = join_all(prom_queries.into_iter().map(|prom_query| {
599            let query_ctx = query_ctx.clone();
600            let handler = handler.clone();
601
602            async move { do_range_query(&handler, prom_query, query_ctx).await }
603        }))
604        .await;
605
606        // Safety: at least one responses, checked above
607        responses
608            .into_iter()
609            .reduce(|mut acc, resp| {
610                acc.append_query_response(resp);
611                acc
612            })
613            .unwrap()
614    } else {
615        do_range_query(&handler, prom_query, query_ctx).await
616    }
617}
618
619/// Executes a single range query and returns response
620async fn do_range_query(
621    handler: &PrometheusHandlerRef,
622    prom_query: ParsedPromQuery,
623    query_ctx: QueryContextRef,
624) -> PrometheusJsonResponse {
625    let (metric_name, _) = retrieve_metric_name_and_result_type(prom_query.expr());
626    let query_id = query_ctx.remote_query_id().map(str::to_string);
627    let result = handler.do_query_parsed(prom_query, query_ctx).await;
628    PrometheusJsonResponse::from_query_result(
629        result,
630        metric_name,
631        ValueType::Matrix,
632        query_id.as_deref(),
633    )
634    .await
635}
636
637#[derive(Debug, Default, Serialize)]
638struct Matches(Vec<String>);
639
640#[derive(Debug, Default, Serialize, Deserialize)]
641pub struct LabelsQuery {
642    start: Option<String>,
643    end: Option<String>,
644    lookback: Option<String>,
645    #[serde(flatten)]
646    matches: Matches,
647    db: Option<String>,
648}
649
650// Custom Deserialize method to support parsing repeated match[]
651impl<'de> Deserialize<'de> for Matches {
652    fn deserialize<D>(deserializer: D) -> std::result::Result<Matches, D::Error>
653    where
654        D: de::Deserializer<'de>,
655    {
656        struct MatchesVisitor;
657
658        impl<'d> Visitor<'d> for MatchesVisitor {
659            type Value = Vec<String>;
660
661            fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result {
662                formatter.write_str("a string")
663            }
664
665            fn visit_map<M>(self, mut access: M) -> std::result::Result<Self::Value, M::Error>
666            where
667                M: MapAccess<'d>,
668            {
669                let mut matches = Vec::new();
670                while let Some((key, value)) = access.next_entry::<String, String>()? {
671                    if key == "match[]" {
672                        matches.push(value);
673                    }
674                }
675                Ok(matches)
676            }
677        }
678        Ok(Matches(deserializer.deserialize_map(MatchesVisitor)?))
679    }
680}
681
682/// Handles schema errors, transforming a Result into an Option.
683/// - If the input is `Ok(v)`, returns `Some(v)`
684/// - If the input is `Err(err)` and the error status code is `TableNotFound` or
685///   `TableColumnNotFound`, returns `None` (ignoring these specific errors)
686/// - If the input is `Err(err)` with any other error code, directly returns a
687///   `PrometheusJsonResponse::error`.
688macro_rules! handle_schema_err {
689    ($result:expr) => {
690        match $result {
691            Ok(v) => Some(v),
692            Err(err) => {
693                if err.status_code() == StatusCode::TableNotFound
694                    || err.status_code() == StatusCode::TableColumnNotFound
695                {
696                    // Prometheus won't report error if querying nonexist label and metric
697                    None
698                } else {
699                    return PrometheusJsonResponse::error(err.status_code(), err.output_msg());
700                }
701            }
702        }
703    };
704}
705
706#[axum_macros::debug_handler]
707#[tracing::instrument(
708    skip_all,
709    fields(protocol = "prometheus", request_type = "labels_query")
710)]
711pub async fn labels_query(
712    State(handler): State<PrometheusHandlerRef>,
713    Query(params): Query<LabelsQuery>,
714    Extension(mut query_ctx): Extension<QueryContext>,
715    Form(form_params): Form<LabelsQuery>,
716) -> PrometheusJsonResponse {
717    let (catalog, schema) = get_catalog_schema(&params.db, &query_ctx);
718    try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
719    let query_ctx = Arc::new(query_ctx);
720
721    let mut queries = params.matches.0;
722    if queries.is_empty() {
723        queries = form_params.matches.0;
724    }
725
726    let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
727        .with_label_values(&[query_ctx.get_db_string().as_str(), "labels_query"])
728        .start_timer();
729
730    let prom_queries = if queries.is_empty() {
731        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
732        Vec::new()
733    } else {
734        let start = params
735            .start
736            .or(form_params.start)
737            .unwrap_or_else(yesterday_rfc3339);
738        let end = params
739            .end
740            .or(form_params.end)
741            .unwrap_or_else(current_time_rfc3339);
742        let lookback = params
743            .lookback
744            .or(form_params.lookback)
745            .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string());
746        let prom_queries = queries
747            .into_iter()
748            .map(|query| PromQuery {
749                query,
750                start: start.clone(),
751                end: end.clone(),
752                step: DEFAULT_LOOKBACK_STRING.to_string(),
753                lookback: lookback.clone(),
754                alias: None,
755            })
756            .collect::<Vec<_>>();
757        let prom_queries = try_call_return_response!(
758            prom_queries
759                .into_iter()
760                .map(|query| ParsedPromQuery::parse(query, &query_ctx))
761                .collect::<Result<Vec<_>>>()
762        );
763        try_call_return_response!(
764            handler
765                .check_query_permission_parsed(&prom_queries, &query_ctx)
766                .await
767        );
768        prom_queries
769    };
770
771    // Fetch all columns if no query matcher is provided
772    if prom_queries.is_empty() {
773        let (mut labels, table_names) =
774            match get_all_column_names(&catalog, &schema, &handler.catalog_manager()).await {
775                Ok(result) => result,
776                Err(e) => return PrometheusJsonResponse::error(e.status_code(), e.output_msg()),
777            };
778        try_call_return_response!(
779            handler
780                .check_query_target_permission(
781                    current_schema_metric_targets(&query_ctx, &table_names),
782                    &query_ctx,
783                )
784                .await
785        );
786        let _ = labels.insert(METRIC_NAME.to_string());
787        let mut labels_vec = labels.into_iter().collect::<Vec<_>>();
788        labels_vec.sort_unstable();
789        return PrometheusJsonResponse::success(PrometheusResponse::Labels(labels_vec));
790    }
791
792    // Fetch tag columns only from the tables checked above.
793    let mut labels = match resolved_promql_targets(&prom_queries, &query_ctx) {
794        Some(targets) => {
795            match get_target_column_names(&targets, &handler.catalog_manager(), &query_ctx).await {
796                Ok(labels) => labels,
797                Err(e) => return PrometheusJsonResponse::error(e.status_code(), e.output_msg()),
798            }
799        }
800        None => match get_all_column_names(&catalog, &schema, &handler.catalog_manager()).await {
801            Ok((labels, _)) => labels,
802            Err(e) => return PrometheusJsonResponse::error(e.status_code(), e.output_msg()),
803        },
804    };
805    let _ = labels.insert(METRIC_NAME.to_string());
806
807    let mut fetched_labels = HashSet::new();
808    let _ = fetched_labels.insert(METRIC_NAME.to_string());
809
810    let mut merge_map = HashMap::new();
811    for prom_query in prom_queries {
812        let result = handler.do_query_parsed(prom_query, query_ctx.clone()).await;
813        handle_schema_err!(
814            retrieve_labels_name_from_query_result(result, &mut fetched_labels, &mut merge_map)
815                .await
816        );
817    }
818
819    // intersect `fetched_labels` with `labels` to filter out non-tag columns
820    fetched_labels.retain(|l| labels.contains(l));
821    let mut sorted_labels: Vec<String> = fetched_labels.into_iter().collect();
822    sorted_labels.sort();
823    let merge_map = merge_map
824        .into_iter()
825        .map(|(k, v)| (k, Value::from(v)))
826        .collect();
827    let mut resp = PrometheusJsonResponse::success(PrometheusResponse::Labels(sorted_labels));
828    resp.resp_metrics = merge_map;
829    resp
830}
831
832/// Get all tag column name of the given schema
833async fn get_all_column_names(
834    catalog: &str,
835    schema: &str,
836    manager: &CatalogManagerRef,
837) -> std::result::Result<(HashSet<String>, Vec<String>), catalog::error::Error> {
838    let mut labels = HashSet::new();
839    let mut table_names = Vec::new();
840    let mut tables = manager.tables(catalog, schema, None);
841    while let Some(table) = tables.try_next().await? {
842        if is_internal_physical_metric_table(&table) {
843            continue;
844        }
845        table_names.push(table.table_info().name.clone());
846        extend_tag_column_names(&mut labels, &table);
847    }
848
849    Ok((labels, table_names))
850}
851
852async fn get_target_column_names(
853    targets: &[PermissionTableTarget],
854    manager: &CatalogManagerRef,
855    query_ctx: &QueryContext,
856) -> std::result::Result<HashSet<String>, catalog::error::Error> {
857    let mut labels = HashSet::new();
858    for target in targets {
859        if let Some(table) = manager
860            .table(
861                &target.catalog,
862                &target.schema,
863                &target.table,
864                Some(query_ctx),
865            )
866            .await?
867            && !is_internal_physical_metric_table(&table)
868        {
869            extend_tag_column_names(&mut labels, &table);
870        }
871    }
872    Ok(labels)
873}
874
875fn extend_tag_column_names(labels: &mut HashSet<String>, table: &TableRef) {
876    labels.extend(table.primary_key_columns().filter_map(|column| {
877        (column.name != DATA_SCHEMA_TABLE_ID_COLUMN_NAME
878            && column.name != DATA_SCHEMA_TSID_COLUMN_NAME)
879            .then_some(column.name)
880    }));
881}
882
883async fn retrieve_series_from_query_result(
884    result: Result<Output>,
885    series: &mut Vec<HashMap<Column, String>>,
886    query_ctx: &QueryContext,
887    target: &PermissionTableTarget,
888    manager: &CatalogManagerRef,
889    metrics: &mut HashMap<String, u64>,
890) -> Result<()> {
891    let result = result?;
892
893    // fetch tag list
894    let table = manager
895        .table(
896            &target.catalog,
897            &target.schema,
898            &target.table,
899            Some(query_ctx),
900        )
901        .await?
902        .with_context(|| TableNotFoundSnafu {
903            catalog: &target.catalog,
904            schema: &target.schema,
905            table: &target.table,
906        })?;
907    let tag_columns = table
908        .primary_key_columns()
909        .map(|c| c.name)
910        .collect::<HashSet<_>>();
911
912    match result.data {
913        OutputData::RecordBatches(batches) => {
914            record_batches_to_series(batches, series, &target.table, &tag_columns)
915        }
916        OutputData::Stream(stream) => {
917            let batches = RecordBatches::try_collect(stream)
918                .await
919                .context(CollectRecordbatchSnafu)?;
920            record_batches_to_series(batches, series, &target.table, &tag_columns)
921        }
922        OutputData::AffectedRows(_) => Err(Error::UnexpectedResult {
923            reason: "expected data result, but got affected rows".to_string(),
924            location: Location::default(),
925        }),
926    }?;
927
928    if let Some(ref plan) = result.meta.plan {
929        collect_plan_metrics(plan, &mut [metrics]);
930    }
931    Ok(())
932}
933
934/// Retrieve labels name from query result
935async fn retrieve_labels_name_from_query_result(
936    result: Result<Output>,
937    labels: &mut HashSet<String>,
938    metrics: &mut HashMap<String, u64>,
939) -> Result<()> {
940    let result = result?;
941    match result.data {
942        OutputData::RecordBatches(batches) => record_batches_to_labels_name(batches, labels),
943        OutputData::Stream(stream) => {
944            let batches = RecordBatches::try_collect(stream)
945                .await
946                .context(CollectRecordbatchSnafu)?;
947            record_batches_to_labels_name(batches, labels)
948        }
949        OutputData::AffectedRows(_) => UnexpectedResultSnafu {
950            reason: "expected data result, but got affected rows".to_string(),
951        }
952        .fail(),
953    }?;
954    if let Some(ref plan) = result.meta.plan {
955        collect_plan_metrics(plan, &mut [metrics]);
956    }
957    Ok(())
958}
959
960fn record_batches_to_series(
961    batches: RecordBatches,
962    series: &mut Vec<HashMap<Column, String>>,
963    table_name: &str,
964    tag_columns: &HashSet<String>,
965) -> Result<()> {
966    for batch in batches.iter() {
967        // project record batch to only contains tag columns
968        let projection = batch
969            .schema
970            .column_schemas()
971            .iter()
972            .enumerate()
973            .filter_map(|(idx, col)| {
974                if tag_columns.contains(&col.name) {
975                    Some(idx)
976                } else {
977                    None
978                }
979            })
980            .collect::<Vec<_>>();
981        let batch = batch
982            .try_project(&projection)
983            .context(CollectRecordbatchSnafu)?;
984
985        let mut writer = RowWriter::new(&batch.schema, table_name);
986        writer.write(batch, series)?;
987    }
988    Ok(())
989}
990
991/// Writer from a row in the record batch to a Prometheus time series:
992///
993/// `{__name__="<metric name>", <label name>="<label value>", ...}`
994///
995/// The metrics name is the table name; label names are the column names and
996/// the label values are the corresponding row values (all are converted to strings).
997struct RowWriter {
998    /// The template that is to produce a Prometheus time series. It is pre-filled with metrics name
999    /// and label names, waiting to be filled by row values afterward.
1000    template: HashMap<Column, Option<String>>,
1001    /// The current filling row.
1002    current: Option<HashMap<Column, Option<String>>>,
1003}
1004
1005impl RowWriter {
1006    fn new(schema: &SchemaRef, table: &str) -> Self {
1007        let mut template = schema
1008            .column_schemas()
1009            .iter()
1010            .map(|x| (x.name.as_str().into(), None))
1011            .collect::<HashMap<Column, Option<String>>>();
1012        template.insert("__name__".into(), Some(table.to_string()));
1013        Self {
1014            template,
1015            current: None,
1016        }
1017    }
1018
1019    fn insert(&mut self, column: ColumnRef, value: impl ToString) {
1020        let current = self.current.get_or_insert_with(|| self.template.clone());
1021        match current.get_mut(&column as &dyn AsColumnRef) {
1022            Some(x) => {
1023                let _ = x.insert(value.to_string());
1024            }
1025            None => {
1026                let _ = current.insert(column.0.into(), Some(value.to_string()));
1027            }
1028        }
1029    }
1030
1031    fn insert_bytes(&mut self, column_schema: &ColumnSchema, bytes: &[u8]) -> Result<()> {
1032        let column_name = column_schema.name.as_str().into();
1033
1034        if column_schema.data_type.is_json() {
1035            let s = jsonb_to_string(bytes).context(ConvertScalarValueSnafu)?;
1036            self.insert(column_name, s);
1037        } else {
1038            let hex = bytes
1039                .iter()
1040                .map(|b| format!("{b:02x}"))
1041                .collect::<Vec<String>>()
1042                .join("");
1043            self.insert(column_name, hex);
1044        }
1045        Ok(())
1046    }
1047
1048    fn finish(&mut self) -> HashMap<Column, String> {
1049        let Some(current) = self.current.take() else {
1050            return HashMap::new();
1051        };
1052        current
1053            .into_iter()
1054            .filter_map(|(k, v)| v.map(|v| (k, v)))
1055            .collect()
1056    }
1057
1058    fn write(
1059        &mut self,
1060        record_batch: RecordBatch,
1061        series: &mut Vec<HashMap<Column, String>>,
1062    ) -> Result<()> {
1063        let schema = record_batch.schema.clone();
1064        let record_batch = record_batch.into_df_record_batch();
1065        for i in 0..record_batch.num_rows() {
1066            for (j, array) in record_batch.columns().iter().enumerate() {
1067                let column = schema.column_name_by_index(j).into();
1068
1069                if array.is_null(i) {
1070                    self.insert(column, "Null");
1071                    continue;
1072                }
1073
1074                match array.data_type() {
1075                    DataType::Null => {
1076                        self.insert(column, "Null");
1077                    }
1078                    DataType::Boolean => {
1079                        let array = array.as_boolean();
1080                        let v = array.value(i);
1081                        self.insert(column, v);
1082                    }
1083                    DataType::UInt8 => {
1084                        let array = array.as_primitive::<UInt8Type>();
1085                        let v = array.value(i);
1086                        self.insert(column, v);
1087                    }
1088                    DataType::UInt16 => {
1089                        let array = array.as_primitive::<UInt16Type>();
1090                        let v = array.value(i);
1091                        self.insert(column, v);
1092                    }
1093                    DataType::UInt32 => {
1094                        let array = array.as_primitive::<UInt32Type>();
1095                        let v = array.value(i);
1096                        self.insert(column, v);
1097                    }
1098                    DataType::UInt64 => {
1099                        let array = array.as_primitive::<UInt64Type>();
1100                        let v = array.value(i);
1101                        self.insert(column, v);
1102                    }
1103                    DataType::Int8 => {
1104                        let array = array.as_primitive::<Int8Type>();
1105                        let v = array.value(i);
1106                        self.insert(column, v);
1107                    }
1108                    DataType::Int16 => {
1109                        let array = array.as_primitive::<Int16Type>();
1110                        let v = array.value(i);
1111                        self.insert(column, v);
1112                    }
1113                    DataType::Int32 => {
1114                        let array = array.as_primitive::<Int32Type>();
1115                        let v = array.value(i);
1116                        self.insert(column, v);
1117                    }
1118                    DataType::Int64 => {
1119                        let array = array.as_primitive::<Int64Type>();
1120                        let v = array.value(i);
1121                        self.insert(column, v);
1122                    }
1123                    DataType::Float32 => {
1124                        let array = array.as_primitive::<Float32Type>();
1125                        let v = array.value(i);
1126                        self.insert(column, v);
1127                    }
1128                    DataType::Float64 => {
1129                        let array = array.as_primitive::<Float64Type>();
1130                        let v = array.value(i);
1131                        self.insert(column, v);
1132                    }
1133                    DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => {
1134                        let v = datatypes::arrow_array::string_array_value(array, i);
1135                        self.insert(column, v);
1136                    }
1137                    DataType::Binary | DataType::LargeBinary | DataType::BinaryView => {
1138                        let v = datatypes::arrow_array::binary_array_value(array, i);
1139                        let column_schema = &schema.column_schemas()[j];
1140                        self.insert_bytes(column_schema, v)?;
1141                    }
1142                    DataType::Date32 => {
1143                        let array = array.as_primitive::<Date32Type>();
1144                        let v = Date::new(array.value(i));
1145                        self.insert(column, v);
1146                    }
1147                    DataType::Date64 => {
1148                        let array = array.as_primitive::<Date64Type>();
1149                        // `Date64` values are milliseconds representation of `Date32` values,
1150                        // according to its specification. So we convert the `Date64` value here to
1151                        // the `Date32` value to process them unified.
1152                        let v = Date::new((array.value(i) / 86_400_000) as i32);
1153                        self.insert(column, v);
1154                    }
1155                    DataType::Timestamp(_, _) => {
1156                        let v = datatypes::arrow_array::timestamp_array_value(array, i);
1157                        self.insert(column, v.to_iso8601_string());
1158                    }
1159                    DataType::Time32(_) | DataType::Time64(_) => {
1160                        let v = datatypes::arrow_array::time_array_value(array, i);
1161                        self.insert(column, v.to_iso8601_string());
1162                    }
1163                    DataType::Interval(interval_unit) => match interval_unit {
1164                        IntervalUnit::YearMonth => {
1165                            let array = array.as_primitive::<IntervalYearMonthType>();
1166                            let v: IntervalYearMonth = array.value(i).into();
1167                            self.insert(column, v.to_iso8601_string());
1168                        }
1169                        IntervalUnit::DayTime => {
1170                            let array = array.as_primitive::<IntervalDayTimeType>();
1171                            let v: IntervalDayTime = array.value(i).into();
1172                            self.insert(column, v.to_iso8601_string());
1173                        }
1174                        IntervalUnit::MonthDayNano => {
1175                            let array = array.as_primitive::<IntervalMonthDayNanoType>();
1176                            let v: IntervalMonthDayNano = array.value(i).into();
1177                            self.insert(column, v.to_iso8601_string());
1178                        }
1179                    },
1180                    DataType::Duration(_) => {
1181                        let d = datatypes::arrow_array::duration_array_value(array, i);
1182                        self.insert(column, d);
1183                    }
1184                    DataType::List(_) => {
1185                        let v = ScalarValue::try_from_array(array, i).context(DataFusionSnafu)?;
1186                        self.insert(column, v);
1187                    }
1188                    DataType::Struct(_) => {
1189                        let v = ScalarValue::try_from_array(array, i).context(DataFusionSnafu)?;
1190                        self.insert(column, v);
1191                    }
1192                    DataType::Decimal128(precision, scale) => {
1193                        let array = array.as_primitive::<Decimal128Type>();
1194                        let v = Decimal128::new(array.value(i), *precision, *scale);
1195                        self.insert(column, v);
1196                    }
1197                    _ => {
1198                        return NotSupportedSnafu {
1199                            feat: format!("convert {} to http value", array.data_type()),
1200                        }
1201                        .fail();
1202                    }
1203                }
1204            }
1205
1206            series.push(self.finish())
1207        }
1208        Ok(())
1209    }
1210}
1211
1212#[derive(Debug, Clone, Copy, PartialEq)]
1213struct ColumnRef<'a>(&'a str);
1214
1215impl<'a> From<&'a str> for ColumnRef<'a> {
1216    fn from(s: &'a str) -> Self {
1217        Self(s)
1218    }
1219}
1220
1221trait AsColumnRef {
1222    fn as_ref(&self) -> ColumnRef<'_>;
1223}
1224
1225impl AsColumnRef for Column {
1226    fn as_ref(&self) -> ColumnRef<'_> {
1227        self.0.as_str().into()
1228    }
1229}
1230
1231impl AsColumnRef for ColumnRef<'_> {
1232    fn as_ref(&self) -> ColumnRef<'_> {
1233        *self
1234    }
1235}
1236
1237impl<'a> PartialEq for dyn AsColumnRef + 'a {
1238    fn eq(&self, other: &Self) -> bool {
1239        self.as_ref() == other.as_ref()
1240    }
1241}
1242
1243impl<'a> Eq for dyn AsColumnRef + 'a {}
1244
1245impl<'a> Hash for dyn AsColumnRef + 'a {
1246    fn hash<H: Hasher>(&self, state: &mut H) {
1247        self.as_ref().0.hash(state);
1248    }
1249}
1250
1251impl<'a> Borrow<dyn AsColumnRef + 'a> for Column {
1252    fn borrow(&self) -> &(dyn AsColumnRef + 'a) {
1253        self
1254    }
1255}
1256
1257/// Retrieve labels name from record batches
1258fn record_batches_to_labels_name(
1259    batches: RecordBatches,
1260    labels: &mut HashSet<String>,
1261) -> Result<()> {
1262    let mut column_indices = Vec::new();
1263    let mut value_column_indices = Vec::new();
1264    for (i, column) in batches.schema().column_schemas().iter().enumerate() {
1265        if is_prometheus_value_column(&column.data_type) {
1266            value_column_indices.push(i);
1267        }
1268        column_indices.push(i);
1269    }
1270
1271    if value_column_indices.is_empty() {
1272        return Err(Error::Internal {
1273            err_msg: "no value column found".to_string(),
1274        });
1275    }
1276
1277    for batch in batches.iter() {
1278        let names = column_indices
1279            .iter()
1280            .map(|c| batches.schema().column_name_by_index(*c).to_string())
1281            .collect::<Vec<_>>();
1282
1283        let value_columns = value_column_indices
1284            .iter()
1285            .map(|i| batch.column(*i))
1286            .collect::<Vec<_>>();
1287
1288        for row_index in 0..batch.num_rows() {
1289            // if all value columns are null, skip this row
1290            if value_columns.iter().all(|c| c.is_null(row_index)) {
1291                continue;
1292            }
1293
1294            // if a value is not null, record the tag name and return
1295            names.iter().for_each(|name| {
1296                let _ = labels.insert(name.clone());
1297            });
1298            return Ok(());
1299        }
1300    }
1301    Ok(())
1302}
1303
1304fn is_prometheus_value_column(data_type: &ConcreteDataType) -> bool {
1305    matches!(data_type, ConcreteDataType::Float64(_)) || is_native_histogram_value_type(data_type)
1306}
1307
1308pub(crate) fn retrieve_metric_name_and_result_type(
1309    promql_expr: &PromqlExpr,
1310) -> (Option<String>, ValueType) {
1311    (
1312        promql_expr_to_metric_name(promql_expr),
1313        promql_expr.value_type(),
1314    )
1315}
1316
1317/// Tries to get catalog and schema from an optional db param. And retrieves
1318/// them from [QueryContext] if they don't present.
1319pub(crate) fn get_catalog_schema(db: &Option<String>, ctx: &QueryContext) -> (String, String) {
1320    if let Some(db) = db {
1321        parse_catalog_and_schema_from_db_string(db)
1322    } else {
1323        (
1324            ctx.current_catalog().to_string(),
1325            ctx.current_schema().clone(),
1326        )
1327    }
1328}
1329
1330/// Update catalog and schema in [QueryContext] if necessary.
1331pub(crate) fn try_update_catalog_schema(ctx: &mut QueryContext, catalog: &str, schema: &str) {
1332    if ctx.current_catalog() != catalog || ctx.current_schema() != schema {
1333        ctx.set_current_catalog(catalog);
1334        ctx.set_current_schema(schema);
1335    }
1336}
1337
1338fn current_schema_metric_targets(
1339    query_ctx: &QueryContext,
1340    metric_names: &[String],
1341) -> PermissionTableTargets {
1342    PermissionTableTargets::resolved(
1343        metric_names
1344            .iter()
1345            .map(|metric| {
1346                PermissionTableTarget::new(
1347                    query_ctx.current_catalog(),
1348                    query_ctx.current_schema(),
1349                    metric,
1350                )
1351            })
1352            .collect(),
1353    )
1354}
1355
1356fn is_internal_physical_metric_table(table: &TableRef) -> bool {
1357    table.table_info().is_physical_table()
1358}
1359
1360fn promql_expr_to_metric_name(expr: &PromqlExpr) -> Option<String> {
1361    let mut metric_names = HashSet::new();
1362    collect_metric_names(expr, &mut metric_names);
1363
1364    // Return the metric name only if there's exactly one unique metric name
1365    if metric_names.len() == 1 {
1366        metric_names.into_iter().next()
1367    } else {
1368        None
1369    }
1370}
1371
1372/// Recursively collect all metric names from a PromQL expression
1373fn collect_metric_names(expr: &PromqlExpr, metric_names: &mut HashSet<String>) {
1374    match expr {
1375        PromqlExpr::Aggregate(AggregateExpr { modifier, expr, .. }) => {
1376            match modifier {
1377                Some(LabelModifier::Include(labels))
1378                    if !labels.labels.contains(&METRIC_NAME.to_string()) =>
1379                {
1380                    metric_names.clear();
1381                    return;
1382                }
1383                Some(LabelModifier::Exclude(labels))
1384                    if labels.labels.contains(&METRIC_NAME.to_string()) =>
1385                {
1386                    metric_names.clear();
1387                    return;
1388                }
1389                _ => {}
1390            }
1391            collect_metric_names(expr, metric_names)
1392        }
1393        PromqlExpr::Unary(UnaryExpr { .. }) => metric_names.clear(),
1394        PromqlExpr::Binary(BinaryExpr { lhs, op, .. }) => {
1395            if matches!(
1396                op.id(),
1397                token::T_LAND // INTERSECT
1398                    | token::T_LOR // UNION
1399                    | token::T_LUNLESS // EXCEPT
1400            ) {
1401                collect_metric_names(lhs, metric_names)
1402            } else {
1403                metric_names.clear()
1404            }
1405        }
1406        PromqlExpr::Paren(ParenExpr { expr }) => collect_metric_names(expr, metric_names),
1407        PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => collect_metric_names(expr, metric_names),
1408        PromqlExpr::VectorSelector(VectorSelector { name, matchers, .. }) => {
1409            if let Some(name) = name {
1410                metric_names.insert(name.clone());
1411            } else if let Some(matcher) = matchers.find_matchers(METRIC_NAME).into_iter().next() {
1412                metric_names.insert(matcher.value);
1413            }
1414        }
1415        PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1416            let VectorSelector { name, matchers, .. } = vs;
1417            if let Some(name) = name {
1418                metric_names.insert(name.clone());
1419            } else if let Some(matcher) = matchers.find_matchers(METRIC_NAME).into_iter().next() {
1420                metric_names.insert(matcher.value);
1421            }
1422        }
1423        PromqlExpr::Call(Call { args, .. }) => {
1424            args.args
1425                .iter()
1426                .for_each(|e| collect_metric_names(e, metric_names));
1427        }
1428        PromqlExpr::NumberLiteral(_) | PromqlExpr::StringLiteral(_) | PromqlExpr::Extension(_) => {}
1429    }
1430}
1431
1432fn find_metric_name_and_matchers<E, F>(expr: &PromqlExpr, f: F) -> Option<E>
1433where
1434    F: Fn(&Option<String>, &Matchers) -> Option<E> + Clone,
1435{
1436    match expr {
1437        PromqlExpr::Aggregate(AggregateExpr { expr, .. }) => find_metric_name_and_matchers(expr, f),
1438        PromqlExpr::Unary(UnaryExpr { expr }) => find_metric_name_and_matchers(expr, f),
1439        PromqlExpr::Binary(BinaryExpr { lhs, rhs, .. }) => {
1440            find_metric_name_and_matchers(lhs, f.clone()).or(find_metric_name_and_matchers(rhs, f))
1441        }
1442        PromqlExpr::Paren(ParenExpr { expr }) => find_metric_name_and_matchers(expr, f),
1443        PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => find_metric_name_and_matchers(expr, f),
1444        PromqlExpr::NumberLiteral(_) => None,
1445        PromqlExpr::StringLiteral(_) => None,
1446        PromqlExpr::Extension(_) => None,
1447        PromqlExpr::VectorSelector(VectorSelector { name, matchers, .. }) => f(name, matchers),
1448        PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1449            let VectorSelector { name, matchers, .. } = vs;
1450
1451            f(name, matchers)
1452        }
1453        PromqlExpr::Call(Call { args, .. }) => args
1454            .args
1455            .iter()
1456            .find_map(|e| find_metric_name_and_matchers(e, f.clone())),
1457    }
1458}
1459
1460struct MetricNameDiscovery {
1461    name_matchers: Vec<Matcher>,
1462    schema: Option<String>,
1463}
1464
1465/// Finds non-equality `__name__` matchers and their unambiguous schema.
1466fn find_metric_name_not_equal_matchers(expr: &PromqlExpr) -> Result<Option<MetricNameDiscovery>> {
1467    if find_metric_name_and_matchers(expr, |name, matchers| {
1468        (name.is_none() && !matchers.or_matchers.is_empty()).then_some(())
1469    })
1470    .is_some()
1471    {
1472        return Ok(None);
1473    }
1474
1475    find_metric_name_and_matchers(expr, |name, matchers| {
1476        // Has name, ignore the matchers
1477        if name.is_some() {
1478            return None;
1479        }
1480
1481        let name_matchers = matchers.find_matchers(METRIC_NAME);
1482        if name_matchers.len() != 1 || name_matchers[0].op == MatchOp::Equal {
1483            return None;
1484        }
1485
1486        Some(
1487            resolve_schema_from_matchers(&matchers.matchers).map(|schema| MetricNameDiscovery {
1488                name_matchers,
1489                schema,
1490            }),
1491        )
1492    })
1493    .transpose()
1494}
1495
1496fn expand_metric_name_queries(
1497    query: &ParsedPromQuery,
1498    metric_names: Vec<String>,
1499) -> Vec<ParsedPromQuery> {
1500    metric_names
1501        .into_iter()
1502        .map(|metric| {
1503            let mut query = query.clone();
1504            query.update_expr(|expr| update_metric_name_matcher(expr, &metric));
1505            query
1506        })
1507        .collect()
1508}
1509
1510fn static_promql_targets(
1511    expr: &PromqlExpr,
1512    query_ctx: &QueryContextRef,
1513) -> Result<(PermissionTableTargets, usize)> {
1514    let mut targets = Vec::new();
1515    let mut unresolved_selectors = 0;
1516    collect_static_promql_targets(expr, query_ctx, &mut targets, &mut unresolved_selectors)?;
1517    Ok((
1518        PermissionTableTargets::resolved(targets),
1519        unresolved_selectors,
1520    ))
1521}
1522
1523fn resolved_promql_targets(
1524    queries: &[ParsedPromQuery],
1525    query_ctx: &QueryContextRef,
1526) -> Option<Vec<PermissionTableTarget>> {
1527    let mut targets = Vec::new();
1528    for query in queries {
1529        let (PermissionTableTargets::Resolved(query_targets), unresolved_selectors) =
1530            static_promql_targets(query.expr(), query_ctx).ok()?
1531        else {
1532            return None;
1533        };
1534        if unresolved_selectors != 0 {
1535            return None;
1536        }
1537        for target in query_targets {
1538            if !targets.contains(&target) {
1539                targets.push(target);
1540            }
1541        }
1542    }
1543    Some(targets)
1544}
1545
1546fn collect_static_promql_targets(
1547    expr: &PromqlExpr,
1548    query_ctx: &QueryContextRef,
1549    targets: &mut Vec<PermissionTableTarget>,
1550    unresolved_selectors: &mut usize,
1551) -> Result<()> {
1552    match expr {
1553        PromqlExpr::Aggregate(AggregateExpr { expr, .. })
1554        | PromqlExpr::Unary(UnaryExpr { expr })
1555        | PromqlExpr::Paren(ParenExpr { expr })
1556        | PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => {
1557            collect_static_promql_targets(expr, query_ctx, targets, unresolved_selectors)?
1558        }
1559        PromqlExpr::Binary(BinaryExpr { lhs, rhs, .. }) => {
1560            collect_static_promql_targets(lhs, query_ctx, targets, unresolved_selectors)?;
1561            collect_static_promql_targets(rhs, query_ctx, targets, unresolved_selectors)?;
1562        }
1563        PromqlExpr::VectorSelector(selector) => {
1564            collect_static_vector_target(selector, query_ctx, targets, unresolved_selectors)?
1565        }
1566        PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1567            collect_static_vector_target(vs, query_ctx, targets, unresolved_selectors)?
1568        }
1569        PromqlExpr::Call(Call { args, .. }) => {
1570            for expr in &args.args {
1571                collect_static_promql_targets(expr, query_ctx, targets, unresolved_selectors)?;
1572            }
1573        }
1574        PromqlExpr::NumberLiteral(_) | PromqlExpr::StringLiteral(_) | PromqlExpr::Extension(_) => {}
1575    }
1576    Ok(())
1577}
1578
1579fn collect_static_vector_target(
1580    selector: &VectorSelector,
1581    query_ctx: &QueryContextRef,
1582    targets: &mut Vec<PermissionTableTarget>,
1583    unresolved_selectors: &mut usize,
1584) -> Result<()> {
1585    if selector.name.is_none() && !selector.matchers.or_matchers.is_empty() {
1586        *unresolved_selectors += 1;
1587        return Ok(());
1588    }
1589
1590    let metric = selector.name.clone().or_else(|| {
1591        let mut matchers = selector.matchers.find_matchers(METRIC_NAME);
1592        if matchers.len() == 1 && matchers[0].op == MatchOp::Equal {
1593            matchers.pop().map(|matcher| matcher.value)
1594        } else {
1595            None
1596        }
1597    });
1598    let Some(metric) = metric else {
1599        *unresolved_selectors += 1;
1600        return Ok(());
1601    };
1602
1603    let schema = resolve_schema_from_matchers(&selector.matchers.matchers)?
1604        .unwrap_or_else(|| query_ctx.current_schema());
1605    targets.push(PermissionTableTarget::new(
1606        query_ctx.current_catalog(),
1607        schema,
1608        metric,
1609    ));
1610    Ok(())
1611}
1612
1613/// Update the `__name__` matchers in expression into special value
1614/// Returns the updated expression.
1615fn update_metric_name_matcher(expr: &mut PromqlExpr, metric_name: &str) {
1616    match expr {
1617        PromqlExpr::Aggregate(AggregateExpr { expr, .. }) => {
1618            update_metric_name_matcher(expr, metric_name)
1619        }
1620        PromqlExpr::Unary(UnaryExpr { expr }) => update_metric_name_matcher(expr, metric_name),
1621        PromqlExpr::Binary(BinaryExpr { lhs, rhs, .. }) => {
1622            update_metric_name_matcher(lhs, metric_name);
1623            update_metric_name_matcher(rhs, metric_name);
1624        }
1625        PromqlExpr::Paren(ParenExpr { expr }) => update_metric_name_matcher(expr, metric_name),
1626        PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => {
1627            update_metric_name_matcher(expr, metric_name)
1628        }
1629        PromqlExpr::VectorSelector(VectorSelector { name, matchers, .. }) => {
1630            if name.is_some() {
1631                return;
1632            }
1633
1634            for m in &mut matchers.matchers {
1635                if m.name == METRIC_NAME && m.op != MatchOp::Equal {
1636                    m.op = MatchOp::Equal;
1637                    m.value = metric_name.to_string();
1638                }
1639            }
1640        }
1641        PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1642            let VectorSelector { name, matchers, .. } = vs;
1643            if name.is_some() {
1644                return;
1645            }
1646
1647            for m in &mut matchers.matchers {
1648                if m.name == METRIC_NAME && m.op != MatchOp::Equal {
1649                    m.op = MatchOp::Equal;
1650                    m.value = metric_name.to_string();
1651                }
1652            }
1653        }
1654        PromqlExpr::Call(Call { args, .. }) => {
1655            args.args.iter_mut().for_each(|e| {
1656                update_metric_name_matcher(e, metric_name);
1657            });
1658        }
1659        PromqlExpr::NumberLiteral(_) | PromqlExpr::StringLiteral(_) | PromqlExpr::Extension(_) => {}
1660    }
1661}
1662
1663#[derive(Debug, Default, Serialize, Deserialize)]
1664pub struct LabelValueQuery {
1665    start: Option<String>,
1666    end: Option<String>,
1667    lookback: Option<String>,
1668    #[serde(flatten)]
1669    matches: Matches,
1670    db: Option<String>,
1671    limit: Option<usize>,
1672}
1673
1674#[axum_macros::debug_handler]
1675#[tracing::instrument(
1676    skip_all,
1677    fields(protocol = "prometheus", request_type = "label_values_query")
1678)]
1679pub async fn label_values_query(
1680    State(handler): State<PrometheusHandlerRef>,
1681    Path(label_name): Path<String>,
1682    Extension(mut query_ctx): Extension<QueryContext>,
1683    Query(params): Query<LabelValueQuery>,
1684) -> PrometheusJsonResponse {
1685    let (catalog, schema) = get_catalog_schema(&params.db, &query_ctx);
1686    try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
1687    let query_ctx = Arc::new(query_ctx);
1688
1689    let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
1690        .with_label_values(&[query_ctx.get_db_string().as_str(), "label_values_query"])
1691        .start_timer();
1692
1693    let matches = params.matches.0;
1694    if label_name == METRIC_NAME_LABEL {
1695        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
1696        let exact_metric_names = matches
1697            .iter()
1698            .filter_map(|selector| retrieve_exact_metric_name_from_promql(selector))
1699            .collect::<Vec<_>>();
1700        try_call_return_response!(
1701            handler
1702                .check_query_target_permission(
1703                    current_schema_metric_targets(&query_ctx, &exact_metric_names),
1704                    &query_ctx,
1705                )
1706                .await
1707        );
1708        let catalog_manager = handler.catalog_manager();
1709
1710        let mut table_names = try_call_return_response!(
1711            retrieve_table_names(&query_ctx, catalog_manager, matches).await
1712        );
1713        table_names = try_call_return_response!(
1714            handler
1715                .filter_metadata_metric_names(
1716                    table_names,
1717                    query_ctx.current_schema().as_str(),
1718                    &query_ctx,
1719                )
1720                .await
1721        );
1722
1723        truncate_results(&mut table_names, params.limit);
1724        return PrometheusJsonResponse::success(PrometheusResponse::LabelValues(table_names));
1725    } else if label_name == FIELD_NAME_LABEL {
1726        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
1727        let enumerate_all = matches.is_empty();
1728        let metric_names = matches
1729            .iter()
1730            .map(|selector| retrieve_exact_metric_name_from_promql(selector))
1731            .collect::<Vec<_>>();
1732        if metric_names.iter().any(Option::is_none) {
1733            try_call_return_response!(
1734                handler
1735                    .check_query_target_permission(PermissionTableTargets::Unresolved, &query_ctx,)
1736                    .await
1737            );
1738        }
1739        let metric_names = metric_names.into_iter().flatten().collect::<Vec<_>>();
1740
1741        try_call_return_response!(
1742            handler
1743                .check_query_target_permission(
1744                    current_schema_metric_targets(&query_ctx, &metric_names),
1745                    &query_ctx,
1746                )
1747                .await
1748        );
1749
1750        let (field_columns, table_names) = if !enumerate_all && metric_names.is_empty() {
1751            (HashSet::new(), Vec::new())
1752        } else {
1753            handle_schema_err!(
1754                retrieve_field_names(&query_ctx, handler.catalog_manager(), metric_names).await
1755            )
1756            .unwrap_or_default()
1757        };
1758        try_call_return_response!(
1759            handler
1760                .check_query_target_permission(
1761                    current_schema_metric_targets(&query_ctx, &table_names),
1762                    &query_ctx,
1763                )
1764                .await
1765        );
1766        let mut field_columns = field_columns.into_iter().collect::<Vec<_>>();
1767        field_columns.sort_unstable();
1768        truncate_results(&mut field_columns, params.limit);
1769        return PrometheusJsonResponse::success(PrometheusResponse::LabelValues(field_columns));
1770    } else if is_database_selection_label(&label_name) {
1771        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
1772        let catalog_manager = handler.catalog_manager();
1773
1774        let (mut schema_names, targets) = try_call_return_response!(
1775            retrieve_schema_names(&query_ctx, catalog_manager, matches).await
1776        );
1777        try_call_return_response!(
1778            handler
1779                .check_query_target_permission(targets, &query_ctx)
1780                .await
1781        );
1782        truncate_results(&mut schema_names, params.limit);
1783        return PrometheusJsonResponse::success(PrometheusResponse::LabelValues(schema_names));
1784    }
1785
1786    let queries = matches;
1787    if queries.is_empty() {
1788        return PrometheusJsonResponse::error(
1789            StatusCode::InvalidArguments,
1790            "match[] parameter is required",
1791        );
1792    }
1793
1794    let start = params.start.unwrap_or_else(yesterday_rfc3339);
1795    let end = params.end.unwrap_or_else(current_time_rfc3339);
1796    let lookback = params
1797        .lookback
1798        .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string());
1799    let prom_queries = queries
1800        .into_iter()
1801        .map(|query| PromQuery {
1802            query,
1803            start: start.clone(),
1804            end: end.clone(),
1805            step: DEFAULT_LOOKBACK_STRING.to_string(),
1806            lookback: lookback.clone(),
1807            alias: None,
1808        })
1809        .collect::<Vec<_>>();
1810    let prom_queries = try_call_return_response!(
1811        prom_queries
1812            .into_iter()
1813            .map(|query| ParsedPromQuery::parse(query, &query_ctx))
1814            .collect::<Result<Vec<_>>>()
1815    );
1816    try_call_return_response!(
1817        handler
1818            .check_query_permission_parsed(&prom_queries, &query_ctx)
1819            .await
1820    );
1821
1822    let mut label_values = HashSet::new();
1823
1824    let Some(first_query) = prom_queries.first() else {
1825        unreachable!("empty match[] is rejected above")
1826    };
1827    let QueryStatement::Promql(eval_stmt, _) = first_query.statement() else {
1828        unreachable!("query is parsed from PromQL")
1829    };
1830    let start = eval_stmt.start;
1831    let end = eval_stmt.end;
1832
1833    for prom_query in prom_queries {
1834        let (_, statement) = prom_query.into_parts();
1835        let QueryStatement::Promql(eval_stmt, _) = statement else {
1836            unreachable!("query is parsed from PromQL")
1837        };
1838        let PromqlExpr::VectorSelector(mut vector_selector) = eval_stmt.expr else {
1839            return PrometheusJsonResponse::error(
1840                StatusCode::InvalidArguments,
1841                "expected vector selector",
1842            );
1843        };
1844        let Some(name) = take_metric_name(&mut vector_selector) else {
1845            return PrometheusJsonResponse::error(
1846                StatusCode::InvalidArguments,
1847                "expected metric name",
1848            );
1849        };
1850        let VectorSelector { matchers, .. } = vector_selector;
1851        // Only use and filter matchers.
1852        let matchers = matchers.matchers;
1853        let result = handler
1854            .query_label_values(name, label_name.clone(), matchers, start, end, &query_ctx)
1855            .await;
1856        if let Some(result) = handle_schema_err!(result) {
1857            label_values.extend(result.into_iter());
1858        }
1859    }
1860
1861    let mut label_values: Vec<_> = label_values.into_iter().collect();
1862    label_values.sort_unstable();
1863    truncate_results(&mut label_values, params.limit);
1864
1865    PrometheusJsonResponse::success(PrometheusResponse::LabelValues(label_values))
1866}
1867
1868fn truncate_results(label_values: &mut Vec<String>, limit: Option<usize>) {
1869    if let Some(limit) = limit
1870        && limit > 0
1871        && label_values.len() >= limit
1872    {
1873        label_values.truncate(limit);
1874    }
1875}
1876
1877/// Take metric name from the [VectorSelector].
1878/// It takes the name in the selector or removes the name matcher.
1879fn take_metric_name(selector: &mut VectorSelector) -> Option<String> {
1880    if let Some(name) = selector.name.take() {
1881        return Some(name);
1882    }
1883
1884    let (pos, matcher) = selector
1885        .matchers
1886        .matchers
1887        .iter()
1888        .find_position(|matcher| matcher.name == "__name__" && matcher.op == MatchOp::Equal)?;
1889    let name = matcher.value.clone();
1890    // We need to remove the name matcher to avoid using it as a filter in query.
1891    selector.matchers.matchers.remove(pos);
1892
1893    Some(name)
1894}
1895
1896async fn retrieve_table_names(
1897    query_ctx: &QueryContext,
1898    catalog_manager: CatalogManagerRef,
1899    matches: Vec<String>,
1900) -> Result<Vec<String>> {
1901    let catalog = query_ctx.current_catalog();
1902    let schema = query_ctx.current_schema();
1903
1904    let mut tables_stream = catalog_manager.tables(catalog, &schema, Some(query_ctx));
1905    let mut table_names = Vec::new();
1906
1907    let name_matchers = matches
1908        .iter()
1909        .map(|selector| {
1910            let expr = promql_parser::parser::parse(selector)
1911                .map_err(|reason| InvalidQuerySnafu { reason }.build())?;
1912            let PromqlExpr::VectorSelector(selector) = expr else {
1913                return InvalidQuerySnafu {
1914                    reason: "expected vector selector".to_string(),
1915                }
1916                .fail();
1917            };
1918            if !selector.matchers.or_matchers.is_empty() {
1919                return Ok(None);
1920            }
1921            if let Some(name) = selector.name {
1922                return Ok(Some(Matcher::new(MatchOp::Equal, METRIC_NAME_LABEL, &name)));
1923            }
1924            Ok(selector
1925                .matchers
1926                .matchers
1927                .into_iter()
1928                .find(|matcher| matcher.name == METRIC_NAME_LABEL))
1929        })
1930        .collect::<Result<Vec<_>>>()?;
1931
1932    while let Some(table) = tables_stream.next().await {
1933        let table = table?;
1934        if !table
1935            .table_info()
1936            .meta
1937            .options
1938            .extra_options
1939            .contains_key(LOGICAL_TABLE_METADATA_KEY)
1940            || is_internal_physical_metric_table(&table)
1941        {
1942            // skip non-prometheus (non-metricengine) tables for __name__ query
1943            continue;
1944        }
1945
1946        let table_name = &table.table_info().name;
1947
1948        if name_matchers.is_empty()
1949            || name_matchers.iter().any(|matcher| match matcher {
1950                None => true,
1951                Some(matcher) => match &matcher.op {
1952                    MatchOp::Equal => table_name == &matcher.value,
1953                    MatchOp::Re(reg) => reg.is_match(table_name),
1954                    // != and !~ are not supported, so include every table and
1955                    // let the caller authorize the complete candidate set.
1956                    _ => true,
1957                },
1958            })
1959        {
1960            table_names.push(table_name.clone());
1961        }
1962    }
1963
1964    table_names.sort_unstable();
1965    Ok(table_names)
1966}
1967
1968async fn retrieve_metric_metadata(
1969    query_ctx: &QueryContext,
1970    manager: CatalogManagerRef,
1971    params: &MetadataQuery,
1972) -> Result<BTreeMap<String, Vec<PromMetadata>>> {
1973    let mut metadata = BTreeMap::new();
1974    if params.limit == Some(0) {
1975        return Ok(metadata);
1976    }
1977
1978    let catalog = query_ctx.current_catalog();
1979    let schema = query_ctx.current_schema();
1980
1981    if let Some(metric) = &params.metric {
1982        let Some(table) = manager
1983            .table(catalog, &schema, metric, Some(query_ctx))
1984            .await?
1985        else {
1986            return Ok(metadata);
1987        };
1988        let table_info = table.table_info();
1989        if is_prometheus_metric_table(table_info.as_ref()) {
1990            metadata.insert(
1991                table_info.name.clone(),
1992                vec![prometheus_metadata_from_table(table_info.as_ref())],
1993            );
1994        }
1995        return Ok(metadata);
1996    }
1997
1998    let mut tables_stream = manager.tables(catalog, &schema, Some(query_ctx));
1999
2000    while let Some(table) = tables_stream.next().await {
2001        let table = table?;
2002        let table_info = table.table_info();
2003        if !is_prometheus_metric_table(table_info.as_ref()) {
2004            continue;
2005        }
2006
2007        metadata.insert(
2008            table_info.name.clone(),
2009            vec![prometheus_metadata_from_table(table_info.as_ref())],
2010        );
2011    }
2012
2013    Ok(metadata)
2014}
2015
2016fn is_prometheus_metric_table(table_info: &TableInfo) -> bool {
2017    table_info
2018        .meta
2019        .options
2020        .extra_options
2021        .contains_key(LOGICAL_TABLE_METADATA_KEY)
2022}
2023
2024fn prometheus_metadata_from_table(table_info: &TableInfo) -> PromMetadata {
2025    let options = &table_info.meta.options.extra_options;
2026    let metric_type = match options.get(SEMANTIC_METRIC_TYPE) {
2027        Some(metric_type)
2028            if options
2029                .get(SEMANTIC_METRIC_TEMPORALITY)
2030                .is_some_and(|temporality| {
2031                    matches!(
2032                        temporality.as_str(),
2033                        METRIC_TEMPORALITY_DELTA | SEMANTIC_VALUE_MIXED
2034                    )
2035                })
2036                && matches!(
2037                    metric_type.as_str(),
2038                    "counter" | "histogram" | "updown_counter"
2039                ) =>
2040        {
2041            "unknown".to_string()
2042        }
2043        Some(metric_type) => match metric_type.as_str() {
2044            "updown_counter" => "gauge".to_string(),
2045            "gauge_histogram" => "gaugehistogram".to_string(),
2046            "mixed" => "unknown".to_string(),
2047            metric_type => metric_type.to_string(),
2048        },
2049        None if table_has_native_histogram_value(table_info) => "histogram".to_string(),
2050        None => String::new(),
2051    };
2052    let unit = options
2053        .get(SEMANTIC_METRIC_UNIT)
2054        .map(|unit| ucum_to_openmetrics_unit(unit))
2055        .unwrap_or_default();
2056
2057    PromMetadata {
2058        metric_type,
2059        unit,
2060        // TODO: Persist and return Prometheus help text and OTLP metric descriptions.
2061        help: String::new(),
2062    }
2063}
2064
2065fn table_has_native_histogram_value(table_info: &TableInfo) -> bool {
2066    table_info
2067        .meta
2068        .schema
2069        .column_schemas()
2070        .iter()
2071        .any(|column| is_native_histogram_value_type(&column.data_type))
2072}
2073
2074async fn retrieve_field_names(
2075    query_ctx: &QueryContext,
2076    manager: CatalogManagerRef,
2077    matches: Vec<String>,
2078) -> Result<(HashSet<String>, Vec<String>)> {
2079    let mut field_columns = HashSet::new();
2080    let mut table_names = Vec::new();
2081    let catalog = query_ctx.current_catalog();
2082    let schema = query_ctx.current_schema();
2083
2084    if matches.is_empty() {
2085        // query all tables if no matcher is provided
2086        let mut tables = manager.tables(catalog, &schema, Some(query_ctx));
2087        while let Some(table) = tables.next().await {
2088            let table = table?;
2089            if is_internal_physical_metric_table(&table) {
2090                continue;
2091            }
2092            table_names.push(table.table_info().name.clone());
2093            for column in table.field_columns() {
2094                field_columns.insert(column.name);
2095            }
2096        }
2097        return Ok((field_columns, table_names));
2098    }
2099
2100    for table_name in matches {
2101        let table = manager
2102            .table(catalog, &schema, &table_name, Some(query_ctx))
2103            .await?
2104            .with_context(|| TableNotFoundSnafu {
2105                catalog: catalog.to_string(),
2106                schema: schema.clone(),
2107                table: table_name.clone(),
2108            })?;
2109
2110        if is_internal_physical_metric_table(&table) {
2111            continue;
2112        }
2113        table_names.push(table.table_info().name.clone());
2114        for column in table.field_columns() {
2115            field_columns.insert(column.name);
2116        }
2117    }
2118    Ok((field_columns, table_names))
2119}
2120
2121async fn retrieve_schema_names(
2122    query_ctx: &QueryContext,
2123    catalog_manager: CatalogManagerRef,
2124    matches: Vec<String>,
2125) -> Result<(Vec<String>, PermissionTableTargets)> {
2126    let mut schemas = Vec::new();
2127    let mut targets = Vec::new();
2128    let catalog = query_ctx.current_catalog();
2129    let metric_names = matches
2130        .iter()
2131        .map(|match_item| retrieve_exact_metric_name_from_promql(match_item))
2132        .collect::<Vec<_>>();
2133    let unresolved = metric_names.is_empty() || metric_names.iter().any(Option::is_none);
2134
2135    let candidate_schemas = catalog_manager
2136        .schema_names(catalog, Some(query_ctx))
2137        .await?;
2138
2139    for schema in candidate_schemas {
2140        let mut found = true;
2141        for table_name in metric_names.iter().flatten() {
2142            let table = catalog_manager
2143                .table(catalog, &schema, table_name, Some(query_ctx))
2144                .await?;
2145            let Some(table) = table else {
2146                found = false;
2147                continue;
2148            };
2149            targets.push(PermissionTableTarget::new(catalog, &schema, table_name));
2150            if is_internal_physical_metric_table(&table) {
2151                found = false;
2152            }
2153        }
2154
2155        if found {
2156            schemas.push(schema);
2157        }
2158    }
2159
2160    schemas.sort_unstable();
2161
2162    let targets = if unresolved {
2163        PermissionTableTargets::Unresolved
2164    } else {
2165        PermissionTableTargets::resolved(targets)
2166    };
2167    Ok((schemas, targets))
2168}
2169
2170fn retrieve_exact_metric_name_from_promql(query: &str) -> Option<String> {
2171    let PromqlExpr::VectorSelector(selector) = promql_parser::parser::parse(query).ok()? else {
2172        return None;
2173    };
2174    if let Some(name) = selector.name {
2175        return (!name.is_empty()).then_some(name);
2176    }
2177
2178    let mut name_matchers = selector.matchers.find_matchers(METRIC_NAME);
2179    if name_matchers.len() != 1 {
2180        return None;
2181    }
2182    let matcher = name_matchers.pop()?;
2183    (matcher.op == MatchOp::Equal && !matcher.value.is_empty()).then_some(matcher.value)
2184}
2185
2186#[derive(Debug, Default, Serialize, Deserialize)]
2187pub struct SeriesQuery {
2188    start: Option<String>,
2189    end: Option<String>,
2190    lookback: Option<String>,
2191    #[serde(flatten)]
2192    matches: Matches,
2193    db: Option<String>,
2194}
2195
2196#[axum_macros::debug_handler]
2197#[tracing::instrument(
2198    skip_all,
2199    fields(protocol = "prometheus", request_type = "series_query")
2200)]
2201pub async fn series_query(
2202    State(handler): State<PrometheusHandlerRef>,
2203    Query(params): Query<SeriesQuery>,
2204    Extension(mut query_ctx): Extension<QueryContext>,
2205    Form(form_params): Form<SeriesQuery>,
2206) -> PrometheusJsonResponse {
2207    let mut queries: Vec<String> = params.matches.0;
2208    if queries.is_empty() {
2209        queries = form_params.matches.0;
2210    }
2211    if queries.is_empty() {
2212        return PrometheusJsonResponse::error(
2213            StatusCode::Unsupported,
2214            "match[] parameter is required",
2215        );
2216    }
2217    let start = params
2218        .start
2219        .or(form_params.start)
2220        .unwrap_or_else(yesterday_rfc3339);
2221    let end = params
2222        .end
2223        .or(form_params.end)
2224        .unwrap_or_else(current_time_rfc3339);
2225    let lookback = params
2226        .lookback
2227        .or(form_params.lookback)
2228        .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string());
2229
2230    // update catalog and schema in query context if necessary
2231    if let Some(db) = &params.db {
2232        let (catalog, schema) = parse_catalog_and_schema_from_db_string(db);
2233        try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
2234    }
2235    let query_ctx = Arc::new(query_ctx);
2236
2237    let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
2238        .with_label_values(&[query_ctx.get_db_string().as_str(), "series_query"])
2239        .start_timer();
2240
2241    let prom_queries = queries
2242        .into_iter()
2243        .map(|query| PromQuery {
2244            query,
2245            start: start.clone(),
2246            end: end.clone(),
2247            // TODO: find a better value for step
2248            step: DEFAULT_LOOKBACK_STRING.to_string(),
2249            lookback: lookback.clone(),
2250            alias: None,
2251        })
2252        .collect::<Vec<_>>();
2253    let mut expanded_queries = Vec::new();
2254    for prom_query in prom_queries {
2255        let prom_query = try_call_return_response!(
2256            @output_msg ParsedPromQuery::parse(prom_query, &query_ctx),
2257            StatusCode::InvalidArguments
2258        );
2259        let promql_expr = prom_query.expr();
2260        let Some(discovery) =
2261            try_call_return_response!(find_metric_name_not_equal_matchers(promql_expr))
2262        else {
2263            expanded_queries.push(prom_query);
2264            continue;
2265        };
2266
2267        try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
2268        let (static_targets, unresolved_selectors) =
2269            try_call_return_response!(static_promql_targets(promql_expr, &query_ctx));
2270        try_call_return_response!(
2271            handler
2272                .check_query_target_permission(static_targets, &query_ctx)
2273                .await
2274        );
2275        if unresolved_selectors != 1 {
2276            try_call_return_response!(
2277                handler
2278                    .check_query_target_permission(PermissionTableTargets::Unresolved, &query_ctx)
2279                    .await
2280            );
2281            expanded_queries.push(prom_query);
2282            continue;
2283        }
2284
2285        let schema = discovery
2286            .schema
2287            .unwrap_or_else(|| query_ctx.current_schema());
2288        let metric_names = try_call_return_response!(
2289            handler
2290                .query_metric_names(discovery.name_matchers, &schema, &query_ctx)
2291                .await
2292        );
2293        expanded_queries.extend(expand_metric_name_queries(&prom_query, metric_names));
2294    }
2295    let prom_queries = expanded_queries;
2296    try_call_return_response!(
2297        handler
2298            .check_query_permission_parsed(&prom_queries, &query_ctx)
2299            .await
2300    );
2301
2302    let mut series = Vec::new();
2303    let mut merge_map = HashMap::new();
2304    for prom_query in prom_queries {
2305        let Some(target) = resolved_promql_targets(std::slice::from_ref(&prom_query), &query_ctx)
2306            .and_then(|targets| match targets.as_slice() {
2307                [target] => Some(target.clone()),
2308                _ => None,
2309            })
2310        else {
2311            return PrometheusJsonResponse::error(
2312                StatusCode::InvalidArguments,
2313                "series selector must resolve exactly one metric table",
2314            );
2315        };
2316        let result = handler.do_query_parsed(prom_query, query_ctx.clone()).await;
2317
2318        handle_schema_err!(
2319            retrieve_series_from_query_result(
2320                result,
2321                &mut series,
2322                &query_ctx,
2323                &target,
2324                &handler.catalog_manager(),
2325                &mut merge_map,
2326            )
2327            .await
2328        );
2329    }
2330    let merge_map = merge_map
2331        .into_iter()
2332        .map(|(k, v)| (k, Value::from(v)))
2333        .collect();
2334    let mut resp = PrometheusJsonResponse::success(PrometheusResponse::Series(series));
2335    resp.resp_metrics = merge_map;
2336    resp
2337}
2338
2339#[derive(Debug, Default, Serialize, Deserialize)]
2340pub struct ParseQuery {
2341    query: Option<String>,
2342    db: Option<String>,
2343}
2344
2345#[axum_macros::debug_handler]
2346#[tracing::instrument(
2347    skip_all,
2348    fields(protocol = "prometheus", request_type = "parse_query")
2349)]
2350pub async fn parse_query(
2351    State(_handler): State<PrometheusHandlerRef>,
2352    Query(params): Query<ParseQuery>,
2353    Extension(_query_ctx): Extension<QueryContext>,
2354    Form(form_params): Form<ParseQuery>,
2355) -> PrometheusJsonResponse {
2356    if let Some(query) = params.query.or(form_params.query) {
2357        let ast = try_call_return_response!(
2358            promql_parser::parser::parse(&query),
2359            StatusCode::InvalidArguments
2360        );
2361        PrometheusJsonResponse::success(PrometheusResponse::ParseResult(ast))
2362    } else {
2363        PrometheusJsonResponse::error(StatusCode::InvalidArguments, "query is required")
2364    }
2365}
2366
2367#[cfg(test)]
2368mod tests {
2369    use std::collections::HashSet;
2370    use std::sync::{Arc, Mutex};
2371
2372    use catalog::memory::MemoryCatalogManager;
2373    use catalog::{RegisterSchemaRequest, RegisterTableRequest};
2374    use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
2375    use common_query::native_histogram::{
2376        CounterResetHint, NativeHistogram, build_histogram_array, native_histogram_value_type,
2377    };
2378    use common_query::prelude::greptime_native_histogram;
2379    use datatypes::prelude::ConcreteDataType;
2380    use datatypes::schema::{ColumnSchema, Schema};
2381    use datatypes::vectors::{StringVector, StructVector};
2382    use promql_parser::parser::value::ValueType;
2383    use table::metadata::{TableInfoBuilder, TableMetaBuilder, TableType, TableVersion};
2384    use table::requests::{
2385        SEMANTIC_METRIC_TEMPORALITY, SEMANTIC_METRIC_TYPE, SEMANTIC_METRIC_UNIT, TableOptions,
2386    };
2387    use table::test_util::EmptyTable;
2388    use table::test_util::table_info::test_table_info;
2389
2390    use super::*;
2391    use crate::prometheus_handler::PrometheusHandler;
2392
2393    struct TestCase {
2394        name: &'static str,
2395        promql: &'static str,
2396        expected_metric: Option<&'static str>,
2397        expected_type: ValueType,
2398        should_error: bool,
2399    }
2400
2401    struct TestPrometheusHandler {
2402        catalog_manager: CatalogManagerRef,
2403        deny_operation: bool,
2404        denied_table: Option<&'static str>,
2405        metric_names: Vec<String>,
2406        queries: Mutex<Vec<String>>,
2407    }
2408
2409    #[async_trait::async_trait]
2410    impl PrometheusHandler for TestPrometheusHandler {
2411        async fn do_query(&self, _: &PromQuery, _: QueryContextRef) -> Result<Output> {
2412            unreachable!("HTTP handlers should execute the parsed query")
2413        }
2414
2415        async fn do_query_parsed(
2416            &self,
2417            query: ParsedPromQuery,
2418            _: QueryContextRef,
2419        ) -> Result<Output> {
2420            self.queries
2421                .lock()
2422                .unwrap()
2423                .push(query.query().query.clone());
2424            Ok(Output::new_with_record_batches(RecordBatches::empty()))
2425        }
2426
2427        async fn check_query_permission(&self, _: &[PromQuery], _: &QueryContextRef) -> Result<()> {
2428            if self.deny_operation {
2429                return auth::error::PermissionDeniedSnafu
2430                    .fail()
2431                    .context(crate::error::AuthSnafu);
2432            }
2433            Ok(())
2434        }
2435
2436        async fn check_query_target_permission(
2437            &self,
2438            targets: PermissionTableTargets,
2439            _: &QueryContextRef,
2440        ) -> Result<()> {
2441            if matches!(
2442                targets,
2443                PermissionTableTargets::Resolved(targets)
2444                    if self.denied_table.is_some_and(|denied| {
2445                        targets.iter().any(|target| target.table == denied)
2446                    })
2447            ) {
2448                return auth::error::PermissionDeniedSnafu
2449                    .fail()
2450                    .context(crate::error::AuthSnafu);
2451            }
2452            Ok(())
2453        }
2454
2455        async fn filter_metadata_metric_names(
2456            &self,
2457            metric_names: Vec<String>,
2458            _: &str,
2459            _: &QueryContextRef,
2460        ) -> Result<Vec<String>> {
2461            Ok(metric_names
2462                .into_iter()
2463                .filter(|metric| metric != "denied")
2464                .collect())
2465        }
2466
2467        async fn query_metric_names(
2468            &self,
2469            _: Vec<Matcher>,
2470            _: &str,
2471            _: &QueryContextRef,
2472        ) -> Result<Vec<String>> {
2473            Ok(self.metric_names.clone())
2474        }
2475
2476        async fn query_label_values(
2477            &self,
2478            _: String,
2479            _: String,
2480            _: Vec<Matcher>,
2481            _: std::time::SystemTime,
2482            _: std::time::SystemTime,
2483            _: &QueryContextRef,
2484        ) -> Result<Vec<String>> {
2485            unreachable!()
2486        }
2487
2488        fn catalog_manager(&self) -> CatalogManagerRef {
2489            self.catalog_manager.clone()
2490        }
2491    }
2492
2493    fn permission_denied_response() -> PrometheusJsonResponse {
2494        let result: auth::error::Result<()> = auth::error::PermissionDeniedSnafu.fail();
2495        try_call_return_response!(result);
2496        unreachable!()
2497    }
2498
2499    fn invalid_promql_response() -> PrometheusJsonResponse {
2500        let result = ParsedPromQuery::parse(
2501            PromQuery {
2502                query: "up{".to_string(),
2503                ..Default::default()
2504            },
2505            &QueryContext::arc(),
2506        );
2507        try_call_return_response!(@output_msg result, StatusCode::InvalidArguments);
2508        unreachable!()
2509    }
2510
2511    #[test]
2512    fn test_try_call_return_response_preserves_status() {
2513        use axum::response::IntoResponse;
2514
2515        let response = permission_denied_response();
2516        assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
2517        assert_eq!(
2518            axum::http::StatusCode::FORBIDDEN,
2519            response.into_response().status()
2520        );
2521    }
2522
2523    #[test]
2524    fn test_try_call_return_response_preserves_error_source() {
2525        let response = invalid_promql_response();
2526        assert_eq!(Some(StatusCode::InvalidArguments), response.status_code);
2527        let error = response.error.unwrap();
2528        assert!(
2529            error.contains("unexpected end of input inside braces"),
2530            "{error}"
2531        );
2532    }
2533
2534    #[tokio::test]
2535    async fn test_promql_timer_records_parse_errors() {
2536        let handler: PrometheusHandlerRef = Arc::new(TestPrometheusHandler {
2537            catalog_manager: MemoryCatalogManager::new(),
2538            deny_operation: false,
2539            denied_table: None,
2540            metric_names: Vec::new(),
2541            queries: Mutex::new(Vec::new()),
2542        });
2543        let query_ctx = QueryContext::with("promql_timer_test", "parse_error");
2544        let db = query_ctx.get_db_string();
2545
2546        let instant_histogram = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
2547            .with_label_values(&[db.as_str(), "instant_query"]);
2548        let instant_count = instant_histogram.get_sample_count();
2549        let response = instant_query(
2550            State(handler.clone()),
2551            Query(InstantQuery {
2552                query: Some("up{".to_string()),
2553                ..Default::default()
2554            }),
2555            Extension(query_ctx),
2556            Form(InstantQuery::default()),
2557        )
2558        .await;
2559        assert_eq!(Some(StatusCode::InvalidArguments), response.status_code);
2560        assert_eq!(instant_count + 1, instant_histogram.get_sample_count());
2561
2562        let range_histogram = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
2563            .with_label_values(&[db.as_str(), "range_query"]);
2564        let range_count = range_histogram.get_sample_count();
2565        let response = range_query(
2566            State(handler),
2567            Query(RangeQuery {
2568                query: Some("up{".to_string()),
2569                start: Some("0".to_string()),
2570                end: Some("1".to_string()),
2571                step: Some("1s".to_string()),
2572                ..Default::default()
2573            }),
2574            Extension(QueryContext::with("promql_timer_test", "parse_error")),
2575            Form(RangeQuery::default()),
2576        )
2577        .await;
2578        assert_eq!(Some(StatusCode::InvalidArguments), response.status_code);
2579        assert_eq!(range_count + 1, range_histogram.get_sample_count());
2580    }
2581
2582    #[tokio::test]
2583    async fn test_field_name_enumeration_fails_closed() {
2584        use axum::response::IntoResponse;
2585
2586        let mut allowed = test_table_info(
2587            1024,
2588            "allowed",
2589            DEFAULT_SCHEMA_NAME,
2590            DEFAULT_CATALOG_NAME,
2591            Arc::new(Schema::new(vec![])),
2592        );
2593        allowed.meta.options.extra_options.insert(
2594            LOGICAL_TABLE_METADATA_KEY.to_string(),
2595            "physical_metrics".to_string(),
2596        );
2597        let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&allowed));
2598        let mut denied = allowed.clone();
2599        denied.ident.table_id = 1025;
2600        denied.name = "denied".to_string();
2601        manager
2602            .register_table_sync(RegisterTableRequest {
2603                catalog: DEFAULT_CATALOG_NAME.to_string(),
2604                schema: DEFAULT_SCHEMA_NAME.to_string(),
2605                table_name: denied.name.clone(),
2606                table_id: denied.table_id(),
2607                table: EmptyTable::from_table_info(&denied),
2608            })
2609            .unwrap();
2610        let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
2611        assert_eq!(
2612            vec!["allowed".to_string(), "denied".to_string()],
2613            retrieve_table_names(&query_ctx, manager.clone(), Vec::new())
2614                .await
2615                .unwrap()
2616        );
2617
2618        let response = label_values_query(
2619            State(Arc::new(TestPrometheusHandler {
2620                catalog_manager: manager,
2621                deny_operation: false,
2622                denied_table: Some("denied"),
2623                metric_names: Vec::new(),
2624                queries: Mutex::new(Vec::new()),
2625            })),
2626            Path(FIELD_NAME_LABEL.to_string()),
2627            Extension(query_ctx),
2628            Query(LabelValueQuery::default()),
2629        )
2630        .await;
2631
2632        assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
2633        assert_eq!(
2634            axum::http::StatusCode::FORBIDDEN,
2635            response.into_response().status()
2636        );
2637    }
2638
2639    #[tokio::test]
2640    async fn test_series_query_expands_metric_name_regex() {
2641        let cpu_user = test_table_info(
2642            1024,
2643            "cpu_user",
2644            DEFAULT_SCHEMA_NAME,
2645            DEFAULT_CATALOG_NAME,
2646            Arc::new(Schema::new(vec![])),
2647        );
2648        let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&cpu_user));
2649        let mut cpu_system = cpu_user.clone();
2650        cpu_system.ident.table_id = 1025;
2651        cpu_system.name = "cpu_system".to_string();
2652        manager
2653            .register_table_sync(RegisterTableRequest {
2654                catalog: DEFAULT_CATALOG_NAME.to_string(),
2655                schema: DEFAULT_SCHEMA_NAME.to_string(),
2656                table_name: cpu_system.name.clone(),
2657                table_id: cpu_system.table_id(),
2658                table: EmptyTable::from_table_info(&cpu_system),
2659            })
2660            .unwrap();
2661
2662        let handler = Arc::new(TestPrometheusHandler {
2663            catalog_manager: manager,
2664            deny_operation: false,
2665            denied_table: None,
2666            metric_names: vec!["cpu_user".to_string(), "cpu_system".to_string()],
2667            queries: Mutex::new(Vec::new()),
2668        });
2669        let state: PrometheusHandlerRef = handler.clone();
2670        let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
2671        let target_ctx = Arc::new(query_ctx.clone());
2672        let response = series_query(
2673            State(state),
2674            Query(SeriesQuery {
2675                matches: Matches(vec![r#"{__name__=~"cpu_.*"}"#.to_string()]),
2676                ..Default::default()
2677            }),
2678            Extension(query_ctx),
2679            Form(SeriesQuery::default()),
2680        )
2681        .await;
2682
2683        assert!(
2684            response.status_code.is_none(),
2685            "status={:?}, error={:?}",
2686            response.status_code,
2687            response.error
2688        );
2689        let mut targets = handler
2690            .queries
2691            .lock()
2692            .unwrap()
2693            .iter()
2694            .map(|query| {
2695                let query = ParsedPromQuery::parse(
2696                    PromQuery {
2697                        query: query.clone(),
2698                        ..Default::default()
2699                    },
2700                    &target_ctx,
2701                )
2702                .unwrap();
2703                resolved_promql_targets(std::slice::from_ref(&query), &target_ctx)
2704                    .unwrap()
2705                    .pop()
2706                    .unwrap()
2707                    .table
2708            })
2709            .collect_vec();
2710        targets.sort_unstable();
2711        assert_eq!(vec!["cpu_system", "cpu_user"], targets);
2712    }
2713
2714    #[test]
2715    fn test_retrieve_metric_name_and_result_type() {
2716        let test_cases = &[
2717            // Single metric cases
2718            TestCase {
2719                name: "simple metric",
2720                promql: "cpu_usage",
2721                expected_metric: Some("cpu_usage"),
2722                expected_type: ValueType::Vector,
2723                should_error: false,
2724            },
2725            TestCase {
2726                name: "metric with selector",
2727                promql: r#"cpu_usage{instance="localhost"}"#,
2728                expected_metric: Some("cpu_usage"),
2729                expected_type: ValueType::Vector,
2730                should_error: false,
2731            },
2732            TestCase {
2733                name: "metric with range selector",
2734                promql: "cpu_usage[5m]",
2735                expected_metric: Some("cpu_usage"),
2736                expected_type: ValueType::Matrix,
2737                should_error: false,
2738            },
2739            TestCase {
2740                name: "metric with __name__ matcher",
2741                promql: r#"{__name__="cpu_usage"}"#,
2742                expected_metric: Some("cpu_usage"),
2743                expected_type: ValueType::Vector,
2744                should_error: false,
2745            },
2746            TestCase {
2747                name: "metric with unary operator",
2748                promql: "-cpu_usage",
2749                expected_metric: None,
2750                expected_type: ValueType::Vector,
2751                should_error: false,
2752            },
2753            // Aggregation and function cases
2754            TestCase {
2755                name: "metric with aggregation",
2756                promql: "sum(cpu_usage)",
2757                expected_metric: Some("cpu_usage"),
2758                expected_type: ValueType::Vector,
2759                should_error: false,
2760            },
2761            TestCase {
2762                name: "complex aggregation",
2763                promql: r#"sum by (instance) (cpu_usage{job="node"})"#,
2764                expected_metric: None,
2765                expected_type: ValueType::Vector,
2766                should_error: false,
2767            },
2768            TestCase {
2769                name: "complex aggregation",
2770                promql: r#"sum by (__name__) (cpu_usage{job="node"})"#,
2771                expected_metric: Some("cpu_usage"),
2772                expected_type: ValueType::Vector,
2773                should_error: false,
2774            },
2775            TestCase {
2776                name: "complex aggregation",
2777                promql: r#"sum without (instance) (cpu_usage{job="node"})"#,
2778                expected_metric: Some("cpu_usage"),
2779                expected_type: ValueType::Vector,
2780                should_error: false,
2781            },
2782            // Same metric binary operations
2783            TestCase {
2784                name: "same metric addition",
2785                promql: "cpu_usage + cpu_usage",
2786                expected_metric: None,
2787                expected_type: ValueType::Vector,
2788                should_error: false,
2789            },
2790            TestCase {
2791                name: "metric with scalar addition",
2792                promql: r#"sum(rate(cpu_usage{job="node"}[5m])) + 100"#,
2793                expected_metric: None,
2794                expected_type: ValueType::Vector,
2795                should_error: false,
2796            },
2797            // Multiple metrics cases
2798            TestCase {
2799                name: "different metrics addition",
2800                promql: "cpu_usage + memory_usage",
2801                expected_metric: None,
2802                expected_type: ValueType::Vector,
2803                should_error: false,
2804            },
2805            TestCase {
2806                name: "different metrics subtraction",
2807                promql: "network_in - network_out",
2808                expected_metric: None,
2809                expected_type: ValueType::Vector,
2810                should_error: false,
2811            },
2812            // Unless operator cases
2813            TestCase {
2814                name: "unless with different metrics",
2815                promql: "cpu_usage unless memory_usage",
2816                expected_metric: Some("cpu_usage"),
2817                expected_type: ValueType::Vector,
2818                should_error: false,
2819            },
2820            TestCase {
2821                name: "unless with same metric",
2822                promql: "cpu_usage unless cpu_usage",
2823                expected_metric: Some("cpu_usage"),
2824                expected_type: ValueType::Vector,
2825                should_error: false,
2826            },
2827            // Subquery cases
2828            TestCase {
2829                name: "basic subquery",
2830                promql: "cpu_usage[5m:1m]",
2831                expected_metric: Some("cpu_usage"),
2832                expected_type: ValueType::Matrix,
2833                should_error: false,
2834            },
2835            TestCase {
2836                name: "subquery with multiple metrics",
2837                promql: "(cpu_usage + memory_usage)[5m:1m]",
2838                expected_metric: None,
2839                expected_type: ValueType::Matrix,
2840                should_error: false,
2841            },
2842            // Literal values
2843            TestCase {
2844                name: "scalar value",
2845                promql: "42",
2846                expected_metric: None,
2847                expected_type: ValueType::Scalar,
2848                should_error: false,
2849            },
2850            TestCase {
2851                name: "string literal",
2852                promql: r#""hello world""#,
2853                expected_metric: None,
2854                expected_type: ValueType::String,
2855                should_error: false,
2856            },
2857            // Error cases
2858            TestCase {
2859                name: "invalid syntax",
2860                promql: "cpu_usage{invalid=",
2861                expected_metric: None,
2862                expected_type: ValueType::Vector,
2863                should_error: true,
2864            },
2865            TestCase {
2866                name: "empty query",
2867                promql: "",
2868                expected_metric: None,
2869                expected_type: ValueType::Vector,
2870                should_error: true,
2871            },
2872            TestCase {
2873                name: "malformed brackets",
2874                promql: "cpu_usage[5m",
2875                expected_metric: None,
2876                expected_type: ValueType::Vector,
2877                should_error: true,
2878            },
2879        ];
2880
2881        for test_case in test_cases {
2882            let result = promql_parser::parser::parse(test_case.promql)
2883                .map(|expr| retrieve_metric_name_and_result_type(&expr));
2884
2885            if test_case.should_error {
2886                assert!(
2887                    result.is_err(),
2888                    "Test '{}' should have failed but succeeded with: {:?}",
2889                    test_case.name,
2890                    result
2891                );
2892            } else {
2893                let (metric_name, value_type) = result.unwrap_or_else(|e| {
2894                    panic!(
2895                        "Test '{}' should have succeeded but failed with error: {}",
2896                        test_case.name, e
2897                    )
2898                });
2899
2900                let expected_metric_name = test_case.expected_metric.map(|s| s.to_string());
2901                assert_eq!(
2902                    metric_name, expected_metric_name,
2903                    "Test '{}': metric name mismatch. Expected: {:?}, Got: {:?}",
2904                    test_case.name, expected_metric_name, metric_name
2905                );
2906
2907                assert_eq!(
2908                    value_type, test_case.expected_type,
2909                    "Test '{}': value type mismatch. Expected: {:?}, Got: {:?}",
2910                    test_case.name, test_case.expected_type, value_type
2911                );
2912            }
2913        }
2914    }
2915
2916    #[tokio::test]
2917    async fn test_get_all_column_names_uses_tag_columns() {
2918        let schema = Arc::new(Schema::new(vec![
2919            ColumnSchema::new(
2920                "greptime_timestamp",
2921                ConcreteDataType::timestamp_millisecond_datatype(),
2922                false,
2923            )
2924            .with_time_index(true),
2925            ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
2926            ColumnSchema::new("region", ConcreteDataType::string_datatype(), false),
2927            ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
2928            ColumnSchema::new(
2929                DATA_SCHEMA_TSID_COLUMN_NAME,
2930                ConcreteDataType::uint64_datatype(),
2931                true,
2932            ),
2933        ]));
2934        let mut options = TableOptions::default();
2935        options.extra_options.insert(
2936            LOGICAL_TABLE_METADATA_KEY.to_string(),
2937            "physical_metrics".to_string(),
2938        );
2939        let meta = TableMetaBuilder::empty()
2940            .schema(schema)
2941            .primary_key_indices(vec![1, 2, 4])
2942            .engine("metric".to_string())
2943            .next_column_id(5)
2944            .options(options)
2945            .build()
2946            .unwrap();
2947        let table_info = TableInfoBuilder::default()
2948            .table_id(1024)
2949            .table_version(0 as TableVersion)
2950            .name("cpu_usage")
2951            .catalog_name(DEFAULT_CATALOG_NAME)
2952            .schema_name(DEFAULT_SCHEMA_NAME)
2953            .table_type(TableType::Base)
2954            .meta(meta)
2955            .build()
2956            .unwrap();
2957        let manager =
2958            MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&table_info));
2959
2960        let physical_schema = Arc::new(Schema::new(vec![
2961            ColumnSchema::new(
2962                "greptime_timestamp",
2963                ConcreteDataType::timestamp_millisecond_datatype(),
2964                false,
2965            )
2966            .with_time_index(true),
2967            ColumnSchema::new(
2968                "physical_only_tag",
2969                ConcreteDataType::string_datatype(),
2970                false,
2971            ),
2972            ColumnSchema::new(
2973                "physical_only_value",
2974                ConcreteDataType::float64_datatype(),
2975                true,
2976            ),
2977        ]));
2978        let mut physical_options = TableOptions::default();
2979        physical_options.extra_options.insert(
2980            store_api::metric_engine_consts::PHYSICAL_TABLE_METADATA_KEY.to_string(),
2981            String::new(),
2982        );
2983        let physical_meta = TableMetaBuilder::empty()
2984            .schema(physical_schema)
2985            .primary_key_indices(vec![1])
2986            .engine("metric".to_string())
2987            .next_column_id(3)
2988            .options(physical_options)
2989            .build()
2990            .unwrap();
2991        let physical_table_info = TableInfoBuilder::default()
2992            .table_id(1025)
2993            .table_version(0 as TableVersion)
2994            .name("physical_metrics")
2995            .catalog_name(DEFAULT_CATALOG_NAME)
2996            .schema_name(DEFAULT_SCHEMA_NAME)
2997            .table_type(TableType::Base)
2998            .meta(physical_meta)
2999            .build()
3000            .unwrap();
3001        let physical_table = EmptyTable::from_table_info(&physical_table_info);
3002        manager
3003            .register_table_sync(RegisterTableRequest {
3004                catalog: DEFAULT_CATALOG_NAME.to_string(),
3005                schema: DEFAULT_SCHEMA_NAME.to_string(),
3006                table_name: "physical_metrics".to_string(),
3007                table_id: 1025,
3008                table: physical_table,
3009            })
3010            .unwrap();
3011        let manager: CatalogManagerRef = manager;
3012
3013        let (labels, table_names) =
3014            get_all_column_names(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, &manager)
3015                .await
3016                .unwrap();
3017
3018        assert_eq!(
3019            labels,
3020            HashSet::from(["host".to_string(), "region".to_string()])
3021        );
3022        assert_eq!(table_names, vec!["cpu_usage"]);
3023
3024        let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3025        let (schemas, targets) =
3026            retrieve_schema_names(&query_ctx, manager.clone(), vec!["cpu_usage".to_string()])
3027                .await
3028                .unwrap();
3029        assert_eq!(schemas, vec![DEFAULT_SCHEMA_NAME]);
3030        assert_eq!(
3031            targets,
3032            PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3033                DEFAULT_CATALOG_NAME,
3034                DEFAULT_SCHEMA_NAME,
3035                "cpu_usage",
3036            )])
3037        );
3038
3039        let (_, targets) = retrieve_schema_names(
3040            &query_ctx,
3041            manager.clone(),
3042            vec![r#"{__name__=~"cpu.*"}"#.to_string()],
3043        )
3044        .await
3045        .unwrap();
3046        assert_eq!(targets, PermissionTableTargets::Unresolved);
3047
3048        let (schemas, targets) = retrieve_schema_names(&query_ctx, manager.clone(), vec![])
3049            .await
3050            .unwrap();
3051        assert!(schemas.contains(&DEFAULT_SCHEMA_NAME.to_string()));
3052        assert_eq!(targets, PermissionTableTargets::Unresolved);
3053
3054        let (schemas, targets) = retrieve_schema_names(
3055            &query_ctx,
3056            manager.clone(),
3057            vec!["physical_metrics".to_string()],
3058        )
3059        .await
3060        .unwrap();
3061        assert!(schemas.is_empty());
3062        assert_eq!(
3063            targets,
3064            PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3065                DEFAULT_CATALOG_NAME,
3066                DEFAULT_SCHEMA_NAME,
3067                "physical_metrics",
3068            )])
3069        );
3070
3071        let (schemas, targets) = retrieve_schema_names(
3072            &query_ctx,
3073            manager.clone(),
3074            vec!["missing".to_string(), "physical_metrics".to_string()],
3075        )
3076        .await
3077        .unwrap();
3078        assert!(schemas.is_empty());
3079        assert_eq!(
3080            targets,
3081            PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3082                DEFAULT_CATALOG_NAME,
3083                DEFAULT_SCHEMA_NAME,
3084                "physical_metrics",
3085            )])
3086        );
3087
3088        assert_eq!(
3089            vec!["cpu_usage".to_string()],
3090            retrieve_table_names(
3091                &query_ctx,
3092                manager.clone(),
3093                vec!["missing".to_string(), r#"{__name__=~"cpu.*"}"#.to_string(),],
3094            )
3095            .await
3096            .unwrap()
3097        );
3098        assert!(
3099            retrieve_table_names(
3100                &query_ctx,
3101                manager.clone(),
3102                vec!["cpu_usage".to_string(), "{".to_string()],
3103            )
3104            .await
3105            .is_err()
3106        );
3107
3108        let (fields, table_names) = retrieve_field_names(
3109            &query_ctx,
3110            manager.clone(),
3111            vec!["physical_metrics".to_string()],
3112        )
3113        .await
3114        .unwrap();
3115        assert!(fields.is_empty());
3116        assert!(table_names.is_empty());
3117
3118        let (fields, table_names) = retrieve_field_names(&query_ctx, manager, vec![])
3119            .await
3120            .unwrap();
3121        assert_eq!(fields, HashSet::from(["value".to_string()]));
3122        assert_eq!(table_names, vec!["cpu_usage"]);
3123    }
3124
3125    #[tokio::test]
3126    async fn test_get_target_column_names_uses_selected_schema() {
3127        let manager = MemoryCatalogManager::with_default_setup();
3128        manager
3129            .register_schema_sync(RegisterSchemaRequest {
3130                catalog: DEFAULT_CATALOG_NAME.to_string(),
3131                schema: "private".to_string(),
3132            })
3133            .unwrap();
3134
3135        for (schema_name, tag_name, table_id) in [
3136            (DEFAULT_SCHEMA_NAME, "public_tag", 2048),
3137            ("private", "private_tag", 2049),
3138        ] {
3139            let schema = Arc::new(Schema::new(vec![
3140                ColumnSchema::new(
3141                    "greptime_timestamp",
3142                    ConcreteDataType::timestamp_millisecond_datatype(),
3143                    false,
3144                )
3145                .with_time_index(true),
3146                ColumnSchema::new(tag_name, ConcreteDataType::string_datatype(), false),
3147                ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
3148            ]));
3149            let meta = TableMetaBuilder::empty()
3150                .schema(schema)
3151                .primary_key_indices(vec![1])
3152                .engine("metric".to_string())
3153                .next_column_id(3)
3154                .build()
3155                .unwrap();
3156            let table_info = TableInfoBuilder::default()
3157                .table_id(table_id)
3158                .table_version(0 as TableVersion)
3159                .name("cpu")
3160                .catalog_name(DEFAULT_CATALOG_NAME)
3161                .schema_name(schema_name)
3162                .table_type(TableType::Base)
3163                .meta(meta)
3164                .build()
3165                .unwrap();
3166            manager
3167                .register_table_sync(RegisterTableRequest {
3168                    catalog: DEFAULT_CATALOG_NAME.to_string(),
3169                    schema: schema_name.to_string(),
3170                    table_name: "cpu".to_string(),
3171                    table_id,
3172                    table: EmptyTable::from_table_info(&table_info),
3173                })
3174                .unwrap();
3175        }
3176
3177        let query_ctx = Arc::new(QueryContext::with(
3178            DEFAULT_CATALOG_NAME,
3179            DEFAULT_SCHEMA_NAME,
3180        ));
3181        let queries = [ParsedPromQuery::parse(
3182            PromQuery {
3183                query: r#"{__name__="cpu",__schema__="private"}"#.to_string(),
3184                ..Default::default()
3185            },
3186            &query_ctx,
3187        )
3188        .unwrap()];
3189        let targets = resolved_promql_targets(&queries, &query_ctx).unwrap();
3190        let manager: CatalogManagerRef = manager;
3191
3192        assert_eq!(
3193            get_target_column_names(&targets, &manager, &query_ctx)
3194                .await
3195                .unwrap(),
3196            HashSet::from(["private_tag".to_string()])
3197        );
3198    }
3199
3200    #[test]
3201    fn test_metric_name_discovery_uses_selected_schema() {
3202        let expr =
3203            promql_parser::parser::parse(r#"{__name__=~"cpu.*",__database__="private"}"#).unwrap();
3204        let discovery = find_metric_name_not_equal_matchers(&expr).unwrap().unwrap();
3205
3206        assert_eq!(discovery.schema.as_deref(), Some("private"));
3207        assert_eq!(discovery.name_matchers.len(), 1);
3208        assert_eq!(discovery.name_matchers[0].value, "cpu.*");
3209        assert!(matches!(discovery.name_matchers[0].op, MatchOp::Re(_)));
3210
3211        let expr = promql_parser::parser::parse(
3212            r#"{__name__=~"cpu.*",__database__="private",__schema__="public"}"#,
3213        )
3214        .unwrap();
3215        let discovery = find_metric_name_not_equal_matchers(&expr).unwrap().unwrap();
3216        assert_eq!(discovery.schema.as_deref(), Some("public"));
3217    }
3218
3219    #[test]
3220    fn test_metric_name_discovery_rejects_unsafe_selector_shapes() {
3221        for promql in [
3222            r#"{__name__="denied",__name__=~"no_such_.*"}"#,
3223            r#"{__name__=~"cpu.*" or job="api"}"#,
3224        ] {
3225            let expr = promql_parser::parser::parse(promql).unwrap();
3226            assert!(
3227                find_metric_name_not_equal_matchers(&expr)
3228                    .unwrap()
3229                    .is_none(),
3230                "{promql}"
3231            );
3232        }
3233
3234        let expr = promql_parser::parser::parse(r#"{__name__=~"cpu.*" or job="api"}"#).unwrap();
3235        assert_eq!(
3236            (PermissionTableTargets::Resolved(vec![]), 1),
3237            static_promql_targets(&expr, &QueryContext::arc()).unwrap()
3238        );
3239
3240        let expr =
3241            promql_parser::parser::parse(r#"allowed{job="api" or instance="host"}"#).unwrap();
3242        assert_eq!(
3243            (
3244                PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3245                    "greptime", "public", "allowed",
3246                )]),
3247                0,
3248            ),
3249            static_promql_targets(&expr, &QueryContext::arc()).unwrap()
3250        );
3251    }
3252
3253    #[test]
3254    fn test_metric_name_discovery_ignores_or_matchers_on_named_selector() {
3255        let expr = promql_parser::parser::parse(
3256            r#"{__name__=~"cpu.*"} + allowed{job="api" or instance="host"}"#,
3257        )
3258        .unwrap();
3259
3260        let discovery = find_metric_name_not_equal_matchers(&expr).unwrap().unwrap();
3261        assert_eq!(discovery.name_matchers[0].value, "cpu.*");
3262    }
3263
3264    #[test]
3265    fn test_retrieve_exact_metric_name_from_selector() {
3266        assert_eq!(
3267            retrieve_exact_metric_name_from_promql(r#"cpu{host="a"}"#).as_deref(),
3268            Some("cpu")
3269        );
3270        assert_eq!(
3271            retrieve_exact_metric_name_from_promql(r#"{__name__="cpu",host="a"}"#).as_deref(),
3272            Some("cpu")
3273        );
3274        assert!(retrieve_exact_metric_name_from_promql(r#"{__name__=~"cpu.*"}"#).is_none());
3275    }
3276
3277    #[test]
3278    fn test_expand_metric_name_queries_preserves_schema_matcher() {
3279        let query = PromQuery {
3280            query: r#"{__name__=~"cpu.*",__database__="private"}"#.to_string(),
3281            ..Default::default()
3282        };
3283        let ctx = QueryContext::arc();
3284        let query = ParsedPromQuery::parse(query, &ctx).unwrap();
3285        let queries = expand_metric_name_queries(
3286            &query,
3287            vec!["cpu_user".to_string(), "cpu_system".to_string()],
3288        );
3289
3290        assert_eq!(queries.len(), 2);
3291        for (query, metric_name) in queries.iter().zip(["cpu_user", "cpu_system"]) {
3292            let PromqlExpr::VectorSelector(selector) = query.expr() else {
3293                panic!("expected vector selector");
3294            };
3295            assert!(selector.matchers.matchers.iter().any(|matcher| {
3296                matcher.name == METRIC_NAME
3297                    && matcher.op == MatchOp::Equal
3298                    && matcher.value == metric_name
3299            }));
3300            assert!(selector.matchers.matchers.iter().any(|matcher| {
3301                matcher.name == "__database__"
3302                    && matcher.op == MatchOp::Equal
3303                    && matcher.value == "private"
3304            }));
3305        }
3306    }
3307
3308    #[test]
3309    fn test_expand_metric_name_queries_preserves_exact_selector() {
3310        let query = PromQuery {
3311            query: r#"denied + {__name__=~"cpu.*"}"#.to_string(),
3312            ..Default::default()
3313        };
3314        let ctx = QueryContext::arc();
3315        let query = ParsedPromQuery::parse(query, &ctx).unwrap();
3316        let queries = expand_metric_name_queries(&query, vec!["cpu_user".to_string()]);
3317
3318        assert_eq!(queries.len(), 1);
3319        let expanded = queries[0].expr();
3320        let ctx = Arc::new(QueryContext::with("greptime", "public"));
3321        let (PermissionTableTargets::Resolved(mut targets), unresolved_selectors) =
3322            static_promql_targets(expanded, &ctx).unwrap()
3323        else {
3324            panic!("expected resolved targets");
3325        };
3326        assert_eq!(unresolved_selectors, 0);
3327        targets.sort_unstable_by(|left, right| left.table.cmp(&right.table));
3328        assert_eq!(
3329            targets,
3330            vec![
3331                PermissionTableTarget::new("greptime", "public", "cpu_user"),
3332                PermissionTableTarget::new("greptime", "public", "denied"),
3333            ]
3334        );
3335    }
3336
3337    #[test]
3338    fn test_static_promql_targets_keep_exact_selector_during_discovery() {
3339        let expr = promql_parser::parser::parse(
3340            r#"denied + {__name__=~"no_match_.*",__schema__="private"}"#,
3341        )
3342        .unwrap();
3343        let ctx = Arc::new(QueryContext::with("greptime", "public"));
3344
3345        assert_eq!(
3346            static_promql_targets(&expr, &ctx).unwrap(),
3347            (
3348                PermissionTableTargets::resolved(vec![PermissionTableTarget::new(
3349                    "greptime", "public", "denied",
3350                )]),
3351                1,
3352            )
3353        );
3354    }
3355
3356    #[test]
3357    fn test_static_promql_targets_count_unresolved_selectors() {
3358        let expr =
3359            promql_parser::parser::parse(r#"{__name__=~"no_such_metric"} or {job="api"}"#).unwrap();
3360        let ctx = Arc::new(QueryContext::with("greptime", "public"));
3361
3362        assert_eq!(
3363            static_promql_targets(&expr, &ctx).unwrap(),
3364            (PermissionTableTargets::resolved(Vec::new()), 2)
3365        );
3366    }
3367
3368    #[test]
3369    fn test_record_batches_to_labels_name_accepts_native_histogram_value() {
3370        let schema = Arc::new(Schema::new(vec![
3371            ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
3372            ColumnSchema::new(
3373                greptime_native_histogram(),
3374                native_histogram_value_type().clone(),
3375                true,
3376            ),
3377        ]));
3378        let batch = RecordBatch::new_empty(schema.clone());
3379        let batches = RecordBatches::try_new(schema.clone(), vec![batch]).unwrap();
3380        let mut labels = HashSet::new();
3381
3382        record_batches_to_labels_name(batches, &mut labels).unwrap();
3383
3384        assert!(labels.is_empty());
3385
3386        let histogram = NativeHistogram {
3387            schema: 0,
3388            zero_threshold: 0.001,
3389            sum: 1.0,
3390            reset_hint: CounterResetHint::Unknown,
3391            start_timestamp: None,
3392            custom_values: vec![],
3393            positive_spans: vec![],
3394            negative_spans: vec![],
3395            count: 1.0,
3396            zero_count: 1.0,
3397            positive_buckets: vec![],
3398            negative_buckets: vec![],
3399        };
3400        let histogram_array = build_histogram_array(&[Some(histogram)]);
3401        let histogram_array = histogram_array
3402            .as_any()
3403            .downcast_ref::<arrow::array::StructArray>()
3404            .unwrap()
3405            .clone();
3406        let ConcreteDataType::Struct(histogram_type) = native_histogram_value_type().clone() else {
3407            unreachable!("native histogram type must be a struct")
3408        };
3409        let batch = RecordBatch::new(
3410            schema.clone(),
3411            vec![
3412                Arc::new(StringVector::from(vec![Some("localhost")])) as _,
3413                Arc::new(StructVector::try_new(histogram_type, histogram_array).unwrap()) as _,
3414            ],
3415        )
3416        .unwrap();
3417        let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
3418
3419        record_batches_to_labels_name(batches, &mut labels).unwrap();
3420
3421        assert_eq!(
3422            labels,
3423            HashSet::from(["host".to_string(), greptime_native_histogram().to_string(),])
3424        );
3425    }
3426
3427    fn prometheus_metric_table_info(table_id: u32, name: &str) -> TableInfo {
3428        let schema = Arc::new(Schema::new(vec![
3429            ColumnSchema::new(
3430                "greptime_timestamp",
3431                ConcreteDataType::timestamp_millisecond_datatype(),
3432                false,
3433            )
3434            .with_time_index(true),
3435            ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
3436            ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
3437        ]));
3438        let mut options = TableOptions::default();
3439        options.extra_options.insert(
3440            LOGICAL_TABLE_METADATA_KEY.to_string(),
3441            "greptime_physical_table".to_string(),
3442        );
3443        options
3444            .extra_options
3445            .insert(SEMANTIC_METRIC_TYPE.to_string(), "counter".to_string());
3446        options
3447            .extra_options
3448            .insert(SEMANTIC_METRIC_UNIT.to_string(), "By".to_string());
3449        let meta = TableMetaBuilder::empty()
3450            .schema(schema)
3451            .primary_key_indices(vec![1])
3452            .engine("metric".to_string())
3453            .next_column_id(3)
3454            .options(options)
3455            .build()
3456            .unwrap();
3457
3458        TableInfoBuilder::default()
3459            .table_id(table_id)
3460            .table_version(0 as TableVersion)
3461            .name(name)
3462            .catalog_name(DEFAULT_CATALOG_NAME)
3463            .schema_name(DEFAULT_SCHEMA_NAME)
3464            .table_type(TableType::Base)
3465            .meta(meta)
3466            .build()
3467            .unwrap()
3468    }
3469
3470    #[tokio::test]
3471    async fn test_retrieve_metric_metadata_uses_semantic_options() {
3472        let table_info = prometheus_metric_table_info(1025, "http_requests_total");
3473        let manager: CatalogManagerRef =
3474            MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&table_info));
3475        let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3476
3477        let metadata = retrieve_metric_metadata(&query_ctx, manager, &MetadataQuery::default())
3478            .await
3479            .unwrap();
3480
3481        assert_eq!(
3482            metadata.get("http_requests_total"),
3483            Some(&vec![PromMetadata {
3484                metric_type: "counter".to_string(),
3485                unit: "bytes".to_string(),
3486                help: String::new(),
3487            }])
3488        );
3489
3490        for (ucum, openmetrics) in [
3491            ("1", "ratios"),
3492            ("ms", "milliseconds"),
3493            ("m/s", "meters_per_second"),
3494            ("USD", "USD"),
3495        ] {
3496            let mut table_info = table_info.clone();
3497            table_info
3498                .meta
3499                .options
3500                .extra_options
3501                .insert(SEMANTIC_METRIC_UNIT.to_string(), ucum.to_string());
3502            assert_eq!(
3503                prometheus_metadata_from_table(&table_info).unit,
3504                openmetrics
3505            );
3506        }
3507
3508        for (metric_type, temporality, expected) in [
3509            ("updown_counter", None, "gauge"),
3510            ("gauge_histogram", None, "gaugehistogram"),
3511            ("summary", None, "summary"),
3512            ("mixed", None, "unknown"),
3513            ("counter", Some("delta"), "unknown"),
3514            ("histogram", Some("delta"), "unknown"),
3515            ("counter", Some("mixed"), "unknown"),
3516            ("histogram", Some("mixed"), "unknown"),
3517            ("updown_counter", Some("mixed"), "unknown"),
3518        ] {
3519            let mut table_info = table_info.clone();
3520            table_info
3521                .meta
3522                .options
3523                .extra_options
3524                .insert(SEMANTIC_METRIC_TYPE.to_string(), metric_type.to_string());
3525            if let Some(temporality) = temporality {
3526                table_info.meta.options.extra_options.insert(
3527                    SEMANTIC_METRIC_TEMPORALITY.to_string(),
3528                    temporality.to_string(),
3529                );
3530            }
3531            assert_eq!(
3532                prometheus_metadata_from_table(&table_info).metric_type,
3533                expected
3534            );
3535        }
3536    }
3537
3538    #[tokio::test]
3539    async fn test_metadata_query_checks_permissions_and_filters_tables() {
3540        let allowed = prometheus_metric_table_info(1025, "allowed");
3541        let denied = prometheus_metric_table_info(1026, "denied");
3542        let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&allowed));
3543        manager
3544            .register_table_sync(RegisterTableRequest {
3545                catalog: DEFAULT_CATALOG_NAME.to_string(),
3546                schema: DEFAULT_SCHEMA_NAME.to_string(),
3547                table_name: denied.name.clone(),
3548                table_id: denied.table_id(),
3549                table: EmptyTable::from_table_info(&denied),
3550            })
3551            .unwrap();
3552        let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3553
3554        let metadata = retrieve_metric_metadata(
3555            &query_ctx,
3556            manager.clone(),
3557            &MetadataQuery {
3558                metric: Some("allowed".to_string()),
3559                ..Default::default()
3560            },
3561        )
3562        .await
3563        .unwrap();
3564        assert_eq!(metadata.len(), 1);
3565        assert!(metadata.contains_key("allowed"));
3566
3567        let response = metadata_query(
3568            State(Arc::new(TestPrometheusHandler {
3569                catalog_manager: manager.clone(),
3570                deny_operation: true,
3571                denied_table: None,
3572                metric_names: Vec::new(),
3573                queries: Mutex::new(Vec::new()),
3574            })),
3575            Query(MetadataQuery::default()),
3576            Extension(query_ctx.clone()),
3577        )
3578        .await;
3579        assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
3580
3581        let handler: PrometheusHandlerRef = Arc::new(TestPrometheusHandler {
3582            catalog_manager: manager,
3583            deny_operation: false,
3584            denied_table: Some("denied"),
3585            metric_names: Vec::new(),
3586            queries: Mutex::new(Vec::new()),
3587        });
3588        let response = metadata_query(
3589            State(handler.clone()),
3590            Query(MetadataQuery::default()),
3591            Extension(query_ctx.clone()),
3592        )
3593        .await;
3594        assert!(response.status_code.is_none());
3595        let PrometheusResponse::Metadata(metadata) = response.data else {
3596            panic!("expected metadata response");
3597        };
3598        assert_eq!(
3599            metadata.keys().cloned().collect::<Vec<_>>(),
3600            vec!["allowed".to_string()]
3601        );
3602
3603        let response = metadata_query(
3604            State(handler),
3605            Query(MetadataQuery {
3606                metric: Some("denied".to_string()),
3607                ..Default::default()
3608            }),
3609            Extension(query_ctx),
3610        )
3611        .await;
3612        assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
3613    }
3614}