1use 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, QueryLanguageParser, 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::{PromSampleValue, PrometheusJsonResponse};
75use crate::error::{
76 CollectRecordbatchSnafu, ConvertScalarValueSnafu, DataFusionSnafu, Error, InvalidQuerySnafu,
77 NotSupportedSnafu, ParseTimestampSnafu, 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#[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#[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, PromSampleValue)>,
104 #[serde(skip_serializing_if = "Vec::is_empty", default)]
105 pub histograms: Vec<(f64, PromNativeHistogram)>,
106}
107
108#[derive(Debug, Default, Serialize, Deserialize, PartialEq)]
110pub struct PromNativeHistogram {
111 pub count: String,
113 pub sum: String,
115 pub buckets: Vec<(u8, String, String, String)>,
121}
122
123#[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#[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#[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 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 }
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#[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(¶ms.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) = ¶ms.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(), ¶ms).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
344macro_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 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 if let Some(db) = ¶ms.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
489async 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 if let Some(db) = ¶ms.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 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
619async 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
630 .do_query_parsed(prom_query.with_unordered_output(), query_ctx)
631 .await;
632 PrometheusJsonResponse::from_query_result(
633 result,
634 metric_name,
635 ValueType::Matrix,
636 query_id.as_deref(),
637 )
638 .await
639}
640
641#[derive(Debug, Default, Serialize)]
642struct Matches(Vec<String>);
643
644#[derive(Debug, Default, Serialize, Deserialize)]
645pub struct LabelsQuery {
646 start: Option<String>,
647 end: Option<String>,
648 lookback: Option<String>,
649 #[serde(flatten)]
650 matches: Matches,
651 db: Option<String>,
652}
653
654impl<'de> Deserialize<'de> for Matches {
656 fn deserialize<D>(deserializer: D) -> std::result::Result<Matches, D::Error>
657 where
658 D: de::Deserializer<'de>,
659 {
660 struct MatchesVisitor;
661
662 impl<'d> Visitor<'d> for MatchesVisitor {
663 type Value = Vec<String>;
664
665 fn expecting(&self, formatter: &mut std::fmt::Formatter) -> std::fmt::Result {
666 formatter.write_str("a string")
667 }
668
669 fn visit_map<M>(self, mut access: M) -> std::result::Result<Self::Value, M::Error>
670 where
671 M: MapAccess<'d>,
672 {
673 let mut matches = Vec::new();
674 while let Some((key, value)) = access.next_entry::<String, String>()? {
675 if key == "match[]" {
676 matches.push(value);
677 }
678 }
679 Ok(matches)
680 }
681 }
682 Ok(Matches(deserializer.deserialize_map(MatchesVisitor)?))
683 }
684}
685
686macro_rules! handle_schema_err {
693 ($result:expr) => {
694 match $result {
695 Ok(v) => Some(v),
696 Err(err) => {
697 if err.status_code() == StatusCode::TableNotFound
698 || err.status_code() == StatusCode::TableColumnNotFound
699 {
700 None
702 } else {
703 return PrometheusJsonResponse::error(err.status_code(), err.output_msg());
704 }
705 }
706 }
707 };
708}
709
710#[axum_macros::debug_handler]
711#[tracing::instrument(
712 skip_all,
713 fields(protocol = "prometheus", request_type = "labels_query")
714)]
715pub async fn labels_query(
716 State(handler): State<PrometheusHandlerRef>,
717 Query(params): Query<LabelsQuery>,
718 Extension(mut query_ctx): Extension<QueryContext>,
719 Form(form_params): Form<LabelsQuery>,
720) -> PrometheusJsonResponse {
721 let (catalog, schema) = get_catalog_schema(¶ms.db, &query_ctx);
722 try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
723 let query_ctx = Arc::new(query_ctx);
724
725 let mut queries = params.matches.0;
726 if queries.is_empty() {
727 queries = form_params.matches.0;
728 }
729
730 let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
731 .with_label_values(&[query_ctx.get_db_string().as_str(), "labels_query"])
732 .start_timer();
733
734 let prom_queries = if queries.is_empty() {
735 try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
736 Vec::new()
737 } else {
738 let start = params
739 .start
740 .or(form_params.start)
741 .unwrap_or_else(yesterday_rfc3339);
742 let end = params
743 .end
744 .or(form_params.end)
745 .unwrap_or_else(current_time_rfc3339);
746 let lookback = params
747 .lookback
748 .or(form_params.lookback)
749 .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string());
750 let prom_queries = queries
751 .into_iter()
752 .map(|query| PromQuery {
753 query,
754 start: start.clone(),
755 end: end.clone(),
756 step: DEFAULT_LOOKBACK_STRING.to_string(),
757 lookback: lookback.clone(),
758 alias: None,
759 })
760 .collect::<Vec<_>>();
761 let prom_queries = try_call_return_response!(
762 prom_queries
763 .into_iter()
764 .map(|query| ParsedPromQuery::parse(query, &query_ctx))
765 .collect::<Result<Vec<_>>>()
766 );
767 try_call_return_response!(
768 handler
769 .check_query_permission_parsed(&prom_queries, &query_ctx)
770 .await
771 );
772 prom_queries
773 };
774
775 if prom_queries.is_empty() {
777 let (mut labels, table_names) =
778 match get_all_column_names(&catalog, &schema, &handler.catalog_manager()).await {
779 Ok(result) => result,
780 Err(e) => return PrometheusJsonResponse::error(e.status_code(), e.output_msg()),
781 };
782 try_call_return_response!(
783 handler
784 .check_query_target_permission(
785 current_schema_metric_targets(&query_ctx, &table_names),
786 &query_ctx,
787 )
788 .await
789 );
790 let _ = labels.insert(METRIC_NAME.to_string());
791 let mut labels_vec = labels.into_iter().collect::<Vec<_>>();
792 labels_vec.sort_unstable();
793 return PrometheusJsonResponse::success(PrometheusResponse::Labels(labels_vec));
794 }
795
796 let mut labels = match resolved_promql_targets(&prom_queries, &query_ctx) {
798 Some(targets) => {
799 match get_target_column_names(&targets, &handler.catalog_manager(), &query_ctx).await {
800 Ok(labels) => labels,
801 Err(e) => return PrometheusJsonResponse::error(e.status_code(), e.output_msg()),
802 }
803 }
804 None => match get_all_column_names(&catalog, &schema, &handler.catalog_manager()).await {
805 Ok((labels, _)) => labels,
806 Err(e) => return PrometheusJsonResponse::error(e.status_code(), e.output_msg()),
807 },
808 };
809 let _ = labels.insert(METRIC_NAME.to_string());
810
811 let mut fetched_labels = HashSet::new();
812 let _ = fetched_labels.insert(METRIC_NAME.to_string());
813
814 let mut merge_map = HashMap::new();
815 for prom_query in prom_queries {
816 let result = handler.do_query_parsed(prom_query, query_ctx.clone()).await;
817 handle_schema_err!(
818 retrieve_labels_name_from_query_result(result, &mut fetched_labels, &mut merge_map)
819 .await
820 );
821 }
822
823 fetched_labels.retain(|l| labels.contains(l));
825 let mut sorted_labels: Vec<String> = fetched_labels.into_iter().collect();
826 sorted_labels.sort();
827 let merge_map = merge_map
828 .into_iter()
829 .map(|(k, v)| (k, Value::from(v)))
830 .collect();
831 let mut resp = PrometheusJsonResponse::success(PrometheusResponse::Labels(sorted_labels));
832 resp.resp_metrics = merge_map;
833 resp
834}
835
836async fn get_all_column_names(
838 catalog: &str,
839 schema: &str,
840 manager: &CatalogManagerRef,
841) -> std::result::Result<(HashSet<String>, Vec<String>), catalog::error::Error> {
842 let mut labels = HashSet::new();
843 let mut table_names = Vec::new();
844 let mut tables = manager.tables(catalog, schema, None);
845 while let Some(table) = tables.try_next().await? {
846 if is_internal_physical_metric_table(&table) {
847 continue;
848 }
849 table_names.push(table.table_info().name.clone());
850 extend_tag_column_names(&mut labels, &table);
851 }
852
853 Ok((labels, table_names))
854}
855
856async fn get_target_column_names(
857 targets: &[PermissionTableTarget],
858 manager: &CatalogManagerRef,
859 query_ctx: &QueryContext,
860) -> std::result::Result<HashSet<String>, catalog::error::Error> {
861 let mut labels = HashSet::new();
862 for target in targets {
863 if let Some(table) = manager
864 .table(
865 &target.catalog,
866 &target.schema,
867 &target.table,
868 Some(query_ctx),
869 )
870 .await?
871 && !is_internal_physical_metric_table(&table)
872 {
873 extend_tag_column_names(&mut labels, &table);
874 }
875 }
876 Ok(labels)
877}
878
879fn extend_tag_column_names(labels: &mut HashSet<String>, table: &TableRef) {
880 labels.extend(table.primary_key_columns().filter_map(|column| {
881 (column.name != DATA_SCHEMA_TABLE_ID_COLUMN_NAME
882 && column.name != DATA_SCHEMA_TSID_COLUMN_NAME)
883 .then_some(column.name)
884 }));
885}
886
887async fn retrieve_series_from_query_result(
888 result: Result<Output>,
889 series: &mut Vec<HashMap<Column, String>>,
890 query_ctx: &QueryContext,
891 target: &PermissionTableTarget,
892 manager: &CatalogManagerRef,
893 metrics: &mut HashMap<String, u64>,
894) -> Result<()> {
895 let result = result?;
896
897 let table = manager
899 .table(
900 &target.catalog,
901 &target.schema,
902 &target.table,
903 Some(query_ctx),
904 )
905 .await?
906 .with_context(|| TableNotFoundSnafu {
907 catalog: &target.catalog,
908 schema: &target.schema,
909 table: &target.table,
910 })?;
911 let tag_columns = table
912 .primary_key_columns()
913 .map(|c| c.name)
914 .collect::<HashSet<_>>();
915
916 match result.data {
917 OutputData::RecordBatches(batches) => {
918 record_batches_to_series(batches, series, &target.table, &tag_columns)
919 }
920 OutputData::Stream(stream) => {
921 let batches = RecordBatches::try_collect(stream)
922 .await
923 .context(CollectRecordbatchSnafu)?;
924 record_batches_to_series(batches, series, &target.table, &tag_columns)
925 }
926 OutputData::AffectedRows(_) => Err(Error::UnexpectedResult {
927 reason: "expected data result, but got affected rows".to_string(),
928 location: Location::default(),
929 }),
930 }?;
931
932 if let Some(ref plan) = result.meta.plan {
933 collect_plan_metrics(plan, &mut [metrics]);
934 }
935 Ok(())
936}
937
938async fn retrieve_labels_name_from_query_result(
940 result: Result<Output>,
941 labels: &mut HashSet<String>,
942 metrics: &mut HashMap<String, u64>,
943) -> Result<()> {
944 let result = result?;
945 match result.data {
946 OutputData::RecordBatches(batches) => record_batches_to_labels_name(batches, labels),
947 OutputData::Stream(stream) => {
948 let batches = RecordBatches::try_collect(stream)
949 .await
950 .context(CollectRecordbatchSnafu)?;
951 record_batches_to_labels_name(batches, labels)
952 }
953 OutputData::AffectedRows(_) => UnexpectedResultSnafu {
954 reason: "expected data result, but got affected rows".to_string(),
955 }
956 .fail(),
957 }?;
958 if let Some(ref plan) = result.meta.plan {
959 collect_plan_metrics(plan, &mut [metrics]);
960 }
961 Ok(())
962}
963
964fn record_batches_to_series(
965 batches: RecordBatches,
966 series: &mut Vec<HashMap<Column, String>>,
967 table_name: &str,
968 tag_columns: &HashSet<String>,
969) -> Result<()> {
970 for batch in batches.iter() {
971 let projection = batch
973 .schema
974 .column_schemas()
975 .iter()
976 .enumerate()
977 .filter_map(|(idx, col)| {
978 if tag_columns.contains(&col.name) {
979 Some(idx)
980 } else {
981 None
982 }
983 })
984 .collect::<Vec<_>>();
985 let batch = batch
986 .try_project(&projection)
987 .context(CollectRecordbatchSnafu)?;
988
989 let mut writer = RowWriter::new(&batch.schema, table_name);
990 writer.write(batch, series)?;
991 }
992 Ok(())
993}
994
995struct RowWriter {
1002 template: HashMap<Column, Option<String>>,
1005 current: Option<HashMap<Column, Option<String>>>,
1007}
1008
1009impl RowWriter {
1010 fn new(schema: &SchemaRef, table: &str) -> Self {
1011 let mut template = schema
1012 .column_schemas()
1013 .iter()
1014 .map(|x| (x.name.as_str().into(), None))
1015 .collect::<HashMap<Column, Option<String>>>();
1016 template.insert("__name__".into(), Some(table.to_string()));
1017 Self {
1018 template,
1019 current: None,
1020 }
1021 }
1022
1023 fn insert(&mut self, column: ColumnRef, value: impl ToString) {
1024 let current = self.current.get_or_insert_with(|| self.template.clone());
1025 match current.get_mut(&column as &dyn AsColumnRef) {
1026 Some(x) => {
1027 let _ = x.insert(value.to_string());
1028 }
1029 None => {
1030 let _ = current.insert(column.0.into(), Some(value.to_string()));
1031 }
1032 }
1033 }
1034
1035 fn insert_bytes(&mut self, column_schema: &ColumnSchema, bytes: &[u8]) -> Result<()> {
1036 let column_name = column_schema.name.as_str().into();
1037
1038 if column_schema.data_type.is_json() {
1039 let s = jsonb_to_string(bytes).context(ConvertScalarValueSnafu)?;
1040 self.insert(column_name, s);
1041 } else {
1042 let hex = bytes
1043 .iter()
1044 .map(|b| format!("{b:02x}"))
1045 .collect::<Vec<String>>()
1046 .join("");
1047 self.insert(column_name, hex);
1048 }
1049 Ok(())
1050 }
1051
1052 fn finish(&mut self) -> HashMap<Column, String> {
1053 let Some(current) = self.current.take() else {
1054 return HashMap::new();
1055 };
1056 current
1057 .into_iter()
1058 .filter_map(|(k, v)| v.map(|v| (k, v)))
1059 .collect()
1060 }
1061
1062 fn write(
1063 &mut self,
1064 record_batch: RecordBatch,
1065 series: &mut Vec<HashMap<Column, String>>,
1066 ) -> Result<()> {
1067 let schema = record_batch.schema.clone();
1068 let record_batch = record_batch.into_df_record_batch();
1069 for i in 0..record_batch.num_rows() {
1070 for (j, array) in record_batch.columns().iter().enumerate() {
1071 let column = schema.column_name_by_index(j).into();
1072
1073 if array.is_null(i) {
1074 self.insert(column, "Null");
1075 continue;
1076 }
1077
1078 match array.data_type() {
1079 DataType::Null => {
1080 self.insert(column, "Null");
1081 }
1082 DataType::Boolean => {
1083 let array = array.as_boolean();
1084 let v = array.value(i);
1085 self.insert(column, v);
1086 }
1087 DataType::UInt8 => {
1088 let array = array.as_primitive::<UInt8Type>();
1089 let v = array.value(i);
1090 self.insert(column, v);
1091 }
1092 DataType::UInt16 => {
1093 let array = array.as_primitive::<UInt16Type>();
1094 let v = array.value(i);
1095 self.insert(column, v);
1096 }
1097 DataType::UInt32 => {
1098 let array = array.as_primitive::<UInt32Type>();
1099 let v = array.value(i);
1100 self.insert(column, v);
1101 }
1102 DataType::UInt64 => {
1103 let array = array.as_primitive::<UInt64Type>();
1104 let v = array.value(i);
1105 self.insert(column, v);
1106 }
1107 DataType::Int8 => {
1108 let array = array.as_primitive::<Int8Type>();
1109 let v = array.value(i);
1110 self.insert(column, v);
1111 }
1112 DataType::Int16 => {
1113 let array = array.as_primitive::<Int16Type>();
1114 let v = array.value(i);
1115 self.insert(column, v);
1116 }
1117 DataType::Int32 => {
1118 let array = array.as_primitive::<Int32Type>();
1119 let v = array.value(i);
1120 self.insert(column, v);
1121 }
1122 DataType::Int64 => {
1123 let array = array.as_primitive::<Int64Type>();
1124 let v = array.value(i);
1125 self.insert(column, v);
1126 }
1127 DataType::Float32 => {
1128 let array = array.as_primitive::<Float32Type>();
1129 let v = array.value(i);
1130 self.insert(column, v);
1131 }
1132 DataType::Float64 => {
1133 let array = array.as_primitive::<Float64Type>();
1134 let v = array.value(i);
1135 self.insert(column, v);
1136 }
1137 DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View => {
1138 let v = datatypes::arrow_array::string_array_value(array, i);
1139 self.insert(column, v);
1140 }
1141 DataType::Binary | DataType::LargeBinary | DataType::BinaryView => {
1142 let v = datatypes::arrow_array::binary_array_value(array, i);
1143 let column_schema = &schema.column_schemas()[j];
1144 self.insert_bytes(column_schema, v)?;
1145 }
1146 DataType::Date32 => {
1147 let array = array.as_primitive::<Date32Type>();
1148 let v = Date::new(array.value(i));
1149 self.insert(column, v);
1150 }
1151 DataType::Date64 => {
1152 let array = array.as_primitive::<Date64Type>();
1153 let v = Date::new((array.value(i) / 86_400_000) as i32);
1157 self.insert(column, v);
1158 }
1159 DataType::Timestamp(_, _) => {
1160 let v = datatypes::arrow_array::timestamp_array_value(array, i);
1161 self.insert(column, v.to_iso8601_string());
1162 }
1163 DataType::Time32(_) | DataType::Time64(_) => {
1164 let v = datatypes::arrow_array::time_array_value(array, i);
1165 self.insert(column, v.to_iso8601_string());
1166 }
1167 DataType::Interval(interval_unit) => match interval_unit {
1168 IntervalUnit::YearMonth => {
1169 let array = array.as_primitive::<IntervalYearMonthType>();
1170 let v: IntervalYearMonth = array.value(i).into();
1171 self.insert(column, v.to_iso8601_string());
1172 }
1173 IntervalUnit::DayTime => {
1174 let array = array.as_primitive::<IntervalDayTimeType>();
1175 let v: IntervalDayTime = array.value(i).into();
1176 self.insert(column, v.to_iso8601_string());
1177 }
1178 IntervalUnit::MonthDayNano => {
1179 let array = array.as_primitive::<IntervalMonthDayNanoType>();
1180 let v: IntervalMonthDayNano = array.value(i).into();
1181 self.insert(column, v.to_iso8601_string());
1182 }
1183 },
1184 DataType::Duration(_) => {
1185 let d = datatypes::arrow_array::duration_array_value(array, i);
1186 self.insert(column, d);
1187 }
1188 DataType::List(_) => {
1189 let v = ScalarValue::try_from_array(array, i).context(DataFusionSnafu)?;
1190 self.insert(column, v);
1191 }
1192 DataType::Struct(_) => {
1193 let v = ScalarValue::try_from_array(array, i).context(DataFusionSnafu)?;
1194 self.insert(column, v);
1195 }
1196 DataType::Decimal128(precision, scale) => {
1197 let array = array.as_primitive::<Decimal128Type>();
1198 let v = Decimal128::new(array.value(i), *precision, *scale);
1199 self.insert(column, v);
1200 }
1201 _ => {
1202 return NotSupportedSnafu {
1203 feat: format!("convert {} to http value", array.data_type()),
1204 }
1205 .fail();
1206 }
1207 }
1208 }
1209
1210 series.push(self.finish())
1211 }
1212 Ok(())
1213 }
1214}
1215
1216#[derive(Debug, Clone, Copy, PartialEq)]
1217struct ColumnRef<'a>(&'a str);
1218
1219impl<'a> From<&'a str> for ColumnRef<'a> {
1220 fn from(s: &'a str) -> Self {
1221 Self(s)
1222 }
1223}
1224
1225trait AsColumnRef {
1226 fn as_ref(&self) -> ColumnRef<'_>;
1227}
1228
1229impl AsColumnRef for Column {
1230 fn as_ref(&self) -> ColumnRef<'_> {
1231 self.0.as_str().into()
1232 }
1233}
1234
1235impl AsColumnRef for ColumnRef<'_> {
1236 fn as_ref(&self) -> ColumnRef<'_> {
1237 *self
1238 }
1239}
1240
1241impl<'a> PartialEq for dyn AsColumnRef + 'a {
1242 fn eq(&self, other: &Self) -> bool {
1243 self.as_ref() == other.as_ref()
1244 }
1245}
1246
1247impl<'a> Eq for dyn AsColumnRef + 'a {}
1248
1249impl<'a> Hash for dyn AsColumnRef + 'a {
1250 fn hash<H: Hasher>(&self, state: &mut H) {
1251 self.as_ref().0.hash(state);
1252 }
1253}
1254
1255impl<'a> Borrow<dyn AsColumnRef + 'a> for Column {
1256 fn borrow(&self) -> &(dyn AsColumnRef + 'a) {
1257 self
1258 }
1259}
1260
1261fn record_batches_to_labels_name(
1263 batches: RecordBatches,
1264 labels: &mut HashSet<String>,
1265) -> Result<()> {
1266 let mut column_indices = Vec::new();
1267 let mut value_column_indices = Vec::new();
1268 for (i, column) in batches.schema().column_schemas().iter().enumerate() {
1269 if is_prometheus_value_column(&column.data_type) {
1270 value_column_indices.push(i);
1271 }
1272 column_indices.push(i);
1273 }
1274
1275 if value_column_indices.is_empty() {
1276 return Err(Error::Internal {
1277 err_msg: "no value column found".to_string(),
1278 });
1279 }
1280
1281 for batch in batches.iter() {
1282 let names = column_indices
1283 .iter()
1284 .map(|c| batches.schema().column_name_by_index(*c).to_string())
1285 .collect::<Vec<_>>();
1286
1287 let value_columns = value_column_indices
1288 .iter()
1289 .map(|i| batch.column(*i))
1290 .collect::<Vec<_>>();
1291
1292 for row_index in 0..batch.num_rows() {
1293 if value_columns.iter().all(|c| c.is_null(row_index)) {
1295 continue;
1296 }
1297
1298 names.iter().for_each(|name| {
1300 let _ = labels.insert(name.clone());
1301 });
1302 return Ok(());
1303 }
1304 }
1305 Ok(())
1306}
1307
1308fn is_prometheus_value_column(data_type: &ConcreteDataType) -> bool {
1309 matches!(data_type, ConcreteDataType::Float64(_)) || is_native_histogram_value_type(data_type)
1310}
1311
1312pub(crate) fn retrieve_metric_name_and_result_type(
1313 promql_expr: &PromqlExpr,
1314) -> (Option<String>, ValueType) {
1315 (
1316 promql_expr_to_metric_name(promql_expr),
1317 promql_expr.value_type(),
1318 )
1319}
1320
1321pub(crate) fn get_catalog_schema(db: &Option<String>, ctx: &QueryContext) -> (String, String) {
1324 if let Some(db) = db {
1325 parse_catalog_and_schema_from_db_string(db)
1326 } else {
1327 (
1328 ctx.current_catalog().to_string(),
1329 ctx.current_schema().clone(),
1330 )
1331 }
1332}
1333
1334pub(crate) fn try_update_catalog_schema(ctx: &mut QueryContext, catalog: &str, schema: &str) {
1336 if ctx.current_catalog() != catalog || ctx.current_schema() != schema {
1337 ctx.set_current_catalog(catalog);
1338 ctx.set_current_schema(schema);
1339 }
1340}
1341
1342fn current_schema_metric_targets(
1343 query_ctx: &QueryContext,
1344 metric_names: &[String],
1345) -> PermissionTableTargets {
1346 PermissionTableTargets::resolved(
1347 metric_names
1348 .iter()
1349 .map(|metric| {
1350 PermissionTableTarget::new(
1351 query_ctx.current_catalog(),
1352 query_ctx.current_schema(),
1353 metric,
1354 )
1355 })
1356 .collect(),
1357 )
1358}
1359
1360fn is_internal_physical_metric_table(table: &TableRef) -> bool {
1361 table.table_info().is_physical_table()
1362}
1363
1364fn promql_expr_to_metric_name(expr: &PromqlExpr) -> Option<String> {
1365 let mut metric_names = HashSet::new();
1366 collect_metric_names(expr, &mut metric_names);
1367
1368 if metric_names.len() == 1 {
1370 metric_names.into_iter().next()
1371 } else {
1372 None
1373 }
1374}
1375
1376fn collect_metric_names(expr: &PromqlExpr, metric_names: &mut HashSet<String>) {
1378 match expr {
1379 PromqlExpr::Aggregate(AggregateExpr { modifier, expr, .. }) => {
1380 match modifier {
1381 Some(LabelModifier::Include(labels))
1382 if !labels.labels.contains(&METRIC_NAME.to_string()) =>
1383 {
1384 metric_names.clear();
1385 return;
1386 }
1387 Some(LabelModifier::Exclude(labels))
1388 if labels.labels.contains(&METRIC_NAME.to_string()) =>
1389 {
1390 metric_names.clear();
1391 return;
1392 }
1393 _ => {}
1394 }
1395 collect_metric_names(expr, metric_names)
1396 }
1397 PromqlExpr::Unary(UnaryExpr { .. }) => metric_names.clear(),
1398 PromqlExpr::Binary(BinaryExpr { lhs, op, .. }) => {
1399 if matches!(
1400 op.id(),
1401 token::T_LAND | token::T_LOR | token::T_LUNLESS ) {
1405 collect_metric_names(lhs, metric_names)
1406 } else {
1407 metric_names.clear()
1408 }
1409 }
1410 PromqlExpr::Paren(ParenExpr { expr }) => collect_metric_names(expr, metric_names),
1411 PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => collect_metric_names(expr, metric_names),
1412 PromqlExpr::VectorSelector(VectorSelector { name, matchers, .. }) => {
1413 if let Some(name) = name {
1414 metric_names.insert(name.clone());
1415 } else if let Some(matcher) = matchers.find_matchers(METRIC_NAME).into_iter().next() {
1416 metric_names.insert(matcher.value);
1417 }
1418 }
1419 PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1420 let VectorSelector { name, matchers, .. } = vs;
1421 if let Some(name) = name {
1422 metric_names.insert(name.clone());
1423 } else if let Some(matcher) = matchers.find_matchers(METRIC_NAME).into_iter().next() {
1424 metric_names.insert(matcher.value);
1425 }
1426 }
1427 PromqlExpr::Call(Call { args, .. }) => {
1428 args.args
1429 .iter()
1430 .for_each(|e| collect_metric_names(e, metric_names));
1431 }
1432 PromqlExpr::NumberLiteral(_) | PromqlExpr::StringLiteral(_) | PromqlExpr::Extension(_) => {}
1433 }
1434}
1435
1436fn find_metric_name_and_matchers<E, F>(expr: &PromqlExpr, f: F) -> Option<E>
1437where
1438 F: Fn(&Option<String>, &Matchers) -> Option<E> + Clone,
1439{
1440 match expr {
1441 PromqlExpr::Aggregate(AggregateExpr { expr, .. }) => find_metric_name_and_matchers(expr, f),
1442 PromqlExpr::Unary(UnaryExpr { expr }) => find_metric_name_and_matchers(expr, f),
1443 PromqlExpr::Binary(BinaryExpr { lhs, rhs, .. }) => {
1444 find_metric_name_and_matchers(lhs, f.clone()).or(find_metric_name_and_matchers(rhs, f))
1445 }
1446 PromqlExpr::Paren(ParenExpr { expr }) => find_metric_name_and_matchers(expr, f),
1447 PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => find_metric_name_and_matchers(expr, f),
1448 PromqlExpr::NumberLiteral(_) => None,
1449 PromqlExpr::StringLiteral(_) => None,
1450 PromqlExpr::Extension(_) => None,
1451 PromqlExpr::VectorSelector(VectorSelector { name, matchers, .. }) => f(name, matchers),
1452 PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1453 let VectorSelector { name, matchers, .. } = vs;
1454
1455 f(name, matchers)
1456 }
1457 PromqlExpr::Call(Call { args, .. }) => args
1458 .args
1459 .iter()
1460 .find_map(|e| find_metric_name_and_matchers(e, f.clone())),
1461 }
1462}
1463
1464struct MetricNameDiscovery {
1465 name_matchers: Vec<Matcher>,
1466 schema: Option<String>,
1467}
1468
1469fn find_metric_name_not_equal_matchers(expr: &PromqlExpr) -> Result<Option<MetricNameDiscovery>> {
1471 if find_metric_name_and_matchers(expr, |name, matchers| {
1472 (name.is_none() && !matchers.or_matchers.is_empty()).then_some(())
1473 })
1474 .is_some()
1475 {
1476 return Ok(None);
1477 }
1478
1479 find_metric_name_and_matchers(expr, |name, matchers| {
1480 if name.is_some() {
1482 return None;
1483 }
1484
1485 let name_matchers = matchers.find_matchers(METRIC_NAME);
1486 if name_matchers.len() != 1 || name_matchers[0].op == MatchOp::Equal {
1487 return None;
1488 }
1489
1490 Some(
1491 resolve_schema_from_matchers(&matchers.matchers).map(|schema| MetricNameDiscovery {
1492 name_matchers,
1493 schema,
1494 }),
1495 )
1496 })
1497 .transpose()
1498}
1499
1500fn expand_metric_name_queries(
1501 query: &ParsedPromQuery,
1502 metric_names: Vec<String>,
1503) -> Vec<ParsedPromQuery> {
1504 metric_names
1505 .into_iter()
1506 .map(|metric| {
1507 let mut query = query.clone();
1508 query.update_expr(|expr| update_metric_name_matcher(expr, &metric));
1509 query
1510 })
1511 .collect()
1512}
1513
1514fn static_promql_targets(
1515 expr: &PromqlExpr,
1516 query_ctx: &QueryContextRef,
1517) -> Result<(PermissionTableTargets, usize)> {
1518 let mut targets = Vec::new();
1519 let mut unresolved_selectors = 0;
1520 collect_static_promql_targets(expr, query_ctx, &mut targets, &mut unresolved_selectors)?;
1521 Ok((
1522 PermissionTableTargets::resolved(targets),
1523 unresolved_selectors,
1524 ))
1525}
1526
1527fn resolved_promql_targets(
1528 queries: &[ParsedPromQuery],
1529 query_ctx: &QueryContextRef,
1530) -> Option<Vec<PermissionTableTarget>> {
1531 let mut targets = Vec::new();
1532 for query in queries {
1533 let (PermissionTableTargets::Resolved(query_targets), unresolved_selectors) =
1534 static_promql_targets(query.expr(), query_ctx).ok()?
1535 else {
1536 return None;
1537 };
1538 if unresolved_selectors != 0 {
1539 return None;
1540 }
1541 for target in query_targets {
1542 if !targets.contains(&target) {
1543 targets.push(target);
1544 }
1545 }
1546 }
1547 Some(targets)
1548}
1549
1550fn collect_static_promql_targets(
1551 expr: &PromqlExpr,
1552 query_ctx: &QueryContextRef,
1553 targets: &mut Vec<PermissionTableTarget>,
1554 unresolved_selectors: &mut usize,
1555) -> Result<()> {
1556 match expr {
1557 PromqlExpr::Aggregate(AggregateExpr { expr, .. })
1558 | PromqlExpr::Unary(UnaryExpr { expr })
1559 | PromqlExpr::Paren(ParenExpr { expr })
1560 | PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => {
1561 collect_static_promql_targets(expr, query_ctx, targets, unresolved_selectors)?
1562 }
1563 PromqlExpr::Binary(BinaryExpr { lhs, rhs, .. }) => {
1564 collect_static_promql_targets(lhs, query_ctx, targets, unresolved_selectors)?;
1565 collect_static_promql_targets(rhs, query_ctx, targets, unresolved_selectors)?;
1566 }
1567 PromqlExpr::VectorSelector(selector) => {
1568 collect_static_vector_target(selector, query_ctx, targets, unresolved_selectors)?
1569 }
1570 PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1571 collect_static_vector_target(vs, query_ctx, targets, unresolved_selectors)?
1572 }
1573 PromqlExpr::Call(Call { args, .. }) => {
1574 for expr in &args.args {
1575 collect_static_promql_targets(expr, query_ctx, targets, unresolved_selectors)?;
1576 }
1577 }
1578 PromqlExpr::NumberLiteral(_) | PromqlExpr::StringLiteral(_) | PromqlExpr::Extension(_) => {}
1579 }
1580 Ok(())
1581}
1582
1583fn collect_static_vector_target(
1584 selector: &VectorSelector,
1585 query_ctx: &QueryContextRef,
1586 targets: &mut Vec<PermissionTableTarget>,
1587 unresolved_selectors: &mut usize,
1588) -> Result<()> {
1589 if selector.name.is_none() && !selector.matchers.or_matchers.is_empty() {
1590 *unresolved_selectors += 1;
1591 return Ok(());
1592 }
1593
1594 let metric = selector.name.clone().or_else(|| {
1595 let mut matchers = selector.matchers.find_matchers(METRIC_NAME);
1596 if matchers.len() == 1 && matchers[0].op == MatchOp::Equal {
1597 matchers.pop().map(|matcher| matcher.value)
1598 } else {
1599 None
1600 }
1601 });
1602 let Some(metric) = metric else {
1603 *unresolved_selectors += 1;
1604 return Ok(());
1605 };
1606
1607 let schema = resolve_schema_from_matchers(&selector.matchers.matchers)?
1608 .unwrap_or_else(|| query_ctx.current_schema());
1609 targets.push(PermissionTableTarget::new(
1610 query_ctx.current_catalog(),
1611 schema,
1612 metric,
1613 ));
1614 Ok(())
1615}
1616
1617fn update_metric_name_matcher(expr: &mut PromqlExpr, metric_name: &str) {
1620 match expr {
1621 PromqlExpr::Aggregate(AggregateExpr { expr, .. }) => {
1622 update_metric_name_matcher(expr, metric_name)
1623 }
1624 PromqlExpr::Unary(UnaryExpr { expr }) => update_metric_name_matcher(expr, metric_name),
1625 PromqlExpr::Binary(BinaryExpr { lhs, rhs, .. }) => {
1626 update_metric_name_matcher(lhs, metric_name);
1627 update_metric_name_matcher(rhs, metric_name);
1628 }
1629 PromqlExpr::Paren(ParenExpr { expr }) => update_metric_name_matcher(expr, metric_name),
1630 PromqlExpr::Subquery(SubqueryExpr { expr, .. }) => {
1631 update_metric_name_matcher(expr, metric_name)
1632 }
1633 PromqlExpr::VectorSelector(VectorSelector { name, matchers, .. }) => {
1634 if name.is_some() {
1635 return;
1636 }
1637
1638 for m in &mut matchers.matchers {
1639 if m.name == METRIC_NAME && m.op != MatchOp::Equal {
1640 m.op = MatchOp::Equal;
1641 m.value = metric_name.to_string();
1642 }
1643 }
1644 }
1645 PromqlExpr::MatrixSelector(MatrixSelector { vs, .. }) => {
1646 let VectorSelector { name, matchers, .. } = vs;
1647 if name.is_some() {
1648 return;
1649 }
1650
1651 for m in &mut matchers.matchers {
1652 if m.name == METRIC_NAME && m.op != MatchOp::Equal {
1653 m.op = MatchOp::Equal;
1654 m.value = metric_name.to_string();
1655 }
1656 }
1657 }
1658 PromqlExpr::Call(Call { args, .. }) => {
1659 args.args.iter_mut().for_each(|e| {
1660 update_metric_name_matcher(e, metric_name);
1661 });
1662 }
1663 PromqlExpr::NumberLiteral(_) | PromqlExpr::StringLiteral(_) | PromqlExpr::Extension(_) => {}
1664 }
1665}
1666
1667#[derive(Debug, Default, Serialize, Deserialize)]
1668pub struct LabelValueQuery {
1669 start: Option<String>,
1670 end: Option<String>,
1671 lookback: Option<String>,
1672 #[serde(flatten)]
1673 matches: Matches,
1674 db: Option<String>,
1675 limit: Option<usize>,
1676}
1677
1678#[axum_macros::debug_handler]
1679#[tracing::instrument(
1680 skip_all,
1681 fields(protocol = "prometheus", request_type = "label_values_query")
1682)]
1683pub async fn label_values_query(
1684 State(handler): State<PrometheusHandlerRef>,
1685 Path(label_name): Path<String>,
1686 Extension(mut query_ctx): Extension<QueryContext>,
1687 Query(params): Query<LabelValueQuery>,
1688) -> PrometheusJsonResponse {
1689 let (catalog, schema) = get_catalog_schema(¶ms.db, &query_ctx);
1690 try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
1691 let query_ctx = Arc::new(query_ctx);
1692
1693 let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
1694 .with_label_values(&[query_ctx.get_db_string().as_str(), "label_values_query"])
1695 .start_timer();
1696
1697 let matches = params.matches.0;
1698 if label_name == METRIC_NAME_LABEL {
1699 try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
1700 let exact_metric_names = matches
1701 .iter()
1702 .filter_map(|selector| retrieve_exact_metric_name_from_promql(selector))
1703 .collect::<Vec<_>>();
1704 try_call_return_response!(
1705 handler
1706 .check_query_target_permission(
1707 current_schema_metric_targets(&query_ctx, &exact_metric_names),
1708 &query_ctx,
1709 )
1710 .await
1711 );
1712 let catalog_manager = handler.catalog_manager();
1713
1714 let enumerate_all = matches.is_empty();
1717 let (metadata_selectors, label_selectors) =
1718 try_call_return_response!(split_selectors_by_label_use(&matches));
1719
1720 let mut table_names = if enumerate_all || !metadata_selectors.is_empty() {
1721 try_call_return_response!(
1722 retrieve_table_names(&query_ctx, catalog_manager, metadata_selectors).await
1723 )
1724 } else {
1725 Vec::new()
1726 };
1727
1728 if !label_selectors.is_empty() {
1729 table_names.extend(try_call_return_response!(
1730 retrieve_table_names_by_labels(
1731 &handler,
1732 label_selectors,
1733 params.start.as_deref(),
1734 params.end.as_deref(),
1735 &query_ctx,
1736 )
1737 .await
1738 ));
1739 table_names.sort_unstable();
1740 table_names.dedup();
1741 }
1742
1743 table_names = try_call_return_response!(
1744 handler
1745 .filter_metadata_metric_names(
1746 table_names,
1747 query_ctx.current_schema().as_str(),
1748 &query_ctx,
1749 )
1750 .await
1751 );
1752
1753 truncate_results(&mut table_names, params.limit);
1754 return PrometheusJsonResponse::success(PrometheusResponse::LabelValues(table_names));
1755 } else if label_name == FIELD_NAME_LABEL {
1756 try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
1757 let enumerate_all = matches.is_empty();
1758 let metric_names = matches
1759 .iter()
1760 .map(|selector| retrieve_exact_metric_name_from_promql(selector))
1761 .collect::<Vec<_>>();
1762 if metric_names.iter().any(Option::is_none) {
1763 try_call_return_response!(
1764 handler
1765 .check_query_target_permission(PermissionTableTargets::Unresolved, &query_ctx,)
1766 .await
1767 );
1768 }
1769 let metric_names = metric_names.into_iter().flatten().collect::<Vec<_>>();
1770
1771 try_call_return_response!(
1772 handler
1773 .check_query_target_permission(
1774 current_schema_metric_targets(&query_ctx, &metric_names),
1775 &query_ctx,
1776 )
1777 .await
1778 );
1779
1780 let (field_columns, table_names) = if !enumerate_all && metric_names.is_empty() {
1781 (HashSet::new(), Vec::new())
1782 } else {
1783 handle_schema_err!(
1784 retrieve_field_names(&query_ctx, handler.catalog_manager(), metric_names).await
1785 )
1786 .unwrap_or_default()
1787 };
1788 try_call_return_response!(
1789 handler
1790 .check_query_target_permission(
1791 current_schema_metric_targets(&query_ctx, &table_names),
1792 &query_ctx,
1793 )
1794 .await
1795 );
1796 let mut field_columns = field_columns.into_iter().collect::<Vec<_>>();
1797 field_columns.sort_unstable();
1798 truncate_results(&mut field_columns, params.limit);
1799 return PrometheusJsonResponse::success(PrometheusResponse::LabelValues(field_columns));
1800 } else if is_database_selection_label(&label_name) {
1801 try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
1802 let catalog_manager = handler.catalog_manager();
1803
1804 let (mut schema_names, targets) = try_call_return_response!(
1805 retrieve_schema_names(&query_ctx, catalog_manager, matches).await
1806 );
1807 try_call_return_response!(
1808 handler
1809 .check_query_target_permission(targets, &query_ctx)
1810 .await
1811 );
1812 truncate_results(&mut schema_names, params.limit);
1813 return PrometheusJsonResponse::success(PrometheusResponse::LabelValues(schema_names));
1814 }
1815
1816 let queries = matches;
1817 if queries.is_empty() {
1818 return PrometheusJsonResponse::error(
1819 StatusCode::InvalidArguments,
1820 "match[] parameter is required",
1821 );
1822 }
1823
1824 let start = params.start.unwrap_or_else(yesterday_rfc3339);
1825 let end = params.end.unwrap_or_else(current_time_rfc3339);
1826 let lookback = params
1827 .lookback
1828 .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string());
1829 let prom_queries = queries
1830 .into_iter()
1831 .map(|query| PromQuery {
1832 query,
1833 start: start.clone(),
1834 end: end.clone(),
1835 step: DEFAULT_LOOKBACK_STRING.to_string(),
1836 lookback: lookback.clone(),
1837 alias: None,
1838 })
1839 .collect::<Vec<_>>();
1840 let prom_queries = try_call_return_response!(
1841 prom_queries
1842 .into_iter()
1843 .map(|query| ParsedPromQuery::parse(query, &query_ctx))
1844 .collect::<Result<Vec<_>>>()
1845 );
1846 try_call_return_response!(
1847 handler
1848 .check_query_permission_parsed(&prom_queries, &query_ctx)
1849 .await
1850 );
1851
1852 let mut label_values = HashSet::new();
1853
1854 let Some(first_query) = prom_queries.first() else {
1855 unreachable!("empty match[] is rejected above")
1856 };
1857 let QueryStatement::Promql(eval_stmt, _) = first_query.statement() else {
1858 unreachable!("query is parsed from PromQL")
1859 };
1860 let start = eval_stmt.start;
1861 let end = eval_stmt.end;
1862
1863 for prom_query in prom_queries {
1864 let (_, statement) = prom_query.into_parts();
1865 let QueryStatement::Promql(eval_stmt, _) = statement else {
1866 unreachable!("query is parsed from PromQL")
1867 };
1868 let PromqlExpr::VectorSelector(mut vector_selector) = eval_stmt.expr else {
1869 return PrometheusJsonResponse::error(
1870 StatusCode::InvalidArguments,
1871 "expected vector selector",
1872 );
1873 };
1874 let Some(name) = take_metric_name(&mut vector_selector) else {
1875 return PrometheusJsonResponse::error(
1876 StatusCode::InvalidArguments,
1877 "expected metric name",
1878 );
1879 };
1880 let VectorSelector { matchers, .. } = vector_selector;
1881 let matchers = matchers.matchers;
1883 let result = handler
1884 .query_label_values(name, label_name.clone(), matchers, start, end, &query_ctx)
1885 .await;
1886 if let Some(result) = handle_schema_err!(result) {
1887 label_values.extend(result.into_iter());
1888 }
1889 }
1890
1891 let mut label_values: Vec<_> = label_values.into_iter().collect();
1892 label_values.sort_unstable();
1893 truncate_results(&mut label_values, params.limit);
1894
1895 PrometheusJsonResponse::success(PrometheusResponse::LabelValues(label_values))
1896}
1897
1898fn truncate_results(label_values: &mut Vec<String>, limit: Option<usize>) {
1899 if let Some(limit) = limit
1900 && limit > 0
1901 && label_values.len() >= limit
1902 {
1903 label_values.truncate(limit);
1904 }
1905}
1906
1907fn take_metric_name(selector: &mut VectorSelector) -> Option<String> {
1910 if let Some(name) = selector.name.take() {
1911 return Some(name);
1912 }
1913
1914 let (pos, matcher) = selector
1915 .matchers
1916 .matchers
1917 .iter()
1918 .find_position(|matcher| matcher.name == "__name__" && matcher.op == MatchOp::Equal)?;
1919 let name = matcher.value.clone();
1920 selector.matchers.matchers.remove(pos);
1922
1923 Some(name)
1924}
1925
1926fn take_metric_name_matchers(selector: &mut VectorSelector) -> Vec<Matcher> {
1930 let mut taken = Vec::new();
1931 if let Some(name) = selector.name.take() {
1932 taken.push(Matcher::new(MatchOp::Equal, METRIC_NAME_LABEL, &name));
1933 }
1934
1935 let (name_matchers, rest) = std::mem::take(&mut selector.matchers.matchers)
1936 .into_iter()
1937 .partition(|matcher| matcher.name == METRIC_NAME_LABEL);
1938 selector.matchers.matchers = rest;
1939 taken.extend(name_matchers);
1940
1941 taken
1942}
1943
1944fn metric_name_matches(table_name: &str, matchers: &[Matcher]) -> bool {
1951 matchers.iter().all(|matcher| match &matcher.op {
1952 MatchOp::Equal => table_name == matcher.value,
1953 MatchOp::NotEqual => table_name != matcher.value,
1954 MatchOp::Re(re) => re.is_match(table_name),
1955 MatchOp::NotRe(re) => !re.is_match(table_name),
1956 })
1957}
1958
1959fn is_ordinary_label_matcher(matcher: &Matcher) -> bool {
1962 matcher.name != METRIC_NAME_LABEL
1963 && matcher.name != FIELD_NAME_LABEL
1964 && !is_database_selection_label(&matcher.name)
1965}
1966
1967fn split_selectors_by_label_use(matches: &[String]) -> Result<(Vec<String>, Vec<VectorSelector>)> {
1974 let mut metadata_only = Vec::new();
1975 let mut with_labels = Vec::new();
1976
1977 for selector in matches {
1978 let expr = promql_parser::parser::parse(selector)
1979 .map_err(|reason| InvalidQuerySnafu { reason }.build())?;
1980 let PromqlExpr::VectorSelector(vector_selector) = expr else {
1981 return InvalidQuerySnafu {
1982 reason: "expected vector selector".to_string(),
1983 }
1984 .fail();
1985 };
1986
1987 let constrains_labels = vector_selector.matchers.or_matchers.is_empty()
1988 && vector_selector
1989 .matchers
1990 .matchers
1991 .iter()
1992 .any(is_ordinary_label_matcher);
1993 if constrains_labels {
1994 with_labels.push(vector_selector);
1995 } else {
1996 metadata_only.push(selector.clone());
1997 }
1998 }
1999
2000 Ok((metadata_only, with_labels))
2001}
2002
2003async fn retrieve_table_names_by_labels(
2007 handler: &PrometheusHandlerRef,
2008 selectors: Vec<VectorSelector>,
2009 start: Option<&str>,
2010 end: Option<&str>,
2011 query_ctx: &QueryContextRef,
2012) -> Result<Vec<String>> {
2013 let start_arg = start.map(str::to_string).unwrap_or_else(yesterday_rfc3339);
2014 let end_arg = end.map(str::to_string).unwrap_or_else(current_time_rfc3339);
2015 let start = QueryLanguageParser::parse_promql_timestamp(&start_arg).with_context(|_| {
2016 ParseTimestampSnafu {
2017 timestamp: start_arg.clone(),
2018 }
2019 })?;
2020 let end = QueryLanguageParser::parse_promql_timestamp(&end_arg).with_context(|_| {
2021 ParseTimestampSnafu {
2022 timestamp: end_arg.clone(),
2023 }
2024 })?;
2025
2026 let schema = query_ctx.current_schema();
2027 let mut table_names = Vec::new();
2028 for mut selector in selectors {
2029 let name_matchers = take_metric_name_matchers(&mut selector);
2030 let label_matchers = selector
2033 .matchers
2034 .matchers
2035 .into_iter()
2036 .filter(is_ordinary_label_matcher)
2037 .collect();
2038 let matched = handler
2039 .query_metric_names_by_labels(label_matchers, &schema, start, end, query_ctx)
2040 .await?;
2041 table_names.extend(
2042 matched
2043 .into_iter()
2044 .filter(|name| metric_name_matches(name, &name_matchers)),
2045 );
2046 }
2047
2048 Ok(table_names)
2049}
2050
2051async fn retrieve_table_names(
2052 query_ctx: &QueryContext,
2053 catalog_manager: CatalogManagerRef,
2054 matches: Vec<String>,
2055) -> Result<Vec<String>> {
2056 let catalog = query_ctx.current_catalog();
2057 let schema = query_ctx.current_schema();
2058
2059 let mut tables_stream = catalog_manager.tables(catalog, &schema, Some(query_ctx));
2060 let mut table_names = Vec::new();
2061
2062 let name_matchers = matches
2063 .iter()
2064 .map(|selector| {
2065 let expr = promql_parser::parser::parse(selector)
2066 .map_err(|reason| InvalidQuerySnafu { reason }.build())?;
2067 let PromqlExpr::VectorSelector(selector) = expr else {
2068 return InvalidQuerySnafu {
2069 reason: "expected vector selector".to_string(),
2070 }
2071 .fail();
2072 };
2073 if !selector.matchers.or_matchers.is_empty() {
2074 return Ok(None);
2075 }
2076 if let Some(name) = selector.name {
2077 return Ok(Some(Matcher::new(MatchOp::Equal, METRIC_NAME_LABEL, &name)));
2078 }
2079 Ok(selector
2080 .matchers
2081 .matchers
2082 .into_iter()
2083 .find(|matcher| matcher.name == METRIC_NAME_LABEL))
2084 })
2085 .collect::<Result<Vec<_>>>()?;
2086
2087 while let Some(table) = tables_stream.next().await {
2088 let table = table?;
2089 if !table
2090 .table_info()
2091 .meta
2092 .options
2093 .extra_options
2094 .contains_key(LOGICAL_TABLE_METADATA_KEY)
2095 || is_internal_physical_metric_table(&table)
2096 {
2097 continue;
2099 }
2100
2101 let table_name = &table.table_info().name;
2102
2103 if name_matchers.is_empty()
2104 || name_matchers.iter().any(|matcher| match matcher {
2105 None => true,
2106 Some(matcher) => match &matcher.op {
2107 MatchOp::Equal => table_name == &matcher.value,
2108 MatchOp::Re(reg) => reg.is_match(table_name),
2109 _ => true,
2112 },
2113 })
2114 {
2115 table_names.push(table_name.clone());
2116 }
2117 }
2118
2119 table_names.sort_unstable();
2120 Ok(table_names)
2121}
2122
2123async fn retrieve_metric_metadata(
2124 query_ctx: &QueryContext,
2125 manager: CatalogManagerRef,
2126 params: &MetadataQuery,
2127) -> Result<BTreeMap<String, Vec<PromMetadata>>> {
2128 let mut metadata = BTreeMap::new();
2129 if params.limit == Some(0) {
2130 return Ok(metadata);
2131 }
2132
2133 let catalog = query_ctx.current_catalog();
2134 let schema = query_ctx.current_schema();
2135
2136 if let Some(metric) = ¶ms.metric {
2137 let Some(table) = manager
2138 .table(catalog, &schema, metric, Some(query_ctx))
2139 .await?
2140 else {
2141 return Ok(metadata);
2142 };
2143 let table_info = table.table_info();
2144 if is_prometheus_metric_table(table_info.as_ref()) {
2145 metadata.insert(
2146 table_info.name.clone(),
2147 vec![prometheus_metadata_from_table(table_info.as_ref())],
2148 );
2149 }
2150 return Ok(metadata);
2151 }
2152
2153 let mut tables_stream = manager.tables(catalog, &schema, Some(query_ctx));
2154
2155 while let Some(table) = tables_stream.next().await {
2156 let table = table?;
2157 let table_info = table.table_info();
2158 if !is_prometheus_metric_table(table_info.as_ref()) {
2159 continue;
2160 }
2161
2162 metadata.insert(
2163 table_info.name.clone(),
2164 vec![prometheus_metadata_from_table(table_info.as_ref())],
2165 );
2166 }
2167
2168 Ok(metadata)
2169}
2170
2171fn is_prometheus_metric_table(table_info: &TableInfo) -> bool {
2172 table_info
2173 .meta
2174 .options
2175 .extra_options
2176 .contains_key(LOGICAL_TABLE_METADATA_KEY)
2177}
2178
2179fn prometheus_metadata_from_table(table_info: &TableInfo) -> PromMetadata {
2180 let options = &table_info.meta.options.extra_options;
2181 let metric_type = match options.get(SEMANTIC_METRIC_TYPE) {
2182 Some(metric_type)
2183 if options
2184 .get(SEMANTIC_METRIC_TEMPORALITY)
2185 .is_some_and(|temporality| {
2186 matches!(
2187 temporality.as_str(),
2188 METRIC_TEMPORALITY_DELTA | SEMANTIC_VALUE_MIXED
2189 )
2190 })
2191 && matches!(
2192 metric_type.as_str(),
2193 "counter" | "histogram" | "updown_counter"
2194 ) =>
2195 {
2196 "unknown".to_string()
2197 }
2198 Some(metric_type) => match metric_type.as_str() {
2199 "updown_counter" => "gauge".to_string(),
2200 "gauge_histogram" => "gaugehistogram".to_string(),
2201 "mixed" => "unknown".to_string(),
2202 metric_type => metric_type.to_string(),
2203 },
2204 None if table_has_native_histogram_value(table_info) => "histogram".to_string(),
2205 None => String::new(),
2206 };
2207 let unit = options
2208 .get(SEMANTIC_METRIC_UNIT)
2209 .map(|unit| ucum_to_openmetrics_unit(unit))
2210 .unwrap_or_default();
2211
2212 PromMetadata {
2213 metric_type,
2214 unit,
2215 help: String::new(),
2217 }
2218}
2219
2220fn table_has_native_histogram_value(table_info: &TableInfo) -> bool {
2221 table_info
2222 .meta
2223 .schema
2224 .column_schemas()
2225 .iter()
2226 .any(|column| is_native_histogram_value_type(&column.data_type))
2227}
2228
2229async fn retrieve_field_names(
2230 query_ctx: &QueryContext,
2231 manager: CatalogManagerRef,
2232 matches: Vec<String>,
2233) -> Result<(HashSet<String>, Vec<String>)> {
2234 let mut field_columns = HashSet::new();
2235 let mut table_names = Vec::new();
2236 let catalog = query_ctx.current_catalog();
2237 let schema = query_ctx.current_schema();
2238
2239 if matches.is_empty() {
2240 let mut tables = manager.tables(catalog, &schema, Some(query_ctx));
2242 while let Some(table) = tables.next().await {
2243 let table = table?;
2244 if is_internal_physical_metric_table(&table) {
2245 continue;
2246 }
2247 table_names.push(table.table_info().name.clone());
2248 for column in table.field_columns() {
2249 field_columns.insert(column.name);
2250 }
2251 }
2252 return Ok((field_columns, table_names));
2253 }
2254
2255 for table_name in matches {
2256 let table = manager
2257 .table(catalog, &schema, &table_name, Some(query_ctx))
2258 .await?
2259 .with_context(|| TableNotFoundSnafu {
2260 catalog: catalog.to_string(),
2261 schema: schema.clone(),
2262 table: table_name.clone(),
2263 })?;
2264
2265 if is_internal_physical_metric_table(&table) {
2266 continue;
2267 }
2268 table_names.push(table.table_info().name.clone());
2269 for column in table.field_columns() {
2270 field_columns.insert(column.name);
2271 }
2272 }
2273 Ok((field_columns, table_names))
2274}
2275
2276async fn retrieve_schema_names(
2277 query_ctx: &QueryContext,
2278 catalog_manager: CatalogManagerRef,
2279 matches: Vec<String>,
2280) -> Result<(Vec<String>, PermissionTableTargets)> {
2281 let mut schemas = Vec::new();
2282 let mut targets = Vec::new();
2283 let catalog = query_ctx.current_catalog();
2284 let metric_names = matches
2285 .iter()
2286 .map(|match_item| retrieve_exact_metric_name_from_promql(match_item))
2287 .collect::<Vec<_>>();
2288 let unresolved = metric_names.is_empty() || metric_names.iter().any(Option::is_none);
2289
2290 let candidate_schemas = catalog_manager
2291 .schema_names(catalog, Some(query_ctx))
2292 .await?;
2293
2294 for schema in candidate_schemas {
2295 let mut found = true;
2296 for table_name in metric_names.iter().flatten() {
2297 let table = catalog_manager
2298 .table(catalog, &schema, table_name, Some(query_ctx))
2299 .await?;
2300 let Some(table) = table else {
2301 found = false;
2302 continue;
2303 };
2304 targets.push(PermissionTableTarget::new(catalog, &schema, table_name));
2305 if is_internal_physical_metric_table(&table) {
2306 found = false;
2307 }
2308 }
2309
2310 if found {
2311 schemas.push(schema);
2312 }
2313 }
2314
2315 schemas.sort_unstable();
2316
2317 let targets = if unresolved {
2318 PermissionTableTargets::Unresolved
2319 } else {
2320 PermissionTableTargets::resolved(targets)
2321 };
2322 Ok((schemas, targets))
2323}
2324
2325fn retrieve_exact_metric_name_from_promql(query: &str) -> Option<String> {
2326 let PromqlExpr::VectorSelector(selector) = promql_parser::parser::parse(query).ok()? else {
2327 return None;
2328 };
2329 if let Some(name) = selector.name {
2330 return (!name.is_empty()).then_some(name);
2331 }
2332
2333 let mut name_matchers = selector.matchers.find_matchers(METRIC_NAME);
2334 if name_matchers.len() != 1 {
2335 return None;
2336 }
2337 let matcher = name_matchers.pop()?;
2338 (matcher.op == MatchOp::Equal && !matcher.value.is_empty()).then_some(matcher.value)
2339}
2340
2341#[derive(Debug, Default, Serialize, Deserialize)]
2342pub struct SeriesQuery {
2343 start: Option<String>,
2344 end: Option<String>,
2345 lookback: Option<String>,
2346 #[serde(flatten)]
2347 matches: Matches,
2348 db: Option<String>,
2349}
2350
2351#[axum_macros::debug_handler]
2352#[tracing::instrument(
2353 skip_all,
2354 fields(protocol = "prometheus", request_type = "series_query")
2355)]
2356pub async fn series_query(
2357 State(handler): State<PrometheusHandlerRef>,
2358 Query(params): Query<SeriesQuery>,
2359 Extension(mut query_ctx): Extension<QueryContext>,
2360 Form(form_params): Form<SeriesQuery>,
2361) -> PrometheusJsonResponse {
2362 let mut queries: Vec<String> = params.matches.0;
2363 if queries.is_empty() {
2364 queries = form_params.matches.0;
2365 }
2366 if queries.is_empty() {
2367 return PrometheusJsonResponse::error(
2368 StatusCode::Unsupported,
2369 "match[] parameter is required",
2370 );
2371 }
2372 let start = params
2373 .start
2374 .or(form_params.start)
2375 .unwrap_or_else(yesterday_rfc3339);
2376 let end = params
2377 .end
2378 .or(form_params.end)
2379 .unwrap_or_else(current_time_rfc3339);
2380 let lookback = params
2381 .lookback
2382 .or(form_params.lookback)
2383 .unwrap_or_else(|| DEFAULT_LOOKBACK_STRING.to_string());
2384
2385 if let Some(db) = ¶ms.db {
2387 let (catalog, schema) = parse_catalog_and_schema_from_db_string(db);
2388 try_update_catalog_schema(&mut query_ctx, &catalog, &schema);
2389 }
2390 let query_ctx = Arc::new(query_ctx);
2391
2392 let _timer = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
2393 .with_label_values(&[query_ctx.get_db_string().as_str(), "series_query"])
2394 .start_timer();
2395
2396 let prom_queries = queries
2397 .into_iter()
2398 .map(|query| PromQuery {
2399 query,
2400 start: start.clone(),
2401 end: end.clone(),
2402 step: DEFAULT_LOOKBACK_STRING.to_string(),
2404 lookback: lookback.clone(),
2405 alias: None,
2406 })
2407 .collect::<Vec<_>>();
2408 let mut expanded_queries = Vec::new();
2409 for prom_query in prom_queries {
2410 let prom_query = try_call_return_response!(
2411 @output_msg ParsedPromQuery::parse(prom_query, &query_ctx),
2412 StatusCode::InvalidArguments
2413 );
2414 let promql_expr = prom_query.expr();
2415 let Some(discovery) =
2416 try_call_return_response!(find_metric_name_not_equal_matchers(promql_expr))
2417 else {
2418 expanded_queries.push(prom_query);
2419 continue;
2420 };
2421
2422 try_call_return_response!(handler.check_query_permission(&[], &query_ctx).await);
2423 let (static_targets, unresolved_selectors) =
2424 try_call_return_response!(static_promql_targets(promql_expr, &query_ctx));
2425 try_call_return_response!(
2426 handler
2427 .check_query_target_permission(static_targets, &query_ctx)
2428 .await
2429 );
2430 if unresolved_selectors != 1 {
2431 try_call_return_response!(
2432 handler
2433 .check_query_target_permission(PermissionTableTargets::Unresolved, &query_ctx)
2434 .await
2435 );
2436 expanded_queries.push(prom_query);
2437 continue;
2438 }
2439
2440 let schema = discovery
2441 .schema
2442 .unwrap_or_else(|| query_ctx.current_schema());
2443 let metric_names = try_call_return_response!(
2444 handler
2445 .query_metric_names(discovery.name_matchers, &schema, &query_ctx)
2446 .await
2447 );
2448 expanded_queries.extend(expand_metric_name_queries(&prom_query, metric_names));
2449 }
2450 let prom_queries = expanded_queries;
2451 try_call_return_response!(
2452 handler
2453 .check_query_permission_parsed(&prom_queries, &query_ctx)
2454 .await
2455 );
2456
2457 let mut series = Vec::new();
2458 let mut merge_map = HashMap::new();
2459 for prom_query in prom_queries {
2460 let Some(target) = resolved_promql_targets(std::slice::from_ref(&prom_query), &query_ctx)
2461 .and_then(|targets| match targets.as_slice() {
2462 [target] => Some(target.clone()),
2463 _ => None,
2464 })
2465 else {
2466 return PrometheusJsonResponse::error(
2467 StatusCode::InvalidArguments,
2468 "series selector must resolve exactly one metric table",
2469 );
2470 };
2471 let result = handler.do_query_parsed(prom_query, query_ctx.clone()).await;
2472
2473 handle_schema_err!(
2474 retrieve_series_from_query_result(
2475 result,
2476 &mut series,
2477 &query_ctx,
2478 &target,
2479 &handler.catalog_manager(),
2480 &mut merge_map,
2481 )
2482 .await
2483 );
2484 }
2485 let merge_map = merge_map
2486 .into_iter()
2487 .map(|(k, v)| (k, Value::from(v)))
2488 .collect();
2489 let mut resp = PrometheusJsonResponse::success(PrometheusResponse::Series(series));
2490 resp.resp_metrics = merge_map;
2491 resp
2492}
2493
2494#[derive(Debug, Default, Serialize, Deserialize)]
2495pub struct ParseQuery {
2496 query: Option<String>,
2497 db: Option<String>,
2498}
2499
2500#[axum_macros::debug_handler]
2501#[tracing::instrument(
2502 skip_all,
2503 fields(protocol = "prometheus", request_type = "parse_query")
2504)]
2505pub async fn parse_query(
2506 State(_handler): State<PrometheusHandlerRef>,
2507 Query(params): Query<ParseQuery>,
2508 Extension(_query_ctx): Extension<QueryContext>,
2509 Form(form_params): Form<ParseQuery>,
2510) -> PrometheusJsonResponse {
2511 if let Some(query) = params.query.or(form_params.query) {
2512 let ast = try_call_return_response!(
2513 promql_parser::parser::parse(&query),
2514 StatusCode::InvalidArguments
2515 );
2516 PrometheusJsonResponse::success(PrometheusResponse::ParseResult(ast))
2517 } else {
2518 PrometheusJsonResponse::error(StatusCode::InvalidArguments, "query is required")
2519 }
2520}
2521
2522#[cfg(test)]
2523mod tests {
2524 use std::collections::HashSet;
2525 use std::sync::{Arc, Mutex};
2526
2527 use catalog::memory::MemoryCatalogManager;
2528 use catalog::{RegisterSchemaRequest, RegisterTableRequest};
2529 use common_catalog::consts::{DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME};
2530 use common_query::native_histogram::{
2531 CounterResetHint, NativeHistogram, build_histogram_array, native_histogram_value_type,
2532 };
2533 use common_query::prelude::greptime_native_histogram;
2534 use datatypes::prelude::ConcreteDataType;
2535 use datatypes::schema::{ColumnSchema, Schema};
2536 use datatypes::vectors::{StringVector, StructVector};
2537 use promql_parser::parser::value::ValueType;
2538 use table::metadata::{TableInfoBuilder, TableMetaBuilder, TableType, TableVersion};
2539 use table::requests::{
2540 SEMANTIC_METRIC_TEMPORALITY, SEMANTIC_METRIC_TYPE, SEMANTIC_METRIC_UNIT, TableOptions,
2541 };
2542 use table::test_util::EmptyTable;
2543 use table::test_util::table_info::test_table_info;
2544
2545 use super::*;
2546 use crate::prometheus_handler::PrometheusHandler;
2547
2548 struct TestCase {
2549 name: &'static str,
2550 promql: &'static str,
2551 expected_metric: Option<&'static str>,
2552 expected_type: ValueType,
2553 should_error: bool,
2554 }
2555
2556 struct TestPrometheusHandler {
2557 catalog_manager: CatalogManagerRef,
2558 deny_operation: bool,
2559 denied_table: Option<&'static str>,
2560 metric_names: Vec<String>,
2561 label_metric_names: Vec<String>,
2564 label_lookups: Mutex<Vec<Vec<Matcher>>>,
2565 queries: Mutex<Vec<String>>,
2566 ordered_outputs: Mutex<Vec<bool>>,
2567 }
2568
2569 #[async_trait::async_trait]
2570 impl PrometheusHandler for TestPrometheusHandler {
2571 async fn do_query(&self, _: &PromQuery, _: QueryContextRef) -> Result<Output> {
2572 unreachable!("HTTP handlers should execute the parsed query")
2573 }
2574
2575 async fn do_query_parsed(
2576 &self,
2577 query: ParsedPromQuery,
2578 _: QueryContextRef,
2579 ) -> Result<Output> {
2580 self.ordered_outputs
2581 .lock()
2582 .unwrap()
2583 .push(query.requires_output_ordering());
2584 self.queries
2585 .lock()
2586 .unwrap()
2587 .push(query.query().query.clone());
2588 Ok(Output::new_with_record_batches(RecordBatches::empty()))
2589 }
2590
2591 async fn check_query_permission(&self, _: &[PromQuery], _: &QueryContextRef) -> Result<()> {
2592 if self.deny_operation {
2593 return auth::error::PermissionDeniedSnafu
2594 .fail()
2595 .context(crate::error::AuthSnafu);
2596 }
2597 Ok(())
2598 }
2599
2600 async fn check_query_target_permission(
2601 &self,
2602 targets: PermissionTableTargets,
2603 _: &QueryContextRef,
2604 ) -> Result<()> {
2605 if matches!(
2606 targets,
2607 PermissionTableTargets::Resolved(targets)
2608 if self.denied_table.is_some_and(|denied| {
2609 targets.iter().any(|target| target.table == denied)
2610 })
2611 ) {
2612 return auth::error::PermissionDeniedSnafu
2613 .fail()
2614 .context(crate::error::AuthSnafu);
2615 }
2616 Ok(())
2617 }
2618
2619 async fn filter_metadata_metric_names(
2620 &self,
2621 metric_names: Vec<String>,
2622 _: &str,
2623 _: &QueryContextRef,
2624 ) -> Result<Vec<String>> {
2625 Ok(metric_names
2626 .into_iter()
2627 .filter(|metric| metric != "denied")
2628 .collect())
2629 }
2630
2631 async fn query_metric_names(
2632 &self,
2633 _: Vec<Matcher>,
2634 _: &str,
2635 _: &QueryContextRef,
2636 ) -> Result<Vec<String>> {
2637 Ok(self.metric_names.clone())
2638 }
2639
2640 async fn query_metric_names_by_labels(
2641 &self,
2642 matchers: Vec<Matcher>,
2643 _: &str,
2644 _: std::time::SystemTime,
2645 _: std::time::SystemTime,
2646 _: &QueryContextRef,
2647 ) -> Result<Vec<String>> {
2648 self.label_lookups.lock().unwrap().push(matchers);
2649 Ok(self.label_metric_names.clone())
2650 }
2651
2652 async fn query_label_values(
2653 &self,
2654 _: String,
2655 _: String,
2656 _: Vec<Matcher>,
2657 _: std::time::SystemTime,
2658 _: std::time::SystemTime,
2659 _: &QueryContextRef,
2660 ) -> Result<Vec<String>> {
2661 unreachable!()
2662 }
2663
2664 fn catalog_manager(&self) -> CatalogManagerRef {
2665 self.catalog_manager.clone()
2666 }
2667 }
2668
2669 fn permission_denied_response() -> PrometheusJsonResponse {
2670 let result: auth::error::Result<()> = auth::error::PermissionDeniedSnafu.fail();
2671 try_call_return_response!(result);
2672 unreachable!()
2673 }
2674
2675 fn invalid_promql_response() -> PrometheusJsonResponse {
2676 let result = ParsedPromQuery::parse(
2677 PromQuery {
2678 query: "up{".to_string(),
2679 ..Default::default()
2680 },
2681 &QueryContext::arc(),
2682 );
2683 try_call_return_response!(@output_msg result, StatusCode::InvalidArguments);
2684 unreachable!()
2685 }
2686
2687 #[test]
2688 fn test_try_call_return_response_preserves_status() {
2689 use axum::response::IntoResponse;
2690
2691 let response = permission_denied_response();
2692 assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
2693 assert_eq!(
2694 axum::http::StatusCode::FORBIDDEN,
2695 response.into_response().status()
2696 );
2697 }
2698
2699 #[test]
2700 fn test_try_call_return_response_preserves_error_source() {
2701 let response = invalid_promql_response();
2702 assert_eq!(Some(StatusCode::InvalidArguments), response.status_code);
2703 let error = response.error.unwrap();
2704 assert!(
2705 error.contains("unexpected end of input inside braces"),
2706 "{error}"
2707 );
2708 }
2709
2710 #[tokio::test]
2711 async fn range_query_does_not_require_execution_output_ordering() {
2712 let handler = Arc::new(TestPrometheusHandler {
2713 catalog_manager: MemoryCatalogManager::new(),
2714 deny_operation: false,
2715 denied_table: None,
2716 metric_names: Vec::new(),
2717 label_metric_names: Vec::new(),
2718 label_lookups: Mutex::new(Vec::new()),
2719 queries: Mutex::new(Vec::new()),
2720 ordered_outputs: Mutex::new(Vec::new()),
2721 });
2722 let state: PrometheusHandlerRef = handler.clone();
2723 instant_query(
2724 State(state.clone()),
2725 Query(InstantQuery {
2726 query: Some("sort(vector(1))".to_string()),
2727 time: Some("0".to_string()),
2728 ..Default::default()
2729 }),
2730 Extension(QueryContext::with(
2731 DEFAULT_CATALOG_NAME,
2732 DEFAULT_SCHEMA_NAME,
2733 )),
2734 Form(InstantQuery::default()),
2735 )
2736 .await;
2737
2738 for end in ["0", "1"] {
2740 range_query(
2741 State(state.clone()),
2742 Query(RangeQuery {
2743 query: Some("sort(vector(1))".to_string()),
2744 start: Some("0".to_string()),
2745 end: Some(end.to_string()),
2746 step: Some("1s".to_string()),
2747 ..Default::default()
2748 }),
2749 Extension(QueryContext::with(
2750 DEFAULT_CATALOG_NAME,
2751 DEFAULT_SCHEMA_NAME,
2752 )),
2753 Form(RangeQuery::default()),
2754 )
2755 .await;
2756 }
2757
2758 assert_eq!(
2760 *handler.ordered_outputs.lock().unwrap(),
2761 vec![true, false, false]
2762 );
2763 }
2764
2765 #[tokio::test]
2766 async fn test_promql_timer_records_parse_errors() {
2767 let handler: PrometheusHandlerRef = Arc::new(TestPrometheusHandler {
2768 catalog_manager: MemoryCatalogManager::new(),
2769 deny_operation: false,
2770 denied_table: None,
2771 metric_names: Vec::new(),
2772 label_metric_names: Vec::new(),
2773 label_lookups: Mutex::new(Vec::new()),
2774 queries: Mutex::new(Vec::new()),
2775 ordered_outputs: Mutex::new(Vec::new()),
2776 });
2777 let query_ctx = QueryContext::with("promql_timer_test", "parse_error");
2778 let db = query_ctx.get_db_string();
2779
2780 let instant_histogram = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
2781 .with_label_values(&[db.as_str(), "instant_query"]);
2782 let instant_count = instant_histogram.get_sample_count();
2783 let response = instant_query(
2784 State(handler.clone()),
2785 Query(InstantQuery {
2786 query: Some("up{".to_string()),
2787 ..Default::default()
2788 }),
2789 Extension(query_ctx),
2790 Form(InstantQuery::default()),
2791 )
2792 .await;
2793 assert_eq!(Some(StatusCode::InvalidArguments), response.status_code);
2794 assert_eq!(instant_count + 1, instant_histogram.get_sample_count());
2795
2796 let range_histogram = crate::metrics::METRIC_HTTP_PROMETHEUS_PROMQL_ELAPSED
2797 .with_label_values(&[db.as_str(), "range_query"]);
2798 let range_count = range_histogram.get_sample_count();
2799 let response = range_query(
2800 State(handler),
2801 Query(RangeQuery {
2802 query: Some("up{".to_string()),
2803 start: Some("0".to_string()),
2804 end: Some("1".to_string()),
2805 step: Some("1s".to_string()),
2806 ..Default::default()
2807 }),
2808 Extension(QueryContext::with("promql_timer_test", "parse_error")),
2809 Form(RangeQuery::default()),
2810 )
2811 .await;
2812 assert_eq!(Some(StatusCode::InvalidArguments), response.status_code);
2813 assert_eq!(range_count + 1, range_histogram.get_sample_count());
2814 }
2815
2816 #[tokio::test]
2817 async fn test_field_name_enumeration_fails_closed() {
2818 use axum::response::IntoResponse;
2819
2820 let mut allowed = test_table_info(
2821 1024,
2822 "allowed",
2823 DEFAULT_SCHEMA_NAME,
2824 DEFAULT_CATALOG_NAME,
2825 Arc::new(Schema::new(vec![])),
2826 );
2827 allowed.meta.options.extra_options.insert(
2828 LOGICAL_TABLE_METADATA_KEY.to_string(),
2829 "physical_metrics".to_string(),
2830 );
2831 let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&allowed));
2832 let mut denied = allowed.clone();
2833 denied.ident.table_id = 1025;
2834 denied.name = "denied".to_string();
2835 manager
2836 .register_table_sync(RegisterTableRequest {
2837 catalog: DEFAULT_CATALOG_NAME.to_string(),
2838 schema: DEFAULT_SCHEMA_NAME.to_string(),
2839 table_name: denied.name.clone(),
2840 table_id: denied.table_id(),
2841 table: EmptyTable::from_table_info(&denied),
2842 })
2843 .unwrap();
2844 let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
2845 assert_eq!(
2846 vec!["allowed".to_string(), "denied".to_string()],
2847 retrieve_table_names(&query_ctx, manager.clone(), Vec::new())
2848 .await
2849 .unwrap()
2850 );
2851
2852 let response = label_values_query(
2853 State(Arc::new(TestPrometheusHandler {
2854 catalog_manager: manager,
2855 deny_operation: false,
2856 denied_table: Some("denied"),
2857 metric_names: Vec::new(),
2858 label_metric_names: Vec::new(),
2859 label_lookups: Mutex::new(Vec::new()),
2860 queries: Mutex::new(Vec::new()),
2861 ordered_outputs: Mutex::new(Vec::new()),
2862 })),
2863 Path(FIELD_NAME_LABEL.to_string()),
2864 Extension(query_ctx),
2865 Query(LabelValueQuery::default()),
2866 )
2867 .await;
2868
2869 assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
2870 assert_eq!(
2871 axum::http::StatusCode::FORBIDDEN,
2872 response.into_response().status()
2873 );
2874 }
2875
2876 fn label_values_handler(label_metric_names: Vec<&str>) -> Arc<TestPrometheusHandler> {
2880 let mut cpu_user = test_table_info(
2881 1024,
2882 "cpu_user",
2883 DEFAULT_SCHEMA_NAME,
2884 DEFAULT_CATALOG_NAME,
2885 Arc::new(Schema::new(vec![])),
2886 );
2887 cpu_user.meta.options.extra_options.insert(
2888 LOGICAL_TABLE_METADATA_KEY.to_string(),
2889 "physical_metrics".to_string(),
2890 );
2891 let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&cpu_user));
2892 let mut cpu_system = cpu_user.clone();
2893 cpu_system.ident.table_id = 1025;
2894 cpu_system.name = "cpu_system".to_string();
2895 manager
2896 .register_table_sync(RegisterTableRequest {
2897 catalog: DEFAULT_CATALOG_NAME.to_string(),
2898 schema: DEFAULT_SCHEMA_NAME.to_string(),
2899 table_name: cpu_system.name.clone(),
2900 table_id: cpu_system.table_id(),
2901 table: EmptyTable::from_table_info(&cpu_system),
2902 })
2903 .unwrap();
2904
2905 Arc::new(TestPrometheusHandler {
2906 catalog_manager: manager,
2907 deny_operation: false,
2908 denied_table: None,
2909 metric_names: Vec::new(),
2910 label_metric_names: label_metric_names.into_iter().map(String::from).collect(),
2911 label_lookups: Mutex::new(Vec::new()),
2912 queries: Mutex::new(Vec::new()),
2913 ordered_outputs: Mutex::new(Vec::new()),
2914 })
2915 }
2916
2917 async fn query_metric_name_values(
2918 handler: Arc<TestPrometheusHandler>,
2919 matches: Vec<&str>,
2920 ) -> Vec<String> {
2921 let state: PrometheusHandlerRef = handler;
2922 let response = label_values_query(
2923 State(state),
2924 Path(METRIC_NAME_LABEL.to_string()),
2925 Extension(QueryContext::with(
2926 DEFAULT_CATALOG_NAME,
2927 DEFAULT_SCHEMA_NAME,
2928 )),
2929 Query(LabelValueQuery {
2930 matches: Matches(matches.into_iter().map(String::from).collect()),
2931 ..Default::default()
2932 }),
2933 )
2934 .await;
2935
2936 assert!(
2937 response.status_code.is_none(),
2938 "status={:?}, error={:?}",
2939 response.status_code,
2940 response.error
2941 );
2942 match response.data {
2943 PrometheusResponse::LabelValues(values) => values,
2944 other => panic!("expected label values, got {other:?}"),
2945 }
2946 }
2947
2948 #[tokio::test]
2949 async fn label_matchers_resolve_metric_names_from_data() {
2950 let handler = label_values_handler(vec!["cpu_user"]);
2951 let values = query_metric_name_values(handler.clone(), vec![r#"{pod="abc"}"#]).await;
2952
2953 assert_eq!(vec!["cpu_user".to_string()], values);
2955
2956 let lookups = handler.label_lookups.lock().unwrap();
2957 assert_eq!(1, lookups.len());
2958 assert_eq!(
2959 vec!["pod".to_string()],
2960 lookups[0]
2961 .iter()
2962 .map(|matcher| matcher.name.clone())
2963 .collect::<Vec<_>>()
2964 );
2965 }
2966
2967 #[tokio::test]
2968 async fn metric_name_matchers_narrow_data_resolved_names() {
2969 let handler = label_values_handler(vec!["cpu_user", "cpu_system"]);
2970 let values =
2971 query_metric_name_values(handler, vec![r#"{__name__=~"cpu_u.*", pod="abc"}"#]).await;
2972
2973 assert_eq!(vec!["cpu_user".to_string()], values);
2974 }
2975
2976 #[tokio::test]
2977 async fn special_matchers_are_stripped_before_the_data_lookup() {
2978 let handler = label_values_handler(vec!["cpu_user"]);
2979 let values = query_metric_name_values(
2980 handler.clone(),
2981 vec![r#"{pod="abc", __field__="value", __database__="public"}"#],
2982 )
2983 .await;
2984
2985 assert_eq!(vec!["cpu_user".to_string()], values);
2986 let lookups = handler.label_lookups.lock().unwrap();
2987 assert_eq!(
2988 vec!["pod".to_string()],
2989 lookups[0]
2990 .iter()
2991 .map(|matcher| matcher.name.clone())
2992 .collect::<Vec<_>>()
2993 );
2994 }
2995
2996 #[tokio::test]
2997 async fn database_and_field_matchers_stay_on_the_metadata_path() {
2998 let handler = label_values_handler(vec!["never_returned"]);
2999 let values = query_metric_name_values(
3000 handler.clone(),
3001 vec![r#"{__name__=~"cpu_.*", __field__="value", __database__="other"}"#],
3002 )
3003 .await;
3004
3005 assert_eq!(
3006 vec!["cpu_system".to_string(), "cpu_user".to_string()],
3007 values
3008 );
3009 assert!(handler.label_lookups.lock().unwrap().is_empty());
3010 }
3011
3012 #[tokio::test]
3013 async fn empty_match_still_enumerates_every_metric() {
3014 let handler = label_values_handler(Vec::new());
3015 let values = query_metric_name_values(handler, Vec::new()).await;
3016
3017 assert_eq!(
3018 vec!["cpu_system".to_string(), "cpu_user".to_string()],
3019 values
3020 );
3021 }
3022
3023 #[tokio::test]
3024 async fn test_series_query_expands_metric_name_regex() {
3025 let cpu_user = test_table_info(
3026 1024,
3027 "cpu_user",
3028 DEFAULT_SCHEMA_NAME,
3029 DEFAULT_CATALOG_NAME,
3030 Arc::new(Schema::new(vec![])),
3031 );
3032 let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&cpu_user));
3033 let mut cpu_system = cpu_user.clone();
3034 cpu_system.ident.table_id = 1025;
3035 cpu_system.name = "cpu_system".to_string();
3036 manager
3037 .register_table_sync(RegisterTableRequest {
3038 catalog: DEFAULT_CATALOG_NAME.to_string(),
3039 schema: DEFAULT_SCHEMA_NAME.to_string(),
3040 table_name: cpu_system.name.clone(),
3041 table_id: cpu_system.table_id(),
3042 table: EmptyTable::from_table_info(&cpu_system),
3043 })
3044 .unwrap();
3045
3046 let handler = Arc::new(TestPrometheusHandler {
3047 catalog_manager: manager,
3048 deny_operation: false,
3049 denied_table: None,
3050 metric_names: vec!["cpu_user".to_string(), "cpu_system".to_string()],
3051 label_metric_names: Vec::new(),
3052 label_lookups: Mutex::new(Vec::new()),
3053 queries: Mutex::new(Vec::new()),
3054 ordered_outputs: Mutex::new(Vec::new()),
3055 });
3056 let state: PrometheusHandlerRef = handler.clone();
3057 let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3058 let target_ctx = Arc::new(query_ctx.clone());
3059 let response = series_query(
3060 State(state),
3061 Query(SeriesQuery {
3062 matches: Matches(vec![r#"{__name__=~"cpu_.*"}"#.to_string()]),
3063 ..Default::default()
3064 }),
3065 Extension(query_ctx),
3066 Form(SeriesQuery::default()),
3067 )
3068 .await;
3069
3070 assert!(
3071 response.status_code.is_none(),
3072 "status={:?}, error={:?}",
3073 response.status_code,
3074 response.error
3075 );
3076 let mut targets = handler
3077 .queries
3078 .lock()
3079 .unwrap()
3080 .iter()
3081 .map(|query| {
3082 let query = ParsedPromQuery::parse(
3083 PromQuery {
3084 query: query.clone(),
3085 ..Default::default()
3086 },
3087 &target_ctx,
3088 )
3089 .unwrap();
3090 resolved_promql_targets(std::slice::from_ref(&query), &target_ctx)
3091 .unwrap()
3092 .pop()
3093 .unwrap()
3094 .table
3095 })
3096 .collect_vec();
3097 targets.sort_unstable();
3098 assert_eq!(vec!["cpu_system", "cpu_user"], targets);
3099 }
3100
3101 #[test]
3102 fn test_retrieve_metric_name_and_result_type() {
3103 let test_cases = &[
3104 TestCase {
3106 name: "simple metric",
3107 promql: "cpu_usage",
3108 expected_metric: Some("cpu_usage"),
3109 expected_type: ValueType::Vector,
3110 should_error: false,
3111 },
3112 TestCase {
3113 name: "metric with selector",
3114 promql: r#"cpu_usage{instance="localhost"}"#,
3115 expected_metric: Some("cpu_usage"),
3116 expected_type: ValueType::Vector,
3117 should_error: false,
3118 },
3119 TestCase {
3120 name: "metric with range selector",
3121 promql: "cpu_usage[5m]",
3122 expected_metric: Some("cpu_usage"),
3123 expected_type: ValueType::Matrix,
3124 should_error: false,
3125 },
3126 TestCase {
3127 name: "metric with __name__ matcher",
3128 promql: r#"{__name__="cpu_usage"}"#,
3129 expected_metric: Some("cpu_usage"),
3130 expected_type: ValueType::Vector,
3131 should_error: false,
3132 },
3133 TestCase {
3134 name: "metric with unary operator",
3135 promql: "-cpu_usage",
3136 expected_metric: None,
3137 expected_type: ValueType::Vector,
3138 should_error: false,
3139 },
3140 TestCase {
3142 name: "metric with aggregation",
3143 promql: "sum(cpu_usage)",
3144 expected_metric: Some("cpu_usage"),
3145 expected_type: ValueType::Vector,
3146 should_error: false,
3147 },
3148 TestCase {
3149 name: "complex aggregation",
3150 promql: r#"sum by (instance) (cpu_usage{job="node"})"#,
3151 expected_metric: None,
3152 expected_type: ValueType::Vector,
3153 should_error: false,
3154 },
3155 TestCase {
3156 name: "complex aggregation",
3157 promql: r#"sum by (__name__) (cpu_usage{job="node"})"#,
3158 expected_metric: Some("cpu_usage"),
3159 expected_type: ValueType::Vector,
3160 should_error: false,
3161 },
3162 TestCase {
3163 name: "complex aggregation",
3164 promql: r#"sum without (instance) (cpu_usage{job="node"})"#,
3165 expected_metric: Some("cpu_usage"),
3166 expected_type: ValueType::Vector,
3167 should_error: false,
3168 },
3169 TestCase {
3171 name: "same metric addition",
3172 promql: "cpu_usage + cpu_usage",
3173 expected_metric: None,
3174 expected_type: ValueType::Vector,
3175 should_error: false,
3176 },
3177 TestCase {
3178 name: "metric with scalar addition",
3179 promql: r#"sum(rate(cpu_usage{job="node"}[5m])) + 100"#,
3180 expected_metric: None,
3181 expected_type: ValueType::Vector,
3182 should_error: false,
3183 },
3184 TestCase {
3186 name: "different metrics addition",
3187 promql: "cpu_usage + memory_usage",
3188 expected_metric: None,
3189 expected_type: ValueType::Vector,
3190 should_error: false,
3191 },
3192 TestCase {
3193 name: "different metrics subtraction",
3194 promql: "network_in - network_out",
3195 expected_metric: None,
3196 expected_type: ValueType::Vector,
3197 should_error: false,
3198 },
3199 TestCase {
3201 name: "unless with different metrics",
3202 promql: "cpu_usage unless memory_usage",
3203 expected_metric: Some("cpu_usage"),
3204 expected_type: ValueType::Vector,
3205 should_error: false,
3206 },
3207 TestCase {
3208 name: "unless with same metric",
3209 promql: "cpu_usage unless cpu_usage",
3210 expected_metric: Some("cpu_usage"),
3211 expected_type: ValueType::Vector,
3212 should_error: false,
3213 },
3214 TestCase {
3216 name: "basic subquery",
3217 promql: "cpu_usage[5m:1m]",
3218 expected_metric: Some("cpu_usage"),
3219 expected_type: ValueType::Matrix,
3220 should_error: false,
3221 },
3222 TestCase {
3223 name: "subquery with multiple metrics",
3224 promql: "(cpu_usage + memory_usage)[5m:1m]",
3225 expected_metric: None,
3226 expected_type: ValueType::Matrix,
3227 should_error: false,
3228 },
3229 TestCase {
3231 name: "scalar value",
3232 promql: "42",
3233 expected_metric: None,
3234 expected_type: ValueType::Scalar,
3235 should_error: false,
3236 },
3237 TestCase {
3238 name: "string literal",
3239 promql: r#""hello world""#,
3240 expected_metric: None,
3241 expected_type: ValueType::String,
3242 should_error: false,
3243 },
3244 TestCase {
3246 name: "invalid syntax",
3247 promql: "cpu_usage{invalid=",
3248 expected_metric: None,
3249 expected_type: ValueType::Vector,
3250 should_error: true,
3251 },
3252 TestCase {
3253 name: "empty query",
3254 promql: "",
3255 expected_metric: None,
3256 expected_type: ValueType::Vector,
3257 should_error: true,
3258 },
3259 TestCase {
3260 name: "malformed brackets",
3261 promql: "cpu_usage[5m",
3262 expected_metric: None,
3263 expected_type: ValueType::Vector,
3264 should_error: true,
3265 },
3266 ];
3267
3268 for test_case in test_cases {
3269 let result = promql_parser::parser::parse(test_case.promql)
3270 .map(|expr| retrieve_metric_name_and_result_type(&expr));
3271
3272 if test_case.should_error {
3273 assert!(
3274 result.is_err(),
3275 "Test '{}' should have failed but succeeded with: {:?}",
3276 test_case.name,
3277 result
3278 );
3279 } else {
3280 let (metric_name, value_type) = result.unwrap_or_else(|e| {
3281 panic!(
3282 "Test '{}' should have succeeded but failed with error: {}",
3283 test_case.name, e
3284 )
3285 });
3286
3287 let expected_metric_name = test_case.expected_metric.map(|s| s.to_string());
3288 assert_eq!(
3289 metric_name, expected_metric_name,
3290 "Test '{}': metric name mismatch. Expected: {:?}, Got: {:?}",
3291 test_case.name, expected_metric_name, metric_name
3292 );
3293
3294 assert_eq!(
3295 value_type, test_case.expected_type,
3296 "Test '{}': value type mismatch. Expected: {:?}, Got: {:?}",
3297 test_case.name, test_case.expected_type, value_type
3298 );
3299 }
3300 }
3301 }
3302
3303 #[tokio::test]
3304 async fn test_get_all_column_names_uses_tag_columns() {
3305 let schema = Arc::new(Schema::new(vec![
3306 ColumnSchema::new(
3307 "greptime_timestamp",
3308 ConcreteDataType::timestamp_millisecond_datatype(),
3309 false,
3310 )
3311 .with_time_index(true),
3312 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
3313 ColumnSchema::new("region", ConcreteDataType::string_datatype(), false),
3314 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
3315 ColumnSchema::new(
3316 DATA_SCHEMA_TSID_COLUMN_NAME,
3317 ConcreteDataType::uint64_datatype(),
3318 true,
3319 ),
3320 ]));
3321 let mut options = TableOptions::default();
3322 options.extra_options.insert(
3323 LOGICAL_TABLE_METADATA_KEY.to_string(),
3324 "physical_metrics".to_string(),
3325 );
3326 let meta = TableMetaBuilder::empty()
3327 .schema(schema)
3328 .primary_key_indices(vec![1, 2, 4])
3329 .engine("metric".to_string())
3330 .next_column_id(5)
3331 .options(options)
3332 .build()
3333 .unwrap();
3334 let table_info = TableInfoBuilder::default()
3335 .table_id(1024)
3336 .table_version(0 as TableVersion)
3337 .name("cpu_usage")
3338 .catalog_name(DEFAULT_CATALOG_NAME)
3339 .schema_name(DEFAULT_SCHEMA_NAME)
3340 .table_type(TableType::Base)
3341 .meta(meta)
3342 .build()
3343 .unwrap();
3344 let manager =
3345 MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&table_info));
3346
3347 let physical_schema = Arc::new(Schema::new(vec![
3348 ColumnSchema::new(
3349 "greptime_timestamp",
3350 ConcreteDataType::timestamp_millisecond_datatype(),
3351 false,
3352 )
3353 .with_time_index(true),
3354 ColumnSchema::new(
3355 "physical_only_tag",
3356 ConcreteDataType::string_datatype(),
3357 false,
3358 ),
3359 ColumnSchema::new(
3360 "physical_only_value",
3361 ConcreteDataType::float64_datatype(),
3362 true,
3363 ),
3364 ]));
3365 let mut physical_options = TableOptions::default();
3366 physical_options.extra_options.insert(
3367 store_api::metric_engine_consts::PHYSICAL_TABLE_METADATA_KEY.to_string(),
3368 String::new(),
3369 );
3370 let physical_meta = TableMetaBuilder::empty()
3371 .schema(physical_schema)
3372 .primary_key_indices(vec![1])
3373 .engine("metric".to_string())
3374 .next_column_id(3)
3375 .options(physical_options)
3376 .build()
3377 .unwrap();
3378 let physical_table_info = TableInfoBuilder::default()
3379 .table_id(1025)
3380 .table_version(0 as TableVersion)
3381 .name("physical_metrics")
3382 .catalog_name(DEFAULT_CATALOG_NAME)
3383 .schema_name(DEFAULT_SCHEMA_NAME)
3384 .table_type(TableType::Base)
3385 .meta(physical_meta)
3386 .build()
3387 .unwrap();
3388 let physical_table = EmptyTable::from_table_info(&physical_table_info);
3389 manager
3390 .register_table_sync(RegisterTableRequest {
3391 catalog: DEFAULT_CATALOG_NAME.to_string(),
3392 schema: DEFAULT_SCHEMA_NAME.to_string(),
3393 table_name: "physical_metrics".to_string(),
3394 table_id: 1025,
3395 table: physical_table,
3396 })
3397 .unwrap();
3398 let manager: CatalogManagerRef = manager;
3399
3400 let (labels, table_names) =
3401 get_all_column_names(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME, &manager)
3402 .await
3403 .unwrap();
3404
3405 assert_eq!(
3406 labels,
3407 HashSet::from(["host".to_string(), "region".to_string()])
3408 );
3409 assert_eq!(table_names, vec!["cpu_usage"]);
3410
3411 let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3412 let (schemas, targets) =
3413 retrieve_schema_names(&query_ctx, manager.clone(), vec!["cpu_usage".to_string()])
3414 .await
3415 .unwrap();
3416 assert_eq!(schemas, vec![DEFAULT_SCHEMA_NAME]);
3417 assert_eq!(
3418 targets,
3419 PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3420 DEFAULT_CATALOG_NAME,
3421 DEFAULT_SCHEMA_NAME,
3422 "cpu_usage",
3423 )])
3424 );
3425
3426 let (_, targets) = retrieve_schema_names(
3427 &query_ctx,
3428 manager.clone(),
3429 vec![r#"{__name__=~"cpu.*"}"#.to_string()],
3430 )
3431 .await
3432 .unwrap();
3433 assert_eq!(targets, PermissionTableTargets::Unresolved);
3434
3435 let (schemas, targets) = retrieve_schema_names(&query_ctx, manager.clone(), vec![])
3436 .await
3437 .unwrap();
3438 assert!(schemas.contains(&DEFAULT_SCHEMA_NAME.to_string()));
3439 assert_eq!(targets, PermissionTableTargets::Unresolved);
3440
3441 let (schemas, targets) = retrieve_schema_names(
3442 &query_ctx,
3443 manager.clone(),
3444 vec!["physical_metrics".to_string()],
3445 )
3446 .await
3447 .unwrap();
3448 assert!(schemas.is_empty());
3449 assert_eq!(
3450 targets,
3451 PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3452 DEFAULT_CATALOG_NAME,
3453 DEFAULT_SCHEMA_NAME,
3454 "physical_metrics",
3455 )])
3456 );
3457
3458 let (schemas, targets) = retrieve_schema_names(
3459 &query_ctx,
3460 manager.clone(),
3461 vec!["missing".to_string(), "physical_metrics".to_string()],
3462 )
3463 .await
3464 .unwrap();
3465 assert!(schemas.is_empty());
3466 assert_eq!(
3467 targets,
3468 PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3469 DEFAULT_CATALOG_NAME,
3470 DEFAULT_SCHEMA_NAME,
3471 "physical_metrics",
3472 )])
3473 );
3474
3475 assert_eq!(
3476 vec!["cpu_usage".to_string()],
3477 retrieve_table_names(
3478 &query_ctx,
3479 manager.clone(),
3480 vec!["missing".to_string(), r#"{__name__=~"cpu.*"}"#.to_string(),],
3481 )
3482 .await
3483 .unwrap()
3484 );
3485 assert!(
3486 retrieve_table_names(
3487 &query_ctx,
3488 manager.clone(),
3489 vec!["cpu_usage".to_string(), "{".to_string()],
3490 )
3491 .await
3492 .is_err()
3493 );
3494
3495 let (fields, table_names) = retrieve_field_names(
3496 &query_ctx,
3497 manager.clone(),
3498 vec!["physical_metrics".to_string()],
3499 )
3500 .await
3501 .unwrap();
3502 assert!(fields.is_empty());
3503 assert!(table_names.is_empty());
3504
3505 let (fields, table_names) = retrieve_field_names(&query_ctx, manager, vec![])
3506 .await
3507 .unwrap();
3508 assert_eq!(fields, HashSet::from(["value".to_string()]));
3509 assert_eq!(table_names, vec!["cpu_usage"]);
3510 }
3511
3512 #[tokio::test]
3513 async fn test_get_target_column_names_uses_selected_schema() {
3514 let manager = MemoryCatalogManager::with_default_setup();
3515 manager
3516 .register_schema_sync(RegisterSchemaRequest {
3517 catalog: DEFAULT_CATALOG_NAME.to_string(),
3518 schema: "private".to_string(),
3519 })
3520 .unwrap();
3521
3522 for (schema_name, tag_name, table_id) in [
3523 (DEFAULT_SCHEMA_NAME, "public_tag", 2048),
3524 ("private", "private_tag", 2049),
3525 ] {
3526 let schema = Arc::new(Schema::new(vec![
3527 ColumnSchema::new(
3528 "greptime_timestamp",
3529 ConcreteDataType::timestamp_millisecond_datatype(),
3530 false,
3531 )
3532 .with_time_index(true),
3533 ColumnSchema::new(tag_name, ConcreteDataType::string_datatype(), false),
3534 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
3535 ]));
3536 let meta = TableMetaBuilder::empty()
3537 .schema(schema)
3538 .primary_key_indices(vec![1])
3539 .engine("metric".to_string())
3540 .next_column_id(3)
3541 .build()
3542 .unwrap();
3543 let table_info = TableInfoBuilder::default()
3544 .table_id(table_id)
3545 .table_version(0 as TableVersion)
3546 .name("cpu")
3547 .catalog_name(DEFAULT_CATALOG_NAME)
3548 .schema_name(schema_name)
3549 .table_type(TableType::Base)
3550 .meta(meta)
3551 .build()
3552 .unwrap();
3553 manager
3554 .register_table_sync(RegisterTableRequest {
3555 catalog: DEFAULT_CATALOG_NAME.to_string(),
3556 schema: schema_name.to_string(),
3557 table_name: "cpu".to_string(),
3558 table_id,
3559 table: EmptyTable::from_table_info(&table_info),
3560 })
3561 .unwrap();
3562 }
3563
3564 let query_ctx = Arc::new(QueryContext::with(
3565 DEFAULT_CATALOG_NAME,
3566 DEFAULT_SCHEMA_NAME,
3567 ));
3568 let queries = [ParsedPromQuery::parse(
3569 PromQuery {
3570 query: r#"{__name__="cpu",__schema__="private"}"#.to_string(),
3571 ..Default::default()
3572 },
3573 &query_ctx,
3574 )
3575 .unwrap()];
3576 let targets = resolved_promql_targets(&queries, &query_ctx).unwrap();
3577 let manager: CatalogManagerRef = manager;
3578
3579 assert_eq!(
3580 get_target_column_names(&targets, &manager, &query_ctx)
3581 .await
3582 .unwrap(),
3583 HashSet::from(["private_tag".to_string()])
3584 );
3585 }
3586
3587 #[test]
3588 fn test_metric_name_discovery_uses_selected_schema() {
3589 let expr =
3590 promql_parser::parser::parse(r#"{__name__=~"cpu.*",__database__="private"}"#).unwrap();
3591 let discovery = find_metric_name_not_equal_matchers(&expr).unwrap().unwrap();
3592
3593 assert_eq!(discovery.schema.as_deref(), Some("private"));
3594 assert_eq!(discovery.name_matchers.len(), 1);
3595 assert_eq!(discovery.name_matchers[0].value, "cpu.*");
3596 assert!(matches!(discovery.name_matchers[0].op, MatchOp::Re(_)));
3597
3598 let expr = promql_parser::parser::parse(
3599 r#"{__name__=~"cpu.*",__database__="private",__schema__="public"}"#,
3600 )
3601 .unwrap();
3602 let discovery = find_metric_name_not_equal_matchers(&expr).unwrap().unwrap();
3603 assert_eq!(discovery.schema.as_deref(), Some("public"));
3604 }
3605
3606 #[test]
3607 fn test_metric_name_discovery_rejects_unsafe_selector_shapes() {
3608 for promql in [
3609 r#"{__name__="denied",__name__=~"no_such_.*"}"#,
3610 r#"{__name__=~"cpu.*" or job="api"}"#,
3611 ] {
3612 let expr = promql_parser::parser::parse(promql).unwrap();
3613 assert!(
3614 find_metric_name_not_equal_matchers(&expr)
3615 .unwrap()
3616 .is_none(),
3617 "{promql}"
3618 );
3619 }
3620
3621 let expr = promql_parser::parser::parse(r#"{__name__=~"cpu.*" or job="api"}"#).unwrap();
3622 assert_eq!(
3623 (PermissionTableTargets::Resolved(vec![]), 1),
3624 static_promql_targets(&expr, &QueryContext::arc()).unwrap()
3625 );
3626
3627 let expr =
3628 promql_parser::parser::parse(r#"allowed{job="api" or instance="host"}"#).unwrap();
3629 assert_eq!(
3630 (
3631 PermissionTableTargets::Resolved(vec![PermissionTableTarget::new(
3632 "greptime", "public", "allowed",
3633 )]),
3634 0,
3635 ),
3636 static_promql_targets(&expr, &QueryContext::arc()).unwrap()
3637 );
3638 }
3639
3640 #[test]
3641 fn test_metric_name_discovery_ignores_or_matchers_on_named_selector() {
3642 let expr = promql_parser::parser::parse(
3643 r#"{__name__=~"cpu.*"} + allowed{job="api" or instance="host"}"#,
3644 )
3645 .unwrap();
3646
3647 let discovery = find_metric_name_not_equal_matchers(&expr).unwrap().unwrap();
3648 assert_eq!(discovery.name_matchers[0].value, "cpu.*");
3649 }
3650
3651 #[test]
3652 fn test_retrieve_exact_metric_name_from_selector() {
3653 assert_eq!(
3654 retrieve_exact_metric_name_from_promql(r#"cpu{host="a"}"#).as_deref(),
3655 Some("cpu")
3656 );
3657 assert_eq!(
3658 retrieve_exact_metric_name_from_promql(r#"{__name__="cpu",host="a"}"#).as_deref(),
3659 Some("cpu")
3660 );
3661 assert!(retrieve_exact_metric_name_from_promql(r#"{__name__=~"cpu.*"}"#).is_none());
3662 }
3663
3664 #[test]
3665 fn test_expand_metric_name_queries_preserves_schema_matcher() {
3666 let query = PromQuery {
3667 query: r#"{__name__=~"cpu.*",__database__="private"}"#.to_string(),
3668 ..Default::default()
3669 };
3670 let ctx = QueryContext::arc();
3671 let query = ParsedPromQuery::parse(query, &ctx).unwrap();
3672 let queries = expand_metric_name_queries(
3673 &query,
3674 vec!["cpu_user".to_string(), "cpu_system".to_string()],
3675 );
3676
3677 assert_eq!(queries.len(), 2);
3678 for (query, metric_name) in queries.iter().zip(["cpu_user", "cpu_system"]) {
3679 let PromqlExpr::VectorSelector(selector) = query.expr() else {
3680 panic!("expected vector selector");
3681 };
3682 assert!(selector.matchers.matchers.iter().any(|matcher| {
3683 matcher.name == METRIC_NAME
3684 && matcher.op == MatchOp::Equal
3685 && matcher.value == metric_name
3686 }));
3687 assert!(selector.matchers.matchers.iter().any(|matcher| {
3688 matcher.name == "__database__"
3689 && matcher.op == MatchOp::Equal
3690 && matcher.value == "private"
3691 }));
3692 }
3693 }
3694
3695 #[test]
3696 fn test_expand_metric_name_queries_preserves_exact_selector() {
3697 let query = PromQuery {
3698 query: r#"denied + {__name__=~"cpu.*"}"#.to_string(),
3699 ..Default::default()
3700 };
3701 let ctx = QueryContext::arc();
3702 let query = ParsedPromQuery::parse(query, &ctx).unwrap();
3703 let queries = expand_metric_name_queries(&query, vec!["cpu_user".to_string()]);
3704
3705 assert_eq!(queries.len(), 1);
3706 let expanded = queries[0].expr();
3707 let ctx = Arc::new(QueryContext::with("greptime", "public"));
3708 let (PermissionTableTargets::Resolved(mut targets), unresolved_selectors) =
3709 static_promql_targets(expanded, &ctx).unwrap()
3710 else {
3711 panic!("expected resolved targets");
3712 };
3713 assert_eq!(unresolved_selectors, 0);
3714 targets.sort_unstable_by(|left, right| left.table.cmp(&right.table));
3715 assert_eq!(
3716 targets,
3717 vec![
3718 PermissionTableTarget::new("greptime", "public", "cpu_user"),
3719 PermissionTableTarget::new("greptime", "public", "denied"),
3720 ]
3721 );
3722 }
3723
3724 #[test]
3725 fn test_static_promql_targets_keep_exact_selector_during_discovery() {
3726 let expr = promql_parser::parser::parse(
3727 r#"denied + {__name__=~"no_match_.*",__schema__="private"}"#,
3728 )
3729 .unwrap();
3730 let ctx = Arc::new(QueryContext::with("greptime", "public"));
3731
3732 assert_eq!(
3733 static_promql_targets(&expr, &ctx).unwrap(),
3734 (
3735 PermissionTableTargets::resolved(vec![PermissionTableTarget::new(
3736 "greptime", "public", "denied",
3737 )]),
3738 1,
3739 )
3740 );
3741 }
3742
3743 #[test]
3744 fn test_static_promql_targets_count_unresolved_selectors() {
3745 let expr =
3746 promql_parser::parser::parse(r#"{__name__=~"no_such_metric"} or {job="api"}"#).unwrap();
3747 let ctx = Arc::new(QueryContext::with("greptime", "public"));
3748
3749 assert_eq!(
3750 static_promql_targets(&expr, &ctx).unwrap(),
3751 (PermissionTableTargets::resolved(Vec::new()), 2)
3752 );
3753 }
3754
3755 #[test]
3756 fn test_record_batches_to_labels_name_accepts_native_histogram_value() {
3757 let schema = Arc::new(Schema::new(vec![
3758 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
3759 ColumnSchema::new(
3760 greptime_native_histogram(),
3761 native_histogram_value_type().clone(),
3762 true,
3763 ),
3764 ]));
3765 let batch = RecordBatch::new_empty(schema.clone());
3766 let batches = RecordBatches::try_new(schema.clone(), vec![batch]).unwrap();
3767 let mut labels = HashSet::new();
3768
3769 record_batches_to_labels_name(batches, &mut labels).unwrap();
3770
3771 assert!(labels.is_empty());
3772
3773 let histogram = NativeHistogram {
3774 schema: 0,
3775 zero_threshold: 0.001,
3776 sum: 1.0,
3777 reset_hint: CounterResetHint::Unknown,
3778 start_timestamp: None,
3779 custom_values: vec![],
3780 positive_spans: vec![],
3781 negative_spans: vec![],
3782 count: 1.0,
3783 zero_count: 1.0,
3784 positive_buckets: vec![],
3785 negative_buckets: vec![],
3786 };
3787 let histogram_array = build_histogram_array(&[Some(histogram)]);
3788 let histogram_array = histogram_array
3789 .as_any()
3790 .downcast_ref::<arrow::array::StructArray>()
3791 .unwrap()
3792 .clone();
3793 let ConcreteDataType::Struct(histogram_type) = native_histogram_value_type().clone() else {
3794 unreachable!("native histogram type must be a struct")
3795 };
3796 let batch = RecordBatch::new(
3797 schema.clone(),
3798 vec![
3799 Arc::new(StringVector::from(vec![Some("localhost")])) as _,
3800 Arc::new(StructVector::try_new(histogram_type, histogram_array).unwrap()) as _,
3801 ],
3802 )
3803 .unwrap();
3804 let batches = RecordBatches::try_new(schema, vec![batch]).unwrap();
3805
3806 record_batches_to_labels_name(batches, &mut labels).unwrap();
3807
3808 assert_eq!(
3809 labels,
3810 HashSet::from(["host".to_string(), greptime_native_histogram().to_string(),])
3811 );
3812 }
3813
3814 fn prometheus_metric_table_info(table_id: u32, name: &str) -> TableInfo {
3815 let schema = Arc::new(Schema::new(vec![
3816 ColumnSchema::new(
3817 "greptime_timestamp",
3818 ConcreteDataType::timestamp_millisecond_datatype(),
3819 false,
3820 )
3821 .with_time_index(true),
3822 ColumnSchema::new("host", ConcreteDataType::string_datatype(), false),
3823 ColumnSchema::new("value", ConcreteDataType::float64_datatype(), true),
3824 ]));
3825 let mut options = TableOptions::default();
3826 options.extra_options.insert(
3827 LOGICAL_TABLE_METADATA_KEY.to_string(),
3828 "greptime_physical_table".to_string(),
3829 );
3830 options
3831 .extra_options
3832 .insert(SEMANTIC_METRIC_TYPE.to_string(), "counter".to_string());
3833 options
3834 .extra_options
3835 .insert(SEMANTIC_METRIC_UNIT.to_string(), "By".to_string());
3836 let meta = TableMetaBuilder::empty()
3837 .schema(schema)
3838 .primary_key_indices(vec![1])
3839 .engine("metric".to_string())
3840 .next_column_id(3)
3841 .options(options)
3842 .build()
3843 .unwrap();
3844
3845 TableInfoBuilder::default()
3846 .table_id(table_id)
3847 .table_version(0 as TableVersion)
3848 .name(name)
3849 .catalog_name(DEFAULT_CATALOG_NAME)
3850 .schema_name(DEFAULT_SCHEMA_NAME)
3851 .table_type(TableType::Base)
3852 .meta(meta)
3853 .build()
3854 .unwrap()
3855 }
3856
3857 #[tokio::test]
3858 async fn test_retrieve_metric_metadata_uses_semantic_options() {
3859 let table_info = prometheus_metric_table_info(1025, "http_requests_total");
3860 let manager: CatalogManagerRef =
3861 MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&table_info));
3862 let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3863
3864 let metadata = retrieve_metric_metadata(&query_ctx, manager, &MetadataQuery::default())
3865 .await
3866 .unwrap();
3867
3868 assert_eq!(
3869 metadata.get("http_requests_total"),
3870 Some(&vec![PromMetadata {
3871 metric_type: "counter".to_string(),
3872 unit: "bytes".to_string(),
3873 help: String::new(),
3874 }])
3875 );
3876
3877 for (ucum, openmetrics) in [
3878 ("1", "ratios"),
3879 ("ms", "milliseconds"),
3880 ("m/s", "meters_per_second"),
3881 ("USD", "USD"),
3882 ] {
3883 let mut table_info = table_info.clone();
3884 table_info
3885 .meta
3886 .options
3887 .extra_options
3888 .insert(SEMANTIC_METRIC_UNIT.to_string(), ucum.to_string());
3889 assert_eq!(
3890 prometheus_metadata_from_table(&table_info).unit,
3891 openmetrics
3892 );
3893 }
3894
3895 for (metric_type, temporality, expected) in [
3896 ("updown_counter", None, "gauge"),
3897 ("gauge_histogram", None, "gaugehistogram"),
3898 ("summary", None, "summary"),
3899 ("mixed", None, "unknown"),
3900 ("counter", Some("delta"), "unknown"),
3901 ("histogram", Some("delta"), "unknown"),
3902 ("counter", Some("mixed"), "unknown"),
3903 ("histogram", Some("mixed"), "unknown"),
3904 ("updown_counter", Some("mixed"), "unknown"),
3905 ] {
3906 let mut table_info = table_info.clone();
3907 table_info
3908 .meta
3909 .options
3910 .extra_options
3911 .insert(SEMANTIC_METRIC_TYPE.to_string(), metric_type.to_string());
3912 if let Some(temporality) = temporality {
3913 table_info.meta.options.extra_options.insert(
3914 SEMANTIC_METRIC_TEMPORALITY.to_string(),
3915 temporality.to_string(),
3916 );
3917 }
3918 assert_eq!(
3919 prometheus_metadata_from_table(&table_info).metric_type,
3920 expected
3921 );
3922 }
3923 }
3924
3925 #[tokio::test]
3926 async fn test_metadata_query_checks_permissions_and_filters_tables() {
3927 let allowed = prometheus_metric_table_info(1025, "allowed");
3928 let denied = prometheus_metric_table_info(1026, "denied");
3929 let manager = MemoryCatalogManager::new_with_table(EmptyTable::from_table_info(&allowed));
3930 manager
3931 .register_table_sync(RegisterTableRequest {
3932 catalog: DEFAULT_CATALOG_NAME.to_string(),
3933 schema: DEFAULT_SCHEMA_NAME.to_string(),
3934 table_name: denied.name.clone(),
3935 table_id: denied.table_id(),
3936 table: EmptyTable::from_table_info(&denied),
3937 })
3938 .unwrap();
3939 let query_ctx = QueryContext::with(DEFAULT_CATALOG_NAME, DEFAULT_SCHEMA_NAME);
3940
3941 let metadata = retrieve_metric_metadata(
3942 &query_ctx,
3943 manager.clone(),
3944 &MetadataQuery {
3945 metric: Some("allowed".to_string()),
3946 ..Default::default()
3947 },
3948 )
3949 .await
3950 .unwrap();
3951 assert_eq!(metadata.len(), 1);
3952 assert!(metadata.contains_key("allowed"));
3953
3954 let response = metadata_query(
3955 State(Arc::new(TestPrometheusHandler {
3956 catalog_manager: manager.clone(),
3957 deny_operation: true,
3958 denied_table: None,
3959 metric_names: Vec::new(),
3960 label_metric_names: Vec::new(),
3961 label_lookups: Mutex::new(Vec::new()),
3962 queries: Mutex::new(Vec::new()),
3963 ordered_outputs: Mutex::new(Vec::new()),
3964 })),
3965 Query(MetadataQuery::default()),
3966 Extension(query_ctx.clone()),
3967 )
3968 .await;
3969 assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
3970
3971 let handler: PrometheusHandlerRef = Arc::new(TestPrometheusHandler {
3972 catalog_manager: manager,
3973 deny_operation: false,
3974 denied_table: Some("denied"),
3975 metric_names: Vec::new(),
3976 label_metric_names: Vec::new(),
3977 label_lookups: Mutex::new(Vec::new()),
3978 queries: Mutex::new(Vec::new()),
3979 ordered_outputs: Mutex::new(Vec::new()),
3980 });
3981 let response = metadata_query(
3982 State(handler.clone()),
3983 Query(MetadataQuery::default()),
3984 Extension(query_ctx.clone()),
3985 )
3986 .await;
3987 assert!(response.status_code.is_none());
3988 let PrometheusResponse::Metadata(metadata) = response.data else {
3989 panic!("expected metadata response");
3990 };
3991 assert_eq!(
3992 metadata.keys().cloned().collect::<Vec<_>>(),
3993 vec!["allowed".to_string()]
3994 );
3995
3996 let response = metadata_query(
3997 State(handler),
3998 Query(MetadataQuery {
3999 metric: Some("denied".to_string()),
4000 ..Default::default()
4001 }),
4002 Extension(query_ctx),
4003 )
4004 .await;
4005 assert_eq!(Some(StatusCode::PermissionDenied), response.status_code);
4006 }
4007}