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