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