Skip to main content

promql/extension_plan/
range_manipulate.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashSet;
16use std::pin::Pin;
17use std::sync::Arc;
18use std::task::{Context, Poll};
19
20use common_telemetry::{debug, warn};
21use datafusion::arrow::array::{Array, ArrayRef, Int64Array, TimestampMillisecondArray};
22use datafusion::arrow::compute;
23use datafusion::arrow::datatypes::{DataType, Field, SchemaRef, TimeUnit};
24use datafusion::arrow::error::ArrowError;
25use datafusion::arrow::record_batch::RecordBatch;
26use datafusion::common::stats::Precision;
27use datafusion::common::tree_node::TreeNodeRecursion;
28use datafusion::common::{DFSchema, DFSchemaRef, TableReference};
29use datafusion::error::{DataFusionError, Result as DataFusionResult};
30use datafusion::execution::context::TaskContext;
31use datafusion::logical_expr::{EmptyRelation, Expr, LogicalPlan, UserDefinedLogicalNodeCore};
32use datafusion::physical_expr::EquivalenceProperties;
33use datafusion::physical_plan::metrics::{
34    BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricBuilder, MetricValue, MetricsSet,
35};
36use datafusion::physical_plan::{
37    ChildStats, DisplayAs, DisplayFormatType, Distribution, ExecutionPlan,
38    InputDistributionRequirements, PhysicalExpr, PlanProperties, RecordBatchStream,
39    SendableRecordBatchStream, Statistics, StatisticsArgs,
40};
41use datafusion_expr::ident;
42use datatypes::timestamp::timestamp_array_to_primitive;
43use futures::{Stream, StreamExt, ready};
44use greptime_proto::substrait_extension as pb;
45use prost::Message;
46use snafu::ResultExt;
47
48use crate::error::{DeserializeSnafu, Result};
49use crate::extension_plan::{
50    METRIC_NUM_SERIES, Millisecond, local_offset, nanoseconds_per_native_tick, resolve_column_name,
51    serialize_column_index, timestamp_unit,
52};
53use crate::metrics::PROMQL_SERIES_COUNT;
54use crate::range_array::RangeArray;
55
56/// Time series manipulator for range function.
57///
58/// This plan will "fold" time index and value columns into [RangeArray]s, and truncate
59/// other columns to the same length with the "folded" [RangeArray] column.
60///
61/// To pass runtime information to the execution plan (or the range function), This plan
62/// will add those extra columns:
63/// - timestamp range with type [RangeArray], which is the folded timestamp column.
64/// - end of current range with the same type as the timestamp column. (todo)
65#[derive(Debug, PartialEq, Eq, Hash)]
66pub struct RangeManipulate {
67    start: Millisecond,
68    end: Millisecond,
69    interval: Millisecond,
70    range: Millisecond,
71    offset: Millisecond,
72    time_index: String,
73    field_columns: Vec<String>,
74    input: LogicalPlan,
75    output_schema: DFSchemaRef,
76    unfix: Option<UnfixIndices>,
77}
78
79#[derive(Debug, Clone, PartialEq, Eq, Hash)]
80struct UnfixIndices {
81    pub time_index_idx: u64,
82    pub tag_column_indices: Vec<u64>,
83}
84
85impl RangeManipulate {
86    #[allow(clippy::too_many_arguments)]
87    pub fn new(
88        start: Millisecond,
89        end: Millisecond,
90        interval: Millisecond,
91        offset: Millisecond,
92        range: Millisecond,
93        time_index: String,
94        field_columns: Vec<String>,
95        input: LogicalPlan,
96    ) -> DataFusionResult<Self> {
97        let output_schema =
98            Self::calculate_output_schema(input.schema(), &time_index, &field_columns)?;
99        Ok(Self {
100            start,
101            end,
102            interval,
103            range,
104            offset,
105            time_index,
106            field_columns,
107            input,
108            output_schema,
109            unfix: None,
110        })
111    }
112
113    pub const fn name() -> &'static str {
114        "RangeManipulate"
115    }
116
117    pub fn build_timestamp_range_name(time_index: &str) -> String {
118        format!("{time_index}_range")
119    }
120
121    pub fn internal_range_end_col_name() -> String {
122        "__internal_range_end".to_string()
123    }
124
125    fn range_timestamp_name(&self) -> String {
126        Self::build_timestamp_range_name(&self.time_index)
127    }
128
129    fn calculate_output_schema(
130        input_schema: &DFSchemaRef,
131        time_index: &str,
132        field_columns: &[String],
133    ) -> DataFusionResult<DFSchemaRef> {
134        let columns = input_schema.fields();
135        let mut new_columns = Vec::with_capacity(columns.len() + 1);
136        for i in 0..columns.len() {
137            let x = input_schema.qualified_field(i);
138            new_columns.push((x.0.cloned(), x.1.clone()));
139        }
140
141        // process time index column
142        // the raw timestamp field is preserved. And a new timestamp_range field is appended to the last.
143        let Some(ts_col_index) = input_schema.index_of_column_by_name(None, time_index) else {
144            return Err(datafusion::common::field_not_found(
145                None::<TableReference>,
146                time_index,
147                input_schema.as_ref(),
148            ));
149        };
150        let ts_col_field = &columns[ts_col_index];
151        let output_time_field = Arc::new(
152            ts_col_field
153                .as_ref()
154                .clone()
155                .with_data_type(DataType::Timestamp(TimeUnit::Millisecond, None)),
156        );
157        new_columns[ts_col_index] = (
158            input_schema.qualified_field(ts_col_index).0.cloned(),
159            output_time_field.clone(),
160        );
161        let timestamp_range_field = Field::new(
162            Self::build_timestamp_range_name(time_index),
163            RangeArray::convert_field(output_time_field.as_ref())
164                .data_type()
165                .clone(),
166            ts_col_field.is_nullable(),
167        );
168        new_columns.push((None, Arc::new(timestamp_range_field)));
169
170        // process value columns
171        for name in field_columns {
172            let Some(index) = input_schema.index_of_column_by_name(None, name) else {
173                return Err(datafusion::common::field_not_found(
174                    None::<TableReference>,
175                    name,
176                    input_schema.as_ref(),
177                ));
178            };
179            new_columns[index] = (None, Arc::new(RangeArray::convert_field(&columns[index])));
180        }
181
182        Ok(Arc::new(DFSchema::new_with_metadata(
183            new_columns,
184            input_schema.metadata().clone(),
185        )?))
186    }
187
188    pub fn to_execution_plan(&self, exec_input: Arc<dyn ExecutionPlan>) -> Arc<dyn ExecutionPlan> {
189        let output_schema: SchemaRef = self.output_schema.inner().clone();
190        let properties = exec_input.properties();
191        let properties = Arc::new(PlanProperties::new(
192            EquivalenceProperties::new(output_schema.clone()),
193            properties.partitioning.clone(),
194            properties.emission_type,
195            properties.boundedness,
196        ));
197        Arc::new(RangeManipulateExec {
198            offset: self.offset,
199            start: self.start,
200            end: self.end,
201            interval: self.interval,
202            range: self.range,
203            time_index_column: self.time_index.clone(),
204            time_range_column: self.range_timestamp_name(),
205            field_columns: self.field_columns.clone(),
206            input: exec_input,
207            output_schema,
208            metric: ExecutionPlanMetricsSet::new(),
209            properties,
210        })
211    }
212
213    pub fn serialize(&self) -> Vec<u8> {
214        let time_index_idx = serialize_column_index(self.input.schema(), &self.time_index);
215
216        let tag_column_indices = self
217            .field_columns
218            .iter()
219            .map(|name| serialize_column_index(self.input.schema(), name))
220            .collect::<Vec<u64>>();
221
222        pb::RangeManipulate {
223            start: self.start,
224            end: self.end,
225            interval: self.interval,
226            range: self.range,
227            time_index_idx,
228            tag_column_indices,
229            ..Default::default()
230        }
231        .encode_to_vec()
232    }
233
234    pub fn deserialize(bytes: &[u8]) -> Result<Self> {
235        let pb_range_manipulate = pb::RangeManipulate::decode(bytes).context(DeserializeSnafu)?;
236        let empty_schema = Arc::new(DFSchema::empty());
237        let placeholder_plan = LogicalPlan::EmptyRelation(EmptyRelation {
238            produce_one_row: false,
239            schema: empty_schema.clone(),
240        });
241
242        let unfix = UnfixIndices {
243            time_index_idx: pb_range_manipulate.time_index_idx,
244            tag_column_indices: pb_range_manipulate.tag_column_indices.clone(),
245        };
246        debug!("RangeManipulate deserialize unfix: {:?}", unfix);
247
248        // Unlike `Self::new()`, this method doesn't check the input schema as it will fail
249        // because the input schema is empty.
250        // But this is Ok since datafusion guarantees to call `with_exprs_and_inputs` for the
251        // deserialized plan.
252        Ok(Self {
253            start: pb_range_manipulate.start,
254            end: pb_range_manipulate.end,
255            interval: pb_range_manipulate.interval,
256            range: pb_range_manipulate.range,
257            offset: 0,
258            time_index: String::new(),
259            field_columns: Vec::new(),
260            input: placeholder_plan,
261            output_schema: empty_schema,
262            unfix: Some(unfix),
263        })
264    }
265}
266
267impl PartialOrd for RangeManipulate {
268    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
269        // Compare fields in order excluding output_schema
270        match self.start.partial_cmp(&other.start) {
271            Some(core::cmp::Ordering::Equal) => {}
272            ord => return ord,
273        }
274        match self.end.partial_cmp(&other.end) {
275            Some(core::cmp::Ordering::Equal) => {}
276            ord => return ord,
277        }
278        match self.interval.partial_cmp(&other.interval) {
279            Some(core::cmp::Ordering::Equal) => {}
280            ord => return ord,
281        }
282        match self.range.partial_cmp(&other.range) {
283            Some(core::cmp::Ordering::Equal) => {}
284            ord => return ord,
285        }
286        match self.offset.partial_cmp(&other.offset) {
287            Some(core::cmp::Ordering::Equal) => {}
288            ord => return ord,
289        }
290        match self.time_index.partial_cmp(&other.time_index) {
291            Some(core::cmp::Ordering::Equal) => {}
292            ord => return ord,
293        }
294        match self.field_columns.partial_cmp(&other.field_columns) {
295            Some(core::cmp::Ordering::Equal) => {}
296            ord => return ord,
297        }
298        self.input.partial_cmp(&other.input)
299    }
300}
301
302impl UserDefinedLogicalNodeCore for RangeManipulate {
303    fn name(&self) -> &str {
304        Self::name()
305    }
306
307    fn inputs(&self) -> Vec<&LogicalPlan> {
308        vec![&self.input]
309    }
310
311    fn schema(&self) -> &DFSchemaRef {
312        &self.output_schema
313    }
314
315    fn expressions(&self) -> Vec<Expr> {
316        if self.unfix.is_some() {
317            return vec![];
318        }
319
320        let mut exprs = Vec::with_capacity(1 + self.field_columns.len());
321        exprs.push(ident(&self.time_index));
322        exprs.extend(self.field_columns.iter().map(ident));
323        exprs
324    }
325
326    fn necessary_children_exprs(&self, output_columns: &[usize]) -> Option<Vec<Vec<usize>>> {
327        if self.unfix.is_some() {
328            return None;
329        }
330
331        let input_schema = self.input.schema();
332        let input_len = input_schema.fields().len();
333        let time_index_idx = input_schema.index_of_column_by_name(None, &self.time_index)?;
334
335        if output_columns.is_empty() {
336            let indices = (0..input_len).collect::<Vec<_>>();
337            return Some(vec![indices]);
338        }
339
340        let mut required = Vec::with_capacity(output_columns.len() + 1 + self.field_columns.len());
341        required.push(time_index_idx);
342        for value_column in &self.field_columns {
343            required.push(input_schema.index_of_column_by_name(None, value_column)?);
344        }
345        for &idx in output_columns {
346            if idx < input_len {
347                required.push(idx);
348            } else if idx == input_len {
349                // Derived timestamp range column.
350                required.push(time_index_idx);
351            } else {
352                warn!(
353                    "Output column index {} is out of bounds for input schema with length {}",
354                    idx, input_len
355                );
356                return None;
357            }
358        }
359
360        required.sort_unstable();
361        required.dedup();
362        Some(vec![required])
363    }
364
365    fn fmt_for_explain(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
366        write!(
367            f,
368            "PromRangeManipulate: req range=[{}..{}], interval=[{}], eval range=[{}], time index=[{}], values={:?}",
369            self.start, self.end, self.interval, self.range, self.time_index, self.field_columns
370        )
371    }
372
373    fn with_exprs_and_inputs(
374        &self,
375        _exprs: Vec<Expr>,
376        mut inputs: Vec<LogicalPlan>,
377    ) -> DataFusionResult<Self> {
378        if inputs.len() != 1 {
379            return Err(DataFusionError::Internal(
380                "RangeManipulate should have at exact one input".to_string(),
381            ));
382        }
383
384        let input: LogicalPlan = inputs.pop().unwrap();
385        let input_schema = input.schema();
386
387        if let Some(unfix) = &self.unfix {
388            // transform indices to names
389            let time_index = resolve_column_name(
390                unfix.time_index_idx,
391                input_schema,
392                "RangeManipulate",
393                "time index",
394            )?;
395
396            let field_columns = unfix
397                .tag_column_indices
398                .iter()
399                .map(|idx| resolve_column_name(*idx, input_schema, "RangeManipulate", "tag"))
400                .collect::<DataFusionResult<Vec<String>>>()?;
401
402            let output_schema =
403                Self::calculate_output_schema(input_schema, &time_index, &field_columns)?;
404
405            Ok(Self {
406                start: self.start,
407                end: self.end,
408                interval: self.interval,
409                range: self.range,
410                offset: local_offset(&input, &time_index),
411                time_index,
412                field_columns,
413                input,
414                output_schema,
415                unfix: None,
416            })
417        } else {
418            let output_schema =
419                Self::calculate_output_schema(input_schema, &self.time_index, &self.field_columns)?;
420
421            Ok(Self {
422                start: self.start,
423                end: self.end,
424                interval: self.interval,
425                range: self.range,
426                offset: self.offset,
427                time_index: self.time_index.clone(),
428                field_columns: self.field_columns.clone(),
429                input,
430                output_schema,
431                unfix: None,
432            })
433        }
434    }
435}
436
437#[derive(Debug)]
438pub struct RangeManipulateExec {
439    offset: Millisecond,
440    start: Millisecond,
441    end: Millisecond,
442    interval: Millisecond,
443    range: Millisecond,
444    time_index_column: String,
445    time_range_column: String,
446    field_columns: Vec<String>,
447
448    input: Arc<dyn ExecutionPlan>,
449    output_schema: SchemaRef,
450    metric: ExecutionPlanMetricsSet,
451    properties: Arc<PlanProperties>,
452}
453
454impl ExecutionPlan for RangeManipulateExec {
455    fn apply_expressions(
456        &self,
457        _f: &mut dyn FnMut(&Arc<dyn PhysicalExpr>) -> datafusion_common::Result<TreeNodeRecursion>,
458    ) -> DataFusionResult<TreeNodeRecursion> {
459        Ok(TreeNodeRecursion::Continue)
460    }
461
462    fn schema(&self) -> SchemaRef {
463        self.output_schema.clone()
464    }
465
466    fn properties(&self) -> &Arc<PlanProperties> {
467        &self.properties
468    }
469
470    fn maintains_input_order(&self) -> Vec<bool> {
471        vec![true; self.children().len()]
472    }
473
474    fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
475        vec![&self.input]
476    }
477
478    fn input_distribution_requirements(&self) -> InputDistributionRequirements {
479        let input_requirement = self
480            .input
481            .input_distribution_requirements()
482            .into_per_child();
483        if input_requirement.is_empty() {
484            // if the input is EmptyMetric, its required_input_distribution() is empty so we can't
485            // use its input distribution.
486            InputDistributionRequirements::new(vec![Distribution::UnspecifiedDistribution])
487        } else {
488            InputDistributionRequirements::new(input_requirement)
489        }
490    }
491
492    fn with_new_children(
493        self: Arc<Self>,
494        children: Vec<Arc<dyn ExecutionPlan>>,
495    ) -> DataFusionResult<Arc<dyn ExecutionPlan>> {
496        assert!(!children.is_empty());
497        let exec_input = children[0].clone();
498        let properties = exec_input.properties();
499        let properties = Arc::new(PlanProperties::new(
500            EquivalenceProperties::new(self.output_schema.clone()),
501            properties.partitioning.clone(),
502            properties.emission_type,
503            properties.boundedness,
504        ));
505        Ok(Arc::new(Self {
506            offset: self.offset,
507            start: self.start,
508            end: self.end,
509            interval: self.interval,
510            range: self.range,
511            time_index_column: self.time_index_column.clone(),
512            time_range_column: self.time_range_column.clone(),
513            field_columns: self.field_columns.clone(),
514            output_schema: self.output_schema.clone(),
515            input: children[0].clone(),
516            metric: self.metric.clone(),
517            properties,
518        }))
519    }
520
521    fn execute(
522        &self,
523        partition: usize,
524        context: Arc<TaskContext>,
525    ) -> DataFusionResult<SendableRecordBatchStream> {
526        let baseline_metric = BaselineMetrics::new(&self.metric, partition);
527        let metrics_builder = MetricBuilder::new(&self.metric);
528        let num_series = Count::new();
529        metrics_builder
530            .with_partition(partition)
531            .build(MetricValue::Count {
532                name: METRIC_NUM_SERIES.into(),
533                count: num_series.clone(),
534            });
535
536        let input = self.input.execute(partition, context)?;
537        let schema = input.schema();
538        let time_index = schema
539            .column_with_name(&self.time_index_column)
540            .unwrap_or_else(|| panic!("time index column {} not found", self.time_index_column))
541            .0;
542        let field_columns = self
543            .field_columns
544            .iter()
545            .map(|value_col| {
546                schema
547                    .column_with_name(value_col)
548                    .unwrap_or_else(|| panic!("value column {value_col} not found",))
549                    .0
550            })
551            .collect();
552        let time_unit = timestamp_unit(schema.field(time_index).data_type())?;
553        let aligned_ts_array =
554            RangeManipulateStream::build_aligned_ts_array(self.start, self.end, self.interval);
555        Ok(Box::pin(RangeManipulateStream {
556            offset: self.offset,
557            start: self.start,
558            end: self.end,
559            interval: self.interval,
560            range: self.range,
561            time_index,
562            time_unit,
563            field_columns,
564            aligned_ts_array,
565            output_schema: self.output_schema.clone(),
566            input,
567            metric: baseline_metric,
568            num_series,
569        }))
570    }
571
572    fn metrics(&self) -> Option<MetricsSet> {
573        Some(self.metric.clone_inner())
574    }
575
576    fn child_stats_requests(&self, partition: Option<usize>) -> Vec<ChildStats> {
577        vec![ChildStats::At(partition)]
578    }
579
580    fn statistics_from_inputs(
581        &self,
582        input_stats: &[Arc<Statistics>],
583        _args: &StatisticsArgs,
584    ) -> DataFusionResult<Arc<Statistics>> {
585        let input_stats = &input_stats[0];
586
587        let estimated_row_num = (self.end - self.start) as f64 / self.interval as f64;
588        let estimated_total_bytes = input_stats
589            .total_byte_size
590            .get_value()
591            .zip(input_stats.num_rows.get_value())
592            .map(|(size, rows)| {
593                Precision::Inexact(((*size as f64 / *rows as f64) * estimated_row_num).floor() as _)
594            })
595            .unwrap_or_default();
596
597        Ok(Arc::new(Statistics {
598            num_rows: Precision::Inexact(estimated_row_num as _),
599            total_byte_size: estimated_total_bytes,
600            // TODO(ruihang): support this column statistics
601            column_statistics: Statistics::unknown_column(&self.schema()),
602        }))
603    }
604
605    fn name(&self) -> &str {
606        "RangeManipulateExec"
607    }
608}
609
610impl DisplayAs for RangeManipulateExec {
611    fn fmt_as(&self, t: DisplayFormatType, f: &mut std::fmt::Formatter) -> std::fmt::Result {
612        match t {
613            DisplayFormatType::Default
614            | DisplayFormatType::Verbose
615            | DisplayFormatType::TreeRender => {
616                write!(
617                    f,
618                    "PromRangeManipulateExec: req range=[{}..{}], interval=[{}], eval range=[{}], time index=[{}]",
619                    self.start, self.end, self.interval, self.range, self.time_index_column
620                )
621            }
622        }
623    }
624}
625
626pub struct RangeManipulateStream {
627    offset: Millisecond,
628    start: Millisecond,
629    end: Millisecond,
630    interval: Millisecond,
631    range: Millisecond,
632    time_index: usize,
633    time_unit: TimeUnit,
634    field_columns: Vec<usize>,
635    aligned_ts_array: ArrayRef,
636
637    output_schema: SchemaRef,
638    input: SendableRecordBatchStream,
639    metric: BaselineMetrics,
640    /// Number of series processed.
641    num_series: Count,
642}
643
644impl RecordBatchStream for RangeManipulateStream {
645    fn schema(&self) -> SchemaRef {
646        self.output_schema.clone()
647    }
648}
649
650impl Stream for RangeManipulateStream {
651    type Item = DataFusionResult<RecordBatch>;
652
653    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
654        let poll = loop {
655            match ready!(self.input.poll_next_unpin(cx)) {
656                Some(Ok(batch)) => {
657                    let timer = std::time::Instant::now();
658                    let result = self.manipulate(batch);
659                    if let Ok(None) = result {
660                        self.metric.elapsed_compute().add_elapsed(timer);
661                        continue;
662                    } else {
663                        self.num_series.add(1);
664                        self.metric.elapsed_compute().add_elapsed(timer);
665                        break Poll::Ready(result.transpose());
666                    }
667                }
668                None => {
669                    PROMQL_SERIES_COUNT.observe(self.num_series.value() as f64);
670                    break Poll::Ready(None);
671                }
672                Some(Err(e)) => break Poll::Ready(Some(Err(e))),
673            }
674        };
675        self.metric.record_poll(poll)
676    }
677}
678
679impl RangeManipulateStream {
680    // Prometheus: https://github.com/prometheus/prometheus/blob/e934d0f01158a1d55fa0ebb035346b195fcc1260/promql/engine.go#L1113-L1198
681    // But they are not exactly the same, because we don't eager-evaluate on the data in this plan.
682    // And the generated timestamp is not aligned to the step. It's expected to do later.
683    pub fn manipulate(&self, input: RecordBatch) -> DataFusionResult<Option<RecordBatch>> {
684        let mut other_columns = (0..input.columns().len()).collect::<HashSet<_>>();
685        // calculate the range
686        let (ranges, (start, end)) = self.calculate_range(&input)?;
687        // ignore this if all ranges are empty
688        if ranges.iter().all(|(_, len)| *len == 0) {
689            return Ok(None);
690        }
691
692        // transform columns
693        let mut new_columns = input.columns().to_vec();
694        for index in self.field_columns.iter() {
695            let _ = other_columns.remove(index);
696            let column = input.column(*index);
697            let new_column = Arc::new(
698                RangeArray::from_ranges(column.clone(), ranges.clone())
699                    .map_err(|e| ArrowError::InvalidArgumentError(e.to_string()))?
700                    .into_dict(),
701            );
702            new_columns[*index] = new_column;
703        }
704
705        // The timestamp range payload is always millisecond ABI. Shift in wide
706        // native precision before truncating toward zero, preserving null validity.
707        let scale = nanoseconds_per_native_tick(self.time_unit);
708        let (timestamps, _) = timestamp_array_to_primitive(input.column(self.time_index))
709            .ok_or_else(|| {
710                DataFusionError::Execution("Time index column is not a timestamp".into())
711            })?;
712        let timestamp_values = timestamps
713            .values()
714            .iter()
715            .enumerate()
716            .map(|(index, timestamp)| {
717                if !input.column(self.time_index).is_valid(index) {
718                    return Ok(None);
719                }
720                let shifted_ns = (*timestamp as i128) * scale + (self.offset as i128) * 1_000_000;
721                i64::try_from(shifted_ns / 1_000_000)
722                    .map(Some)
723                    .map_err(|_| {
724                        ArrowError::ComputeError(
725                            "RangeManipulate timestamp payload overflow".into(),
726                        )
727                    })
728            })
729            .collect::<std::result::Result<Vec<_>, _>>()?;
730        let timestamp_values = TimestampMillisecondArray::from(timestamp_values);
731        let ts_range_column = RangeArray::from_ranges(Arc::new(timestamp_values), ranges.clone())
732            .map_err(|e| ArrowError::InvalidArgumentError(e.to_string()))?
733            .into_dict();
734        new_columns.push(Arc::new(ts_range_column));
735
736        // truncate other columns
737        let take_indices = Int64Array::from(vec![0; ranges.len()]);
738        for index in other_columns.into_iter() {
739            new_columns[index] = compute::take(&input.column(index), &take_indices, None)?;
740        }
741        // replace timestamp with the aligned one
742        let new_time_index = if ranges.len() != self.aligned_ts_array.len() {
743            Self::build_aligned_ts_array(start, end, self.interval)
744        } else {
745            self.aligned_ts_array.clone()
746        };
747        new_columns[self.time_index] = new_time_index;
748
749        RecordBatch::try_new(self.output_schema.clone(), new_columns)
750            .map(Some)
751            .map_err(|e| DataFusionError::ArrowError(Box::new(e), None))
752    }
753
754    fn build_aligned_ts_array(start: i64, end: i64, interval: i64) -> ArrayRef {
755        Arc::new(TimestampMillisecondArray::from_iter_values(
756            (start..=end).step_by(interval as _),
757        ))
758    }
759
760    /// Return values:
761    /// - A vector of tuples where each tuple contains the start index and length of the range.
762    /// - A tuple of the actual start/end timestamp used to calculate the range.
763    #[allow(clippy::type_complexity)]
764    fn calculate_range(
765        &self,
766        input: &RecordBatch,
767    ) -> DataFusionResult<(Vec<(u32, u32)>, (i64, i64))> {
768        let ts_column = input.column(self.time_index);
769        let scale = nanoseconds_per_native_tick(self.time_unit);
770        let (timestamps, _) = timestamp_array_to_primitive(ts_column).ok_or_else(|| {
771            DataFusionError::Execution("Time index column is not a timestamp".into())
772        })?;
773        let timestamps = timestamps.values();
774        let timestamp =
775            |index| (timestamps[index] as i128) * scale + (self.offset as i128) * 1_000_000;
776        let len = timestamps.len();
777        if len == 0 {
778            return Ok((vec![], (self.start, self.end)));
779        }
780
781        // Shorten the range using wide arithmetic so timestamps near the native
782        // type limits retain every query-aligned evaluation point.
783        let query_start = self.start as i128;
784        let query_end = self.end as i128;
785        let interval = self.interval as i128;
786        let first_ts = timestamp(0).div_euclid(1_000_000);
787        // Preserve the query's alignment pattern when optimizing start time.
788        let remainder = (first_ts - query_start).rem_euclid(interval);
789        let first_ts_aligned = first_ts + (interval - remainder).rem_euclid(interval);
790        let last_ts_with_range =
791            (timestamp(len - 1) + (self.range as i128) * 1_000_000).div_euclid(1_000_000);
792        let remainder = (last_ts_with_range - query_start).rem_euclid(interval);
793        let last_ts_aligned = last_ts_with_range - remainder;
794        let start = query_start.max(first_ts_aligned);
795        let end = query_end.min(last_ts_aligned);
796        if start > end {
797            return Ok((vec![], (self.start, self.end)));
798        }
799        // The intersection is within the declared i64 query bounds.
800        let start = start as i64;
801        let end = end as i64;
802        let mut ranges = Vec::new();
803
804        // Range membership is decided on shifted native ticks, before the
805        // timestamp-range payload is converted to its millisecond ABI. This
806        // keeps sub-millisecond samples distinct in a range; equal millisecond
807        // payload values are not a reason to deduplicate input samples.
808        //
809        // Calculate for every aligned timestamp (`curr_ts`), assuming ordered timestamps.
810        let mut left = 0usize;
811        let mut right = 0usize;
812        for curr_ts in (start..=end).step_by(self.interval as _) {
813            let start_ts = (curr_ts as i128) * 1_000_000 - (self.range as i128) * 1_000_000;
814
815            while left < len && timestamp(left) <= start_ts {
816                left += 1;
817            }
818            right = right.max(left);
819            while right < len && timestamp(right) <= (curr_ts as i128) * 1_000_000 {
820                right += 1;
821            }
822
823            if left == right {
824                ranges.push((0, 0));
825            } else {
826                ranges.push((left as _, (right - left) as _));
827            }
828        }
829
830        Ok((ranges, (start, end)))
831    }
832}
833
834#[cfg(test)]
835mod test {
836    use datafusion::arrow::array::{
837        ArrayRef, DictionaryArray, Float64Array, StringArray, TimestampMicrosecondArray,
838        TimestampNanosecondArray, TimestampSecondArray,
839    };
840    use datafusion::arrow::buffer::NullBuffer;
841    use datafusion::arrow::datatypes::{
842        ArrowPrimitiveType, DataType, Field, Int64Type, Schema, TimestampMillisecondType,
843    };
844    use datafusion::common::ToDFSchema;
845    use datafusion::datasource::memory::MemorySourceConfig;
846    use datafusion::datasource::source::DataSourceExec;
847    use datafusion::logical_expr::{
848        EmptyRelation, Extension, LogicalPlan, UserDefinedLogicalNodeCore,
849    };
850    use datafusion::physical_expr::Partitioning;
851    use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
852    use datafusion::physical_plan::memory::MemoryStream;
853    use datafusion::physical_plan::{ChildrenPropertiesMode, ReplaceChildrenOptions};
854    use datafusion::prelude::SessionContext;
855    use datatypes::arrow::array::TimestampMillisecondArray;
856    use futures::FutureExt;
857
858    use super::*;
859
860    const TIME_INDEX_COLUMN: &str = "timestamp";
861
862    fn project_batch(batch: &RecordBatch, indices: &[usize]) -> RecordBatch {
863        let fields = indices
864            .iter()
865            .map(|&idx| batch.schema().field(idx).clone())
866            .collect::<Vec<_>>();
867        let columns = indices
868            .iter()
869            .map(|&idx| batch.column(idx).clone())
870            .collect::<Vec<_>>();
871        let schema = Arc::new(Schema::new(fields));
872        RecordBatch::try_new(schema, columns).unwrap()
873    }
874
875    fn prepare_test_data() -> DataSourceExec {
876        let schema = Arc::new(Schema::new(vec![
877            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
878            Field::new("value_1", DataType::Float64, true),
879            Field::new("value_2", DataType::Float64, true),
880            Field::new("path", DataType::Utf8, true),
881        ]));
882        let timestamp_column = Arc::new(TimestampMillisecondArray::from(vec![
883            0, 30_000, 60_000, 90_000, 120_000, // every 30s
884            180_000, 240_000, // every 60s
885            241_000, 271_000, 291_000, // others
886        ])) as _;
887        let field_column: ArrayRef = Arc::new(Float64Array::from(vec![1.0; 10])) as _;
888        let path_column = Arc::new(StringArray::from(vec!["foo"; 10])) as _;
889        let data = RecordBatch::try_new(
890            schema.clone(),
891            vec![
892                timestamp_column,
893                field_column.clone(),
894                field_column,
895                path_column,
896            ],
897        )
898        .unwrap();
899
900        DataSourceExec::new(Arc::new(
901            MemorySourceConfig::try_new(&[vec![data]], schema, None).unwrap(),
902        ))
903    }
904
905    async fn do_normalize_test(
906        start: Millisecond,
907        end: Millisecond,
908        interval: Millisecond,
909        range: Millisecond,
910        expected: String,
911    ) {
912        let memory_exec = Arc::new(prepare_test_data());
913        let time_index = TIME_INDEX_COLUMN.to_string();
914        let field_columns = vec!["value_1".to_string(), "value_2".to_string()];
915        let manipulate_output_schema = SchemaRef::new(
916            RangeManipulate::calculate_output_schema(
917                &memory_exec.schema().to_dfschema_ref().unwrap(),
918                &time_index,
919                &field_columns,
920            )
921            .unwrap()
922            .as_arrow()
923            .clone(),
924        );
925        let properties = Arc::new(PlanProperties::new(
926            EquivalenceProperties::new(manipulate_output_schema.clone()),
927            Partitioning::UnknownPartitioning(1),
928            EmissionType::Incremental,
929            Boundedness::Bounded,
930        ));
931        let normalize_exec = Arc::new(RangeManipulateExec {
932            offset: 0,
933            start,
934            end,
935            interval,
936            range,
937            field_columns,
938            output_schema: manipulate_output_schema,
939            time_range_column: RangeManipulate::build_timestamp_range_name(&time_index),
940            time_index_column: time_index,
941            input: memory_exec,
942            metric: ExecutionPlanMetricsSet::new(),
943            properties,
944        });
945        let session_context = SessionContext::default();
946        let result = datafusion::physical_plan::collect(normalize_exec, session_context.task_ctx())
947            .await
948            .unwrap();
949        // DirectoryArray from RangeArray cannot be print as normal arrays.
950        let result_literal: String = result
951            .into_iter()
952            .filter_map(|batch| {
953                batch
954                    .columns()
955                    .iter()
956                    .map(|array| {
957                        if matches!(array.data_type(), &DataType::Dictionary(..)) {
958                            let dict_array = array
959                                .as_any()
960                                .downcast_ref::<DictionaryArray<Int64Type>>()
961                                .unwrap()
962                                .clone();
963                            format!("{:?}", RangeArray::try_new(dict_array).unwrap())
964                        } else {
965                            format!("{array:?}")
966                        }
967                    })
968                    .reduce(|lhs, rhs| lhs + "\n" + &rhs)
969            })
970            .reduce(|lhs, rhs| lhs + "\n\n" + &rhs)
971            .unwrap();
972
973        assert_eq!(result_literal, expected);
974    }
975
976    #[tokio::test]
977    async fn native_timestamps_preserve_range_membership_and_ms_payload() {
978        for (unit, ticks_per_ms) in [
979            (TimeUnit::Microsecond, 1_000_i64),
980            (TimeUnit::Nanosecond, 1_000_000_i64),
981        ] {
982            let lower = 1_000 * ticks_per_ms;
983            let upper = 1_001 * ticks_per_ms;
984            // Exclude the lower boundary and future sample; retain both native
985            // samples in the same millisecond bucket and the exact upper sample.
986            let timestamps = vec![lower, lower + 1, lower + 2, upper, upper + 1];
987            let time: ArrayRef = match unit {
988                TimeUnit::Microsecond => Arc::new(TimestampMicrosecondArray::from(timestamps)),
989                TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(timestamps)),
990                _ => unreachable!(),
991            };
992            let schema = Arc::new(Schema::new(vec![
993                Field::new(TIME_INDEX_COLUMN, DataType::Timestamp(unit, None), false),
994                Field::new("value", DataType::Float64, true),
995            ]));
996            let batch = RecordBatch::try_new(
997                schema.clone(),
998                vec![
999                    time,
1000                    Arc::new(Float64Array::from(vec![10.0, 20.0, 30.0, 40.0, 50.0])),
1001                ],
1002            )
1003            .unwrap();
1004            let logical_input = LogicalPlan::EmptyRelation(EmptyRelation {
1005                produce_one_row: false,
1006                schema: schema.clone().to_dfschema_ref().unwrap(),
1007            });
1008            let plan = RangeManipulate::new(
1009                1_001,
1010                1_001,
1011                1,
1012                0,
1013                1,
1014                TIME_INDEX_COLUMN.to_string(),
1015                vec!["value".to_string()],
1016                logical_input.clone(),
1017            )
1018            .unwrap();
1019            let output_time = Field::new(
1020                TIME_INDEX_COLUMN,
1021                DataType::Timestamp(TimeUnit::Millisecond, None),
1022                false,
1023            );
1024            let output_schema = Arc::new(Schema::new(vec![
1025                output_time.clone(),
1026                RangeArray::convert_field(&Field::new("value", DataType::Float64, true)),
1027                Field::new(
1028                    RangeManipulate::build_timestamp_range_name(TIME_INDEX_COLUMN),
1029                    RangeArray::convert_field(&output_time).data_type().clone(),
1030                    false,
1031                ),
1032            ]));
1033            assert_eq!(plan.schema().as_arrow(), output_schema.as_ref());
1034
1035            let rebuilt = RangeManipulate::deserialize(&plan.serialize())
1036                .unwrap()
1037                .with_exprs_and_inputs(vec![], vec![logical_input])
1038                .unwrap();
1039            assert_eq!(rebuilt.schema(), plan.schema());
1040            assert_eq!(rebuilt.input.schema().as_arrow(), schema.as_ref());
1041
1042            let input = Arc::new(DataSourceExec::new(Arc::new(
1043                MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(),
1044            )));
1045            let exec = rebuilt.to_execution_plan(input);
1046            assert_eq!(exec.schema(), output_schema);
1047            assert_eq!(exec.children()[0].schema(), schema);
1048
1049            let batches =
1050                datafusion::physical_plan::collect(exec, SessionContext::default().task_ctx())
1051                    .await
1052                    .unwrap();
1053            assert_eq!(batches.len(), 1, "{unit:?}");
1054            let output = &batches[0];
1055            assert_eq!(output.schema(), output_schema);
1056            assert_eq!(output.num_rows(), 1);
1057            assert_eq!(
1058                output
1059                    .column(0)
1060                    .as_any()
1061                    .downcast_ref::<TimestampMillisecondArray>()
1062                    .unwrap()
1063                    .values()
1064                    .as_ref(),
1065                &[1_001]
1066            );
1067
1068            // RangeArray packs offset/length into dictionary keys; Arrow dictionary
1069            // equality treats those packed keys as indices and cannot compare them.
1070            let values = RangeArray::try_new(
1071                output
1072                    .column(1)
1073                    .as_any()
1074                    .downcast_ref::<DictionaryArray<Int64Type>>()
1075                    .unwrap()
1076                    .clone(),
1077            )
1078            .unwrap();
1079            assert_eq!(values.get_offset_length(0), Some((1, 3)));
1080            assert_eq!(
1081                values.get(0).unwrap().to_data(),
1082                Float64Array::from(vec![20.0, 30.0, 40.0]).to_data()
1083            );
1084            let timestamps = RangeArray::try_new(
1085                output
1086                    .column(2)
1087                    .as_any()
1088                    .downcast_ref::<DictionaryArray<Int64Type>>()
1089                    .unwrap()
1090                    .clone(),
1091            )
1092            .unwrap();
1093            assert_eq!(timestamps.get_offset_length(0), Some((1, 3)));
1094            assert_eq!(
1095                timestamps.get(0).unwrap().to_data(),
1096                TimestampMillisecondArray::from(vec![1_000, 1_000, 1_001]).to_data()
1097            );
1098        }
1099    }
1100
1101    #[test]
1102    fn logical_offset_participates_in_ordering() {
1103        let input = LogicalPlan::EmptyRelation(EmptyRelation {
1104            produce_one_row: false,
1105            schema: prepare_test_data().schema().to_dfschema_ref().unwrap(),
1106        });
1107        let first = RangeManipulate::new(
1108            0,
1109            0,
1110            0,
1111            0,
1112            0,
1113            TIME_INDEX_COLUMN.to_string(),
1114            vec!["value_1".to_string()],
1115            input.clone(),
1116        )
1117        .unwrap();
1118        let second = RangeManipulate::new(
1119            0,
1120            0,
1121            0,
1122            1,
1123            0,
1124            TIME_INDEX_COLUMN.to_string(),
1125            vec!["value_1".to_string()],
1126            input,
1127        )
1128        .unwrap();
1129        assert_ne!(first, second);
1130        assert_eq!(first.partial_cmp(&second), Some(std::cmp::Ordering::Less));
1131    }
1132
1133    #[tokio::test]
1134    async fn logical_normalize_offset_survives_rebuild_and_executes() {
1135        for (name, time_unit, raw, offset, start, range, expected_payload) in [
1136            (
1137                "millisecond offset",
1138                TimeUnit::Millisecond,
1139                0,
1140                1_000,
1141                1_000,
1142                1_000,
1143                1_000,
1144            ),
1145            (
1146                "negative native lower limit with positive window",
1147                TimeUnit::Nanosecond,
1148                -9_223_112_837_000_000_000,
1149                -259_200_000,
1150                -9_223_372_037_000,
1151                300_000,
1152                -9_223_372_037_000,
1153            ),
1154            (
1155                "second timestamp with negative fractional offset",
1156                TimeUnit::Second,
1157                1,
1158                -500,
1159                1_000,
1160                1_000,
1161                500,
1162            ),
1163        ] {
1164            let schema = Arc::new(Schema::new(vec![
1165                Field::new(
1166                    TIME_INDEX_COLUMN,
1167                    DataType::Timestamp(time_unit, None),
1168                    false,
1169                ),
1170                Field::new("value", DataType::Float64, true),
1171            ]));
1172            let input = LogicalPlan::EmptyRelation(EmptyRelation {
1173                produce_one_row: false,
1174                schema: schema.clone().to_dfschema_ref().unwrap(),
1175            });
1176            let normalize = crate::extension_plan::SeriesNormalize::new(
1177                offset,
1178                TIME_INDEX_COLUMN,
1179                false,
1180                Vec::new(),
1181                input.clone(),
1182            );
1183            let normalize =
1184                crate::extension_plan::SeriesNormalize::deserialize(&normalize.serialize())
1185                    .unwrap()
1186                    .with_exprs_and_inputs(vec![], vec![input.clone()])
1187                    .unwrap();
1188            let normalized = LogicalPlan::Extension(Extension {
1189                node: Arc::new(normalize),
1190            });
1191            let fresh = RangeManipulate::new(
1192                start,
1193                start,
1194                1,
1195                offset,
1196                range,
1197                TIME_INDEX_COLUMN.to_string(),
1198                vec!["value".to_string()],
1199                input.clone(),
1200            )
1201            .unwrap()
1202            .with_exprs_and_inputs(vec![], vec![input.clone()])
1203            .unwrap();
1204            let serialized = RangeManipulate::new(
1205                start,
1206                start,
1207                1,
1208                offset,
1209                range,
1210                TIME_INDEX_COLUMN.to_string(),
1211                vec!["value".to_string()],
1212                normalized.clone(),
1213            )
1214            .unwrap();
1215            let decoded = RangeManipulate::deserialize(&serialized.serialize())
1216                .unwrap()
1217                .with_exprs_and_inputs(vec![], vec![normalized])
1218                .unwrap();
1219            let timestamp: ArrayRef = match time_unit {
1220                TimeUnit::Millisecond => Arc::new(TimestampMillisecondArray::from(vec![raw])),
1221                TimeUnit::Nanosecond => Arc::new(TimestampNanosecondArray::from(vec![raw])),
1222                TimeUnit::Second => Arc::new(TimestampSecondArray::from(vec![raw])),
1223                _ => unreachable!(),
1224            };
1225            let batch = RecordBatch::try_new(
1226                schema.clone(),
1227                vec![timestamp, Arc::new(Float64Array::from(vec![7.0]))],
1228            )
1229            .unwrap();
1230            for (mode, rebuilt) in [("fresh", fresh), ("decoded", decoded)] {
1231                let rebuilt = rebuilt
1232                    .with_exprs_and_inputs(vec![], vec![input.clone()])
1233                    .unwrap();
1234                assert_eq!(rebuilt.offset, offset, "{name}: {mode}");
1235                assert_eq!(rebuilt.input.schema(), input.schema(), "{name}: {mode}");
1236
1237                let empty_exec_input = Arc::new(DataSourceExec::new(Arc::new(
1238                    MemorySourceConfig::try_new(&[vec![]], schema.clone(), None).unwrap(),
1239                )));
1240                let exec_input = Arc::new(DataSourceExec::new(Arc::new(
1241                    MemorySourceConfig::try_new(&[vec![batch.clone()]], schema.clone(), None)
1242                        .unwrap(),
1243                )));
1244                let exec = rebuilt
1245                    .to_execution_plan(empty_exec_input)
1246                    .replace_children(
1247                        vec![exec_input],
1248                        ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute),
1249                    )
1250                    .unwrap();
1251                let output =
1252                    datafusion::physical_plan::collect(exec, SessionContext::default().task_ctx())
1253                        .await
1254                        .unwrap();
1255                assert_eq!(output.len(), 1, "{name}: {mode}");
1256                let output = &output[0];
1257                assert_eq!(output.num_rows(), 1, "{name}: {mode}");
1258                assert_eq!(
1259                    output
1260                        .column(0)
1261                        .as_any()
1262                        .downcast_ref::<TimestampMillisecondArray>()
1263                        .unwrap()
1264                        .value(0),
1265                    start,
1266                    "{name}: {mode}"
1267                );
1268                let values = RangeArray::try_new(
1269                    output
1270                        .column(1)
1271                        .as_any()
1272                        .downcast_ref::<DictionaryArray<Int64Type>>()
1273                        .unwrap()
1274                        .clone(),
1275                )
1276                .unwrap();
1277                assert_eq!(values.get_offset_length(0), Some((0, 1)), "{name}: {mode}");
1278                assert_eq!(
1279                    values.get(0).unwrap().to_data(),
1280                    Float64Array::from(vec![7.0]).to_data(),
1281                    "{name}: {mode}"
1282                );
1283                let timestamps = RangeArray::try_new(
1284                    output
1285                        .column(2)
1286                        .as_any()
1287                        .downcast_ref::<DictionaryArray<Int64Type>>()
1288                        .unwrap()
1289                        .clone(),
1290                )
1291                .unwrap();
1292                assert_eq!(
1293                    timestamps.get_offset_length(0),
1294                    Some((0, 1)),
1295                    "{name}: {mode}"
1296                );
1297                assert_eq!(
1298                    timestamps.get(0).unwrap().to_data(),
1299                    TimestampMillisecondArray::from(vec![expected_payload]).to_data(),
1300                    "{name}: {mode}"
1301                );
1302            }
1303        }
1304    }
1305
1306    #[tokio::test]
1307    async fn range_payload_preserves_null_timestamp_and_rejects_offset_overflow() {
1308        let schema = Arc::new(Schema::new(vec![
1309            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1310            Field::new("value", DataType::Float64, true),
1311        ]));
1312        let null_timestamp =
1313            TimestampMillisecondArray::new(vec![1_000].into(), Some(NullBuffer::from(vec![false])));
1314        let batch = RecordBatch::try_new(
1315            schema.clone(),
1316            vec![
1317                Arc::new(null_timestamp),
1318                Arc::new(Float64Array::from(vec![7.0])),
1319            ],
1320        )
1321        .unwrap();
1322        let input = Arc::new(DataSourceExec::new(Arc::new(
1323            MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(),
1324        )));
1325        let plan = RangeManipulate::new(
1326            1_000,
1327            1_000,
1328            1,
1329            0,
1330            1,
1331            TIME_INDEX_COLUMN.to_string(),
1332            vec!["value".to_string()],
1333            LogicalPlan::EmptyRelation(EmptyRelation {
1334                produce_one_row: false,
1335                schema: schema.clone().to_dfschema_ref().unwrap(),
1336            }),
1337        )
1338        .unwrap();
1339        let output = datafusion::physical_plan::collect(
1340            plan.to_execution_plan(input),
1341            SessionContext::default().task_ctx(),
1342        )
1343        .await
1344        .unwrap();
1345        let timestamps = RangeArray::try_new(
1346            output[0]
1347                .column(2)
1348                .as_any()
1349                .downcast_ref::<DictionaryArray<Int64Type>>()
1350                .unwrap()
1351                .clone(),
1352        )
1353        .unwrap();
1354        let payload = timestamps.get(0).unwrap();
1355        let payload = payload
1356            .as_any()
1357            .downcast_ref::<TimestampMillisecondArray>()
1358            .unwrap();
1359        assert_eq!(payload.len(), 1);
1360        assert!(!payload.is_valid(0));
1361
1362        let batch = RecordBatch::try_new(
1363            schema.clone(),
1364            vec![
1365                Arc::new(TimestampMillisecondArray::from(vec![0, i64::MAX])),
1366                Arc::new(Float64Array::from(vec![7.0, 8.0])),
1367            ],
1368        )
1369        .unwrap();
1370        let input = Arc::new(DataSourceExec::new(Arc::new(
1371            MemorySourceConfig::try_new(&[vec![batch]], schema.clone(), None).unwrap(),
1372        )));
1373        let normalized = crate::extension_plan::SeriesNormalize::new(
1374            1,
1375            TIME_INDEX_COLUMN,
1376            false,
1377            Vec::new(),
1378            LogicalPlan::EmptyRelation(EmptyRelation {
1379                produce_one_row: false,
1380                schema: schema.to_dfschema_ref().unwrap(),
1381            }),
1382        );
1383        let plan = RangeManipulate::new(
1384            1,
1385            1,
1386            1,
1387            1,
1388            1,
1389            TIME_INDEX_COLUMN.to_string(),
1390            vec!["value".to_string()],
1391            LogicalPlan::Extension(Extension {
1392                node: Arc::new(normalized),
1393            }),
1394        )
1395        .unwrap();
1396        let error = datafusion::physical_plan::collect(
1397            plan.to_execution_plan(input),
1398            SessionContext::default().task_ctx(),
1399        )
1400        .await
1401        .unwrap_err();
1402        assert!(error.to_string().contains("timestamp payload overflow"));
1403    }
1404
1405    #[tokio::test]
1406    async fn pruning_should_keep_time_and_value_columns_for_exec() {
1407        let schema = Arc::new(Schema::new(vec![
1408            Field::new(TIME_INDEX_COLUMN, TimestampMillisecondType::DATA_TYPE, true),
1409            Field::new("value_1", DataType::Float64, true),
1410            Field::new("value_2", DataType::Float64, true),
1411            Field::new("path", DataType::Utf8, true),
1412        ]));
1413        let df_schema = schema.clone().to_dfschema_ref().unwrap();
1414        let input = LogicalPlan::EmptyRelation(EmptyRelation {
1415            produce_one_row: false,
1416            schema: df_schema,
1417        });
1418        let plan = RangeManipulate::new(
1419            0,
1420            310_000,
1421            30_000,
1422            0,
1423            90_000,
1424            TIME_INDEX_COLUMN.to_string(),
1425            vec!["value_1".to_string(), "value_2".to_string()],
1426            input,
1427        )
1428        .unwrap();
1429
1430        // Simulate a parent projection requesting only the `path` column.
1431        let output_columns = [3usize];
1432        let required = plan.necessary_children_exprs(&output_columns).unwrap();
1433        let required = &required[0];
1434        assert_eq!(required.as_slice(), &[0, 1, 2, 3]);
1435
1436        let timestamp_column = Arc::new(TimestampMillisecondArray::from(vec![
1437            0, 30_000, 60_000, 90_000, 120_000, // every 30s
1438            180_000, 240_000, // every 60s
1439            241_000, 271_000, 291_000, // others
1440        ])) as _;
1441        let field_column: ArrayRef = Arc::new(Float64Array::from(vec![1.0; 10])) as _;
1442        let path_column = Arc::new(StringArray::from(vec!["foo"; 10])) as _;
1443        let input_batch = RecordBatch::try_new(
1444            schema,
1445            vec![
1446                timestamp_column,
1447                field_column.clone(),
1448                field_column,
1449                path_column,
1450            ],
1451        )
1452        .unwrap();
1453
1454        let projected = project_batch(&input_batch, required);
1455        let projected_schema = projected.schema();
1456        let memory_exec = Arc::new(DataSourceExec::new(Arc::new(
1457            MemorySourceConfig::try_new(&[vec![projected]], projected_schema, None).unwrap(),
1458        )));
1459        let range_exec = plan.to_execution_plan(memory_exec);
1460        let session_context = SessionContext::default();
1461        let output_batches =
1462            datafusion::physical_plan::collect(range_exec, session_context.task_ctx())
1463                .await
1464                .unwrap();
1465        assert_eq!(output_batches.len(), 1);
1466
1467        let output_batch = &output_batches[0];
1468        let path = output_batch
1469            .column(3)
1470            .as_any()
1471            .downcast_ref::<StringArray>()
1472            .unwrap();
1473        assert!(path.iter().all(|v| v == Some("foo")));
1474
1475        // Simulate the pre-fix pruning behavior: omit the timestamp/value columns from the child.
1476        let broken_required = [3usize];
1477        let broken = project_batch(&input_batch, &broken_required);
1478        let broken_schema = broken.schema();
1479        let broken_exec = Arc::new(DataSourceExec::new(Arc::new(
1480            MemorySourceConfig::try_new(&[vec![broken]], broken_schema, None).unwrap(),
1481        )));
1482        let broken_range_exec = plan.to_execution_plan(broken_exec);
1483        let session_context = SessionContext::default();
1484        let broken_result = std::panic::AssertUnwindSafe(async {
1485            datafusion::physical_plan::collect(broken_range_exec, session_context.task_ctx()).await
1486        })
1487        .catch_unwind()
1488        .await;
1489        assert!(broken_result.is_err());
1490    }
1491
1492    #[tokio::test]
1493    async fn interval_30s_range_90s() {
1494        let expected = String::from(
1495            "PrimitiveArray<Timestamp(ms)>\n[\n  \
1496                1970-01-01T00:00:00,\n  \
1497                1970-01-01T00:00:30,\n  \
1498                1970-01-01T00:01:00,\n  \
1499                1970-01-01T00:01:30,\n  \
1500                1970-01-01T00:02:00,\n  \
1501                1970-01-01T00:02:30,\n  \
1502                1970-01-01T00:03:00,\n  \
1503                1970-01-01T00:03:30,\n  \
1504                1970-01-01T00:04:00,\n  \
1505                1970-01-01T00:04:30,\n  \
1506                1970-01-01T00:05:00,\n\
1507            ]\nRangeArray { \
1508                base array: PrimitiveArray<Float64>\n[\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n], \
1509                ranges: [Some(0..1), Some(0..2), Some(0..3), Some(1..4), Some(2..5), Some(3..5), Some(4..6), Some(5..6), Some(5..7), Some(6..8), Some(6..10)] \
1510            }\nRangeArray { \
1511                base array: PrimitiveArray<Float64>\n[\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n], \
1512                ranges: [Some(0..1), Some(0..2), Some(0..3), Some(1..4), Some(2..5), Some(3..5), Some(4..6), Some(5..6), Some(5..7), Some(6..8), Some(6..10)] \
1513            }\nStringArray\n[\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n]\n\
1514            RangeArray { \
1515                base array: PrimitiveArray<Timestamp(ms)>\n[\n  1970-01-01T00:00:00,\n  1970-01-01T00:00:30,\n  1970-01-01T00:01:00,\n  1970-01-01T00:01:30,\n  1970-01-01T00:02:00,\n  1970-01-01T00:03:00,\n  1970-01-01T00:04:00,\n  1970-01-01T00:04:01,\n  1970-01-01T00:04:31,\n  1970-01-01T00:04:51,\n], \
1516                ranges: [Some(0..1), Some(0..2), Some(0..3), Some(1..4), Some(2..5), Some(3..5), Some(4..6), Some(5..6), Some(5..7), Some(6..8), Some(6..10)] \
1517            }",
1518        );
1519        do_normalize_test(0, 310_000, 30_000, 90_000, expected.clone()).await;
1520
1521        // dump large range
1522        do_normalize_test(-300000, 310_000, 30_000, 90_000, expected).await;
1523    }
1524
1525    #[tokio::test]
1526    async fn small_empty_range() {
1527        let expected = String::from(
1528            "PrimitiveArray<Timestamp(ms)>\n[\n  \
1529            1970-01-01T00:00:00.001,\n  \
1530            1970-01-01T00:00:03.001,\n  \
1531            1970-01-01T00:00:06.001,\n  \
1532            1970-01-01T00:00:09.001,\n\
1533        ]\nRangeArray { \
1534            base array: PrimitiveArray<Float64>\n[\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n], \
1535            ranges: [Some(0..1), Some(0..0), Some(0..0), Some(0..0)] \
1536        }\nRangeArray { \
1537            base array: PrimitiveArray<Float64>\n[\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n  1.0,\n], \
1538            ranges: [Some(0..1), Some(0..0), Some(0..0), Some(0..0)] \
1539        }\nStringArray\n[\n  \"foo\",\n  \"foo\",\n  \"foo\",\n  \"foo\",\n]\n\
1540        RangeArray { \
1541            base array: PrimitiveArray<Timestamp(ms)>\n[\n  1970-01-01T00:00:00,\n  1970-01-01T00:00:30,\n  1970-01-01T00:01:00,\n  1970-01-01T00:01:30,\n  1970-01-01T00:02:00,\n  1970-01-01T00:03:00,\n  1970-01-01T00:04:00,\n  1970-01-01T00:04:01,\n  1970-01-01T00:04:31,\n  1970-01-01T00:04:51,\n], \
1542            ranges: [Some(0..1), Some(0..0), Some(0..0), Some(0..0)] \
1543        }",
1544        );
1545        do_normalize_test(1, 10_001, 3_000, 1_000, expected).await;
1546    }
1547
1548    #[test]
1549    fn test_calculate_range_preserves_alignment() {
1550        // Test case: query starts at timestamp ending in 4000, step is 30s
1551        // Data starts at different alignment - should preserve query's 4000 pattern
1552        let schema = Arc::new(Schema::new(vec![Field::new(
1553            "timestamp",
1554            TimestampMillisecondType::DATA_TYPE,
1555            false,
1556        )]));
1557        let empty_stream = MemoryStream::try_new(vec![], schema.clone(), None).unwrap();
1558
1559        let stream = RangeManipulateStream {
1560            offset: 0,
1561            start: 1758093274000, // ends in 4000
1562            end: 1758093334000,   // ends in 4000
1563            interval: 30000,      // 30s step
1564            range: 60000,         // 60s lookback
1565            time_index: 0,
1566            time_unit: TimeUnit::Millisecond,
1567            field_columns: vec![],
1568            aligned_ts_array: Arc::new(TimestampMillisecondArray::from(vec![0i64; 0])),
1569            output_schema: schema.clone(),
1570            input: Box::pin(empty_stream),
1571            metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
1572            num_series: Count::new(),
1573        };
1574
1575        // Create test data with timestamps not aligned to query pattern
1576        let test_timestamps = vec![
1577            1758093260000, // ends in 0000 (different alignment)
1578            1758093290000, // ends in 0000
1579            1758093320000, // ends in 0000
1580        ];
1581        let ts_array = TimestampMillisecondArray::from(test_timestamps);
1582        let test_schema = Arc::new(Schema::new(vec![Field::new(
1583            "timestamp",
1584            TimestampMillisecondType::DATA_TYPE,
1585            false,
1586        )]));
1587        let batch = RecordBatch::try_new(test_schema, vec![Arc::new(ts_array)]).unwrap();
1588
1589        let (ranges, (start, end)) = stream.calculate_range(&batch).unwrap();
1590
1591        // Verify the optimized start preserves query alignment (should end in 4000)
1592        assert_eq!(
1593            start % 30000,
1594            1758093274000 % 30000,
1595            "Optimized start should preserve query alignment pattern"
1596        );
1597
1598        // Verify we generate correct number of ranges for the alignment
1599        let expected_timestamps: Vec<i64> = (start..=end).step_by(30000).collect();
1600        assert_eq!(ranges.len(), expected_timestamps.len());
1601
1602        // Verify all generated timestamps maintain the same alignment pattern
1603        for ts in expected_timestamps {
1604            assert_eq!(
1605                ts % 30000,
1606                1758093274000 % 30000,
1607                "All timestamps should maintain query alignment pattern"
1608            );
1609        }
1610    }
1611
1612    #[tokio::test]
1613    async fn no_intersection_batch_is_skipped_and_stream_continues() {
1614        let schema = Arc::new(Schema::new(vec![
1615            Field::new(
1616                TIME_INDEX_COLUMN,
1617                TimestampMillisecondType::DATA_TYPE,
1618                false,
1619            ),
1620            Field::new("value", DataType::Float64, false),
1621        ]));
1622        let input = LogicalPlan::EmptyRelation(EmptyRelation {
1623            produce_one_row: false,
1624            schema: schema.clone().to_dfschema_ref().unwrap(),
1625        });
1626        let plan = RangeManipulate::new(
1627            0,
1628            50,
1629            10,
1630            0,
1631            1,
1632            TIME_INDEX_COLUMN.to_string(),
1633            vec!["value".to_string()],
1634            input,
1635        )
1636        .unwrap();
1637        let no_intersection = RecordBatch::try_new(
1638            schema.clone(),
1639            vec![
1640                Arc::new(TimestampMillisecondArray::from(vec![100])),
1641                Arc::new(Float64Array::from(vec![1.0])),
1642            ],
1643        )
1644        .unwrap();
1645        let intersection = RecordBatch::try_new(
1646            schema.clone(),
1647            vec![
1648                Arc::new(TimestampMillisecondArray::from(vec![20])),
1649                Arc::new(Float64Array::from(vec![2.0])),
1650            ],
1651        )
1652        .unwrap();
1653        let input = Arc::new(DataSourceExec::new(Arc::new(
1654            MemorySourceConfig::try_new(&[vec![no_intersection, intersection]], schema, None)
1655                .unwrap(),
1656        )));
1657
1658        let batches = datafusion::physical_plan::collect(
1659            plan.to_execution_plan(input),
1660            SessionContext::default().task_ctx(),
1661        )
1662        .await
1663        .unwrap();
1664
1665        assert_eq!(batches.len(), 1);
1666        assert_eq!(batches[0].num_rows(), 1);
1667        let timestamps = batches[0]
1668            .column(0)
1669            .as_any()
1670            .downcast_ref::<TimestampMillisecondArray>()
1671            .unwrap();
1672        assert_eq!(timestamps.value(0), 20);
1673    }
1674
1675    fn calculate_range_for_test(
1676        query_start: i64,
1677        query_end: i64,
1678        interval: i64,
1679        range: i64,
1680        timestamps: &[i64],
1681    ) -> (Vec<(u32, u32)>, (i64, i64)) {
1682        let schema = Arc::new(Schema::new(vec![Field::new(
1683            TIME_INDEX_COLUMN,
1684            TimestampMillisecondType::DATA_TYPE,
1685            false,
1686        )]));
1687        let empty_stream = MemoryStream::try_new(vec![], schema.clone(), None).unwrap();
1688        let stream = RangeManipulateStream {
1689            offset: 0,
1690            start: query_start,
1691            end: query_end,
1692            interval,
1693            range,
1694            time_index: 0,
1695            time_unit: TimeUnit::Millisecond,
1696            field_columns: vec![],
1697            aligned_ts_array: Arc::new(TimestampMillisecondArray::from(vec![0i64; 0])),
1698            output_schema: schema.clone(),
1699            input: Box::pin(empty_stream),
1700            metric: BaselineMetrics::new(&ExecutionPlanMetricsSet::new(), 0),
1701            num_series: Count::new(),
1702        };
1703        let batch = RecordBatch::try_new(
1704            schema,
1705            vec![Arc::new(TimestampMillisecondArray::from(
1706                timestamps.to_vec(),
1707            ))],
1708        )
1709        .unwrap();
1710
1711        stream.calculate_range(&batch).unwrap()
1712    }
1713
1714    #[test]
1715    fn calculate_range_keeps_query_aligned_tail() {
1716        let (ranges, bounds) = calculate_range_for_test(4, 94, 30, 15, &[20, 50, 80]);
1717
1718        assert_eq!(bounds, (34, 94));
1719        assert_eq!(ranges, vec![(0, 1), (1, 1), (2, 1)]);
1720    }
1721
1722    /// Calculates exact offsets directly from the range predicate used by PromQL.
1723    ///
1724    /// Input timestamps are sorted and non-null. Interval is positive, range is
1725    /// nonnegative, and test values are chosen to avoid `i64` overflow.
1726    fn calculate_range_oracle(
1727        timestamps: &[i64],
1728        start: i64,
1729        end: i64,
1730        interval: i64,
1731        range: i64,
1732    ) -> Vec<(u32, u32)> {
1733        // Match `calculate_range`'s explicit empty-input early return.
1734        if timestamps.is_empty() || start > end {
1735            return vec![];
1736        }
1737
1738        (start..=end)
1739            .step_by(interval as usize)
1740            .map(|curr| {
1741                let mut offset = None;
1742                let mut length = 0;
1743                for (index, &ts) in timestamps.iter().enumerate() {
1744                    if ts > curr - range && ts <= curr {
1745                        offset.get_or_insert(index);
1746                        length += 1;
1747                    }
1748                }
1749                (offset.unwrap_or(0) as u32, length)
1750            })
1751            .collect()
1752    }
1753
1754    #[test]
1755    fn calculate_range_characterizes_returned_bounds() {
1756        let cases = [
1757            (
1758                "positive non-aligned last timestamp plus range",
1759                0,
1760                100,
1761                10,
1762                9,
1763                vec![13, 26],
1764                (20, 30),
1765            ),
1766            (
1767                "negative non-aligned last timestamp plus range",
1768                -50,
1769                50,
1770                10,
1771                9,
1772                vec![-37, -26],
1773                (-30, -20),
1774            ),
1775            (
1776                "query alignment not based on epoch",
1777                4,
1778                94,
1779                30,
1780                15,
1781                vec![20, 50, 80],
1782                (34, 94),
1783            ),
1784            ("leading data", 0, 100, 10, 10, vec![-10, 15], (0, 20)),
1785            (
1786                "trailing data past query end",
1787                0,
1788                100,
1789                10,
1790                10,
1791                vec![35, 45, 110],
1792                (40, 100),
1793            ),
1794            (
1795                "optimized start after query end",
1796                0,
1797                50,
1798                10,
1799                0,
1800                vec![100],
1801                (0, 50),
1802            ),
1803        ];
1804
1805        for (name, query_start, query_end, interval, range, timestamps, expected_bounds) in cases {
1806            let (_, bounds) =
1807                calculate_range_for_test(query_start, query_end, interval, range, &timestamps);
1808            assert_eq!(bounds, expected_bounds, "{name}");
1809        }
1810    }
1811
1812    #[test]
1813    fn calculate_range_keeps_extreme_range_tail() {
1814        let (ranges, bounds) =
1815            calculate_range_for_test(i64::MAX - 1, i64::MAX, 1, i64::MAX, &[i64::MAX]);
1816
1817        assert_eq!(bounds, (i64::MAX, i64::MAX));
1818        assert_eq!(ranges, vec![(0, 1)]);
1819    }
1820
1821    #[test]
1822    fn calculate_range_matches_bruteforce_oracle_for_deterministic_cases() {
1823        let cases = vec![
1824            (
1825                "duplicate lower and upper bounds",
1826                vec![0, 10, 10, 20, 20, 30],
1827                10,
1828                20,
1829                10,
1830                10,
1831            ),
1832            (
1833                "zero range excludes duplicates at current timestamp",
1834                vec![10, 10, 10],
1835                10,
1836                10,
1837                1,
1838                0,
1839            ),
1840            (
1841                "consecutive nonempty empty nonempty ranges",
1842                vec![10, 30],
1843                10,
1844                30,
1845                10,
1846                5,
1847            ),
1848            ("step smaller than range", vec![0, 4, 8, 12], 0, 12, 3, 5),
1849            ("step equal to range", vec![0, 5, 10, 15], 0, 15, 5, 5),
1850            ("step greater than range", vec![0, 7, 14, 21], 0, 21, 7, 3),
1851            (
1852                "negative sparse/tail timestamps",
1853                vec![-30, -20, -10, 0],
1854                -25,
1855                5,
1856                5,
1857                7,
1858            ),
1859            ("one sample", vec![42], 0, 100, 10, 15),
1860            ("empty input", vec![], -20, 20, 5, 10),
1861        ];
1862
1863        for (name, timestamps, query_start, query_end, interval, range) in cases {
1864            let (actual, (start, end)) =
1865                calculate_range_for_test(query_start, query_end, interval, range, &timestamps);
1866            let expected = calculate_range_oracle(&timestamps, start, end, interval, range);
1867            assert_eq!(actual, expected, "{name}");
1868        }
1869    }
1870
1871    #[test]
1872    fn calculate_range_positive_time_translated_regression() {
1873        let timestamps = [0, 10, 20, 30];
1874        let expected = vec![(0, 1), (1, 1), (1, 1), (2, 1), (2, 1), (3, 1), (3, 1)];
1875        let (actual, bounds) = calculate_range_for_test(5, 35, 5, 7, &timestamps);
1876
1877        assert_eq!(bounds, (5, 35));
1878        assert_eq!(
1879            calculate_range_oracle(&timestamps, bounds.0, bounds.1, 5, 7),
1880            expected
1881        );
1882        assert_eq!(actual, expected);
1883    }
1884
1885    #[test]
1886    fn calculate_range_matches_oracle_for_dense_positive_time_windows() {
1887        let timestamps = (0..=3_600).step_by(15).collect::<Vec<i64>>();
1888
1889        for (name, range) in [
1890            ("one minute", 60),
1891            ("five minutes", 300),
1892            ("one hour", 3_600),
1893        ] {
1894            let (actual, (start, end)) = calculate_range_for_test(0, 3_600, 15, range, &timestamps);
1895            assert_eq!((start, end), (0, 3_600), "{name} bounds");
1896            assert_eq!(
1897                actual,
1898                calculate_range_oracle(&timestamps, start, end, 15, range),
1899                "{name} window"
1900            );
1901        }
1902    }
1903
1904    #[derive(Clone, Copy)]
1905    struct TinyPrng(u64);
1906
1907    impl TinyPrng {
1908        fn next_u64(&mut self) -> u64 {
1909            self.0 ^= self.0 << 13;
1910            self.0 ^= self.0 >> 7;
1911            self.0 ^= self.0 << 17;
1912            self.0
1913        }
1914
1915        fn next_i64(&mut self, min: i64, max: i64) -> i64 {
1916            min + (self.next_u64() % (max - min + 1) as u64) as i64
1917        }
1918    }
1919
1920    #[test]
1921    fn calculate_range_matches_bruteforce_oracle_for_seeded_matrix() {
1922        let mut prng = TinyPrng(0x5eed_cafe_f00d_baad);
1923
1924        for case in 0..512 {
1925            let interval = prng.next_i64(1, 11);
1926            let range = prng.next_i64(0, 25);
1927            let query_start = prng.next_i64(-200, 200);
1928            let query_end = query_start + interval * prng.next_i64(0, 20);
1929            let mut timestamps = Vec::new();
1930            let mut timestamp = prng.next_i64(-250, 250);
1931            for _ in 0..prng.next_i64(0, 20) {
1932                timestamp += prng.next_i64(0, 7);
1933                timestamps.push(timestamp);
1934            }
1935
1936            let (actual, (start, end)) =
1937                calculate_range_for_test(query_start, query_end, interval, range, &timestamps);
1938            let expected = calculate_range_oracle(&timestamps, start, end, interval, range);
1939            let expected = if actual.is_empty() && !expected.is_empty() {
1940                assert!(
1941                    expected.iter().all(|(_, len)| *len == 0),
1942                    "case={case}, timestamps={timestamps:?}, query=({query_start}, {query_end}), \
1943                     interval={interval}, range={range}, bounds=({start}, {end}): \
1944                     no-intersection output must have no selected samples"
1945                );
1946                vec![]
1947            } else {
1948                expected
1949            };
1950            assert_eq!(
1951                actual, expected,
1952                "case={case}, timestamps={timestamps:?}, query=({query_start}, {query_end}), \
1953                 interval={interval}, range={range}, bounds=({start}, {end})"
1954            );
1955        }
1956    }
1957}