1mod absent;
16mod empty_metric;
17mod histogram_fold;
18mod instant_manipulate;
19mod normalize;
20mod planner;
21mod range_manipulate;
22mod scalar_calculate;
23mod series_divide;
24#[cfg(test)]
25mod test_util;
26mod union_distinct_on;
27
28pub use absent::{Absent, AbsentExec, AbsentStream};
29use common_query::native_histogram::{SUM_FIELD, native_histogram_value_type};
30use common_query::prometheus::is_prometheus_stale_nan;
31use datafusion::arrow::array::{Array, Float64Array, StructArray};
32use datafusion::arrow::datatypes::{
33 ArrowPrimitiveType, DataType, TimeUnit, TimestampMillisecondType,
34};
35use datafusion::common::{Column, DFSchemaRef};
36use datafusion::error::{DataFusionError, Result as DataFusionResult};
37use datafusion::logical_expr::{Expr, Extension, LogicalPlan};
38use datatypes::data_type::DataType as _;
39pub use empty_metric::{EmptyMetric, EmptyMetricExec, EmptyMetricStream, build_special_time_expr};
40pub use histogram_fold::{
41 HistogramFold, HistogramFoldExec, HistogramFoldOperation, HistogramFoldStream,
42};
43pub use instant_manipulate::{InstantManipulate, InstantManipulateExec, InstantManipulateStream};
44pub use normalize::{SeriesNormalize, SeriesNormalizeExec, SeriesNormalizeStream};
45pub use planner::PromExtensionPlanner;
46pub use range_manipulate::{RangeManipulate, RangeManipulateExec, RangeManipulateStream};
47pub use scalar_calculate::ScalarCalculate;
48pub use series_divide::{SeriesDivide, SeriesDivideExec, SeriesDivideStream};
49pub use union_distinct_on::{UnionDistinctOn, UnionDistinctOnExec, UnionDistinctOnStream};
50
51pub type Millisecond = <TimestampMillisecondType as ArrowPrimitiveType>::Native;
52
53pub(crate) fn timestamp_unit(data_type: &DataType) -> datafusion::error::Result<TimeUnit> {
54 match data_type {
55 DataType::Timestamp(unit, _) => Ok(*unit),
56 _ => Err(datafusion::error::DataFusionError::Execution(
57 "Time index column is not a timestamp".into(),
58 )),
59 }
60}
61
62pub(crate) fn nanoseconds_per_native_tick(unit: TimeUnit) -> i128 {
63 match unit {
64 TimeUnit::Second => 1_000_000_000,
65 TimeUnit::Millisecond => 1_000_000,
66 TimeUnit::Microsecond => 1_000,
67 TimeUnit::Nanosecond => 1,
68 }
69}
70
71pub(crate) fn local_offset(plan: &LogicalPlan, time_index: &str) -> Millisecond {
77 let Some(index) = plan.schema().index_of_column_by_name(None, time_index) else {
78 return 0;
79 };
80 let (qualifier, field) = plan.schema().qualified_field(index);
81 let mut time_index = Column::new(qualifier.cloned(), field.name().clone());
82 let mut plan = plan;
83
84 loop {
85 match plan {
86 LogicalPlan::Extension(Extension { node }) => {
87 return node
88 .as_any()
89 .downcast_ref::<SeriesNormalize>()
90 .and_then(|normalize| normalize.offset_for_time_index(&time_index))
91 .unwrap_or_default();
92 }
93 LogicalPlan::Projection(projection) => {
94 let Some(output_index) = projection.schema.maybe_index_of_column(&time_index)
95 else {
96 return 0;
97 };
98 let expr = &projection.expr[output_index];
99 let source = match expr {
100 Expr::Column(column) => column,
101 Expr::Alias(alias) => {
102 let Expr::Column(column) = alias.expr.as_ref() else {
103 return 0;
104 };
105 if alias.name != column.name {
106 return 0;
107 }
108 column
109 }
110 _ => return 0,
111 };
112 let Some(input_index) = projection.input.schema().maybe_index_of_column(source)
113 else {
114 return 0;
115 };
116 let (qualifier, field) = projection.input.schema().qualified_field(input_index);
117 time_index = Column::new(qualifier.cloned(), field.name().clone());
118 plan = projection.input.as_ref();
119 }
120 _ => return 0,
121 }
122 }
123}
124
125const METRIC_NUM_SERIES: &str = "num_series";
126
127fn prometheus_stale_sample_column(column: &dyn Array) -> Option<(&dyn Array, &Float64Array)> {
128 let values = if let Some(values) = column.as_any().downcast_ref::<Float64Array>() {
129 values
130 } else {
131 let histograms = column.as_any().downcast_ref::<StructArray>()?;
132 if histograms.data_type() != &native_histogram_value_type().as_arrow_type() {
133 return None;
134 }
135 histograms
136 .column_by_name(SUM_FIELD)?
137 .as_any()
138 .downcast_ref::<Float64Array>()?
139 };
140 Some((column, values))
141}
142
143fn is_prometheus_stale_sample((column, values): (&dyn Array, &Float64Array), row: usize) -> bool {
144 column.is_valid(row) && values.is_valid(row) && is_prometheus_stale_nan(values.value(row))
145}
146
147pub fn serialize_column_index(schema: &DFSchemaRef, column_name: &str) -> u64 {
150 schema
151 .index_of_column_by_name(None, column_name)
152 .map(|idx| idx as u64)
153 .unwrap_or(u64::MAX) }
155
156pub fn resolve_column_name(
158 index: u64,
159 schema: &DFSchemaRef,
160 context: &str,
161 column_type: &str,
162) -> DataFusionResult<String> {
163 let columns = schema.columns();
164 columns
165 .get(index as usize)
166 .ok_or_else(|| {
167 DataFusionError::Internal(format!(
168 "Failed to get {} column at idx {} during unfixing {} with columns:{:?}",
169 column_type, index, context, columns
170 ))
171 })
172 .map(|field| field.name().to_string())
173}
174
175pub fn resolve_column_names(
177 indices: &[u64],
178 schema: &DFSchemaRef,
179 context: &str,
180 column_type: &str,
181) -> DataFusionResult<Vec<String>> {
182 indices
183 .iter()
184 .map(|idx| resolve_column_name(*idx, schema, context, column_type))
185 .collect()
186}
187
188#[cfg(test)]
189mod tests {
190 use std::sync::Arc;
191
192 use datafusion::arrow::datatypes::{DataType, Field, Schema, TimeUnit};
193 use datafusion::common::ToDFSchema;
194 use datafusion::logical_expr::{EmptyRelation, Extension, LogicalPlan, Projection};
195 use datafusion_expr::col;
196
197 use super::*;
198
199 fn input() -> LogicalPlan {
200 LogicalPlan::EmptyRelation(EmptyRelation {
201 produce_one_row: false,
202 schema: Arc::new(Schema::new(vec![
203 Field::new(
204 "timestamp",
205 DataType::Timestamp(TimeUnit::Millisecond, None),
206 false,
207 ),
208 Field::new(
209 "other_ts",
210 DataType::Timestamp(TimeUnit::Millisecond, None),
211 false,
212 ),
213 Field::new("value", DataType::Float64, true),
214 ]))
215 .to_dfschema_ref()
216 .unwrap(),
217 })
218 }
219
220 fn normalized() -> LogicalPlan {
221 LogicalPlan::Extension(Extension {
222 node: Arc::new(SeriesNormalize::new(
223 1_000,
224 "timestamp",
225 false,
226 Vec::new(),
227 input(),
228 )),
229 })
230 }
231
232 #[test]
233 fn local_offset_tracks_identity_preserving_projections() {
234 let projection =
235 Projection::try_new(vec![col("timestamp"), col("value")], Arc::new(normalized()))
236 .unwrap();
237 let projection = Projection::try_new(
238 vec![col("timestamp").alias("timestamp"), col("value")],
239 Arc::new(LogicalPlan::Projection(projection)),
240 )
241 .unwrap();
242
243 assert_eq!(
244 1_000,
245 local_offset(&LogicalPlan::Projection(projection), "timestamp")
246 );
247 }
248
249 #[test]
250 fn local_offset_rejects_a_different_timestamp_or_manipulator() {
251 let renamed = Projection::try_new(
252 vec![col("other_ts").alias("timestamp"), col("value")],
253 Arc::new(normalized()),
254 )
255 .unwrap();
256 assert_eq!(
257 0,
258 local_offset(&LogicalPlan::Projection(renamed), "timestamp")
259 );
260
261 let divide = LogicalPlan::Extension(Extension {
262 node: Arc::new(SeriesDivide::new(
263 Vec::new(),
264 "timestamp".to_string(),
265 normalized(),
266 )),
267 });
268 assert_eq!(0, local_offset(÷, "timestamp"));
269 }
270}