query/promql/
label_values.rs1use std::time::{SystemTime, UNIX_EPOCH};
16
17use common_time::Timestamp;
18use common_time::timestamp::TimeUnit;
19use datafusion_common::{Column, ScalarValue};
20use datafusion_expr::utils::conjunction;
21use datafusion_expr::{Expr, LogicalPlan, LogicalPlanBuilder, col};
22use snafu::{OptionExt, ResultExt};
23use table::TableRef;
24
25use crate::promql::error::{
26 DataFusionPlanningSnafu, Result, SystemTimeOutOfRangeSnafu, TimeIndexNotFoundSnafu,
27 TimestampOutOfRangeSnafu,
28};
29
30fn millis_since_epoch(time: SystemTime) -> Result<Timestamp> {
36 let (millis, before_epoch) = match time.duration_since(UNIX_EPOCH) {
37 Ok(duration) => (duration.as_millis(), false),
38 Err(earlier) => (earlier.duration().as_millis(), true),
39 };
40 let millis =
41 signed_millis(millis, before_epoch).with_context(|| SystemTimeOutOfRangeSnafu { time })?;
42 Ok(Timestamp::new_millisecond(millis))
43}
44
45fn signed_millis(millis: u128, before_epoch: bool) -> Option<i64> {
47 let millis = i64::try_from(millis).ok()?;
48 Some(if before_epoch { -millis } else { millis })
49}
50
51fn build_time_filter(time_index_expr: Expr, start: Timestamp, end: Timestamp) -> Expr {
52 time_index_expr
53 .clone()
54 .gt_eq(Expr::Literal(timestamp_to_scalar_value(start), None))
55 .and(time_index_expr.lt_eq(Expr::Literal(timestamp_to_scalar_value(end), None)))
56}
57
58fn timestamp_to_scalar_value(timestamp: Timestamp) -> ScalarValue {
59 let value = timestamp.value();
60 match timestamp.unit() {
61 TimeUnit::Second => ScalarValue::TimestampSecond(Some(value), None),
62 TimeUnit::Millisecond => ScalarValue::TimestampMillisecond(Some(value), None),
63 TimeUnit::Microsecond => ScalarValue::TimestampMicrosecond(Some(value), None),
64 TimeUnit::Nanosecond => ScalarValue::TimestampNanosecond(Some(value), None),
65 }
66}
67
68pub fn rewrite_label_values_query(
70 table: TableRef,
71 scan_plan: LogicalPlan,
72 mut conditions: Vec<Expr>,
73 label_name: String,
74 start: SystemTime,
75 end: SystemTime,
76) -> Result<LogicalPlan> {
77 let schema = table.schema();
78 let ts_column = schema
79 .timestamp_column()
80 .with_context(|| TimeIndexNotFoundSnafu {
81 table: table.table_info().full_table_name(),
82 })?;
83 let unit = ts_column
84 .data_type
85 .as_timestamp()
86 .map(|data_type| data_type.unit())
87 .with_context(|| TimeIndexNotFoundSnafu {
88 table: table.table_info().full_table_name(),
89 })?;
90
91 let start = millis_since_epoch(start)?;
93 let start = start.convert_to(unit).context(TimestampOutOfRangeSnafu {
94 timestamp: start.value(),
95 unit,
96 })?;
97 let end = millis_since_epoch(end)?;
98 let end = end.convert_to(unit).context(TimestampOutOfRangeSnafu {
99 timestamp: end.value(),
100 unit,
101 })?;
102 let time_index_expr = col(Column::from_name(ts_column.name.clone()));
103
104 conditions.push(build_time_filter(time_index_expr, start, end));
105 let filter = conjunction(conditions).unwrap();
107
108 let logical_plan = LogicalPlanBuilder::from(scan_plan)
110 .filter(filter)
111 .context(DataFusionPlanningSnafu)?
112 .project(vec![col(Column::from_name(label_name))])
113 .context(DataFusionPlanningSnafu)?
114 .distinct()
115 .context(DataFusionPlanningSnafu)?
116 .build()
117 .context(DataFusionPlanningSnafu)?;
118
119 Ok(logical_plan)
120}
121
122#[cfg(test)]
123mod tests {
124 use std::time::Duration;
125
126 use super::*;
127
128 #[test]
129 fn millis_before_the_epoch_are_negative() {
130 let time = UNIX_EPOCH - Duration::from_millis(1);
133 assert_eq!(millis_since_epoch(time).unwrap().value(), -1);
134 assert_eq!(millis_since_epoch(UNIX_EPOCH).unwrap().value(), 0);
135
136 let time = UNIX_EPOCH + Duration::from_millis(1);
137 assert_eq!(millis_since_epoch(time).unwrap().value(), 1);
138 }
139
140 #[test]
141 fn millis_beyond_i64_are_rejected() {
142 let limit = i64::MAX as u128;
143 assert_eq!(signed_millis(limit, false), Some(i64::MAX));
144 assert_eq!(signed_millis(limit, true), Some(-i64::MAX));
145 assert_eq!(signed_millis(limit + 1, false), None);
146 assert_eq!(signed_millis(limit + 1, true), None);
147
148 if let Some(time) = UNIX_EPOCH.checked_add(Duration::from_millis(i64::MAX as u64 + 1)) {
151 assert!(millis_since_epoch(time).is_err());
152 }
153 }
154}