Skip to main content

timeseries_table_format/table/
scan.rs

1//! Range scan implementation for `TimeSeriesTable`.
2//!
3//! This module wires the public `scan_range` entry point to the underlying
4//! segment metadata and Parquet readers:
5//! - Pick candidate segments whose ordered `index_min`/`index_max` intersects the requested
6//!   half-open range `[start, end)`.
7//! - Visit candidate segments deterministically with Parquet's native async,
8//!   file-backed reader.
9//! - Filter each batch by its ordered-index column with native Arrow scalar
10//!   comparisons and half-open semantics.
11//!
12//! The filtering path uses Arrow's scalar comparison kernels to avoid
13//! allocating full-length bound arrays, and treats null index values as
14//! "drop row" via `filter_record_batch`. Input rows need not be ordered, and
15//! the returned batches and rows have no ordering guarantee.
16use std::path::Path;
17
18use arrow::array::{Datum, Scalar};
19use arrow::array::{
20    Int64Array, RecordBatch, TimestampMicrosecondArray, TimestampMillisecondArray,
21    TimestampNanosecondArray, TimestampSecondArray, UInt64Array,
22};
23use arrow::compute::filter_record_batch;
24use arrow::compute::kernels::{boolean as boolean_kernels, cmp as cmp_kernels};
25use arrow::datatypes::{DataType, Field, TimeUnit};
26use chrono::{DateTime, Utc};
27use futures::{StreamExt, TryStreamExt, future};
28use parquet::arrow::async_reader::{AsyncFileReader, ParquetRecordBatchStreamBuilder};
29use snafu::prelude::*;
30
31use crate::metadata::{
32    segments::SegmentMeta,
33    table_metadata::{IndexValue, IndexValueError, validate_index_range},
34};
35use crate::storage::{self, TableLocation};
36use crate::transaction_log::TableState;
37
38use super::error::{InvalidRangeSnafu, InvalidSegmentBoundsSnafu, StorageSnafu, TableError};
39use super::{TimeSeriesScan, TimeSeriesTable};
40
41const SCAN_BATCH_SIZE: usize = 8_192;
42
43fn segments_for_range(
44    state: &TableState,
45    start: &IndexValue,
46    end: &IndexValue,
47) -> Result<Vec<SegmentMeta>, IndexValueError> {
48    let mut candidates = Vec::new();
49    for segment in state.segments_sorted_by_index()? {
50        if !segment.index_max.compare(start)?.is_lt() && segment.index_min.compare(end)?.is_lt() {
51            candidates.push(segment.clone());
52        }
53    }
54    Ok(candidates)
55}
56
57/// Filter one batch with native scalar comparisons. Arrow broadcasts each
58/// scalar bound without allocating a batch-sized bound column.
59fn filter_index_batch(
60    batch: RecordBatch,
61    index_idx: usize,
62    start: &dyn Datum,
63    end: &dyn Datum,
64    path: &str,
65    index_column: &str,
66) -> Result<Option<RecordBatch>, TableError> {
67    let index_array = batch.column(index_idx).as_ref();
68    let arrow_error = |source| TableError::Arrow {
69        path: path.to_string(),
70        column: index_column.to_string(),
71        source,
72    };
73    let ge_mask = cmp_kernels::gt_eq(&index_array, start).map_err(&arrow_error)?;
74    let lt_mask = cmp_kernels::lt(&index_array, end).map_err(&arrow_error)?;
75    let mask = boolean_kernels::and(&ge_mask, &lt_mask).map_err(&arrow_error)?;
76    let filtered = filter_record_batch(&batch, &mask).map_err(arrow_error)?;
77    Ok((filtered.num_rows() > 0).then_some(filtered))
78}
79
80#[derive(Clone, Copy)]
81enum ScanBounds {
82    Timestamp(i64, i64),
83    Int64(i64, i64),
84    UInt64(u64, u64),
85}
86
87fn to_bounds_i64(
88    field: &Field,
89    path: &str,
90    column: &str,
91    ts_start: DateTime<Utc>,
92    ts_end: DateTime<Utc>,
93) -> Result<(i64, i64), TableError> {
94    let ceil_bound = |dt: DateTime<Utc>, floor: i64, nanos_per_unit: u32| {
95        if dt.timestamp_subsec_nanos().is_multiple_of(nanos_per_unit) {
96            Ok(floor)
97        } else {
98            floor
99                .checked_add(1)
100                .ok_or_else(|| TableError::TimeConversionOverflow {
101                    path: path.to_string(),
102                    column: column.to_string(),
103                    timestamp: dt,
104                })
105        }
106    };
107    let to_ns = |dt: DateTime<Utc>| {
108        dt.timestamp()
109            .checked_mul(1_000_000_000)
110            .and_then(|secs| secs.checked_add(dt.timestamp_subsec_nanos() as i64))
111            .ok_or_else(|| TableError::TimeConversionOverflow {
112                path: path.to_string(),
113                column: column.to_string(),
114                timestamp: dt,
115            })
116    };
117
118    match field.data_type() {
119        DataType::Timestamp(TimeUnit::Second, _) => Ok((
120            ceil_bound(ts_start, ts_start.timestamp(), 1_000_000_000)?,
121            ceil_bound(ts_end, ts_end.timestamp(), 1_000_000_000)?,
122        )),
123
124        DataType::Timestamp(TimeUnit::Millisecond, _) => Ok((
125            ceil_bound(ts_start, ts_start.timestamp_millis(), 1_000_000)?,
126            ceil_bound(ts_end, ts_end.timestamp_millis(), 1_000_000)?,
127        )),
128
129        DataType::Timestamp(TimeUnit::Microsecond, _) => Ok((
130            ceil_bound(ts_start, ts_start.timestamp_micros(), 1_000)?,
131            ceil_bound(ts_end, ts_end.timestamp_micros(), 1_000)?,
132        )),
133
134        DataType::Timestamp(TimeUnit::Nanosecond, _) => Ok((to_ns(ts_start)?, to_ns(ts_end)?)),
135
136        other => Err(TableError::IndexColumnTypeMismatch {
137            path: path.to_string(),
138            column: column.to_string(),
139            expected: "timestamp",
140            datatype: other.clone(),
141        }),
142    }
143}
144
145async fn segment_range_stream<S, E>(
146    reader: impl AsyncFileReader + Unpin + 'static,
147    path: String,
148    index_column: &str,
149    start: S,
150    end: E,
151) -> Result<TimeSeriesScan, TableError>
152where
153    S: Into<IndexValue>,
154    E: Into<IndexValue>,
155{
156    let start = start.into();
157    let end = end.into();
158    let expected = start.kind_name();
159
160    let builder = ParquetRecordBatchStreamBuilder::new(reader)
161        .await
162        .map_err(|source| TableError::ParquetRead {
163            path: path.clone(),
164            source,
165        })?;
166
167    // Locate the index column and compute native bounds before moving the
168    // builder into the directly-polled record-batch stream.
169    let schema = builder.schema();
170    let index_idx = schema
171        .index_of(index_column)
172        .map_err(|_| TableError::MissingIndexColumn {
173            path: path.clone(),
174            column: index_column.to_string(),
175        })?;
176    let index_field = schema.field(index_idx).clone();
177    let bounds = match (&start, &end, index_field.data_type()) {
178        (IndexValue::Timestamp(start), IndexValue::Timestamp(end), DataType::Timestamp(_, _)) => {
179            let (start, end) = to_bounds_i64(&index_field, &path, index_column, *start, *end)?;
180            ScanBounds::Timestamp(start, end)
181        }
182        (IndexValue::Int64(start), IndexValue::Int64(end), DataType::Int64) => {
183            ScanBounds::Int64(*start, *end)
184        }
185        (IndexValue::UInt64(start), IndexValue::UInt64(end), DataType::UInt64) => {
186            ScanBounds::UInt64(*start, *end)
187        }
188        _ => {
189            return Err(TableError::IndexColumnTypeMismatch {
190                path,
191                column: index_column.to_string(),
192                expected,
193                datatype: index_field.data_type().clone(),
194            });
195        }
196    };
197
198    let reader = builder
199        .with_batch_size(SCAN_BATCH_SIZE)
200        .build()
201        .map_err(|source| TableError::ParquetRead {
202            path: path.clone(),
203            source,
204        })?;
205    let index_column = index_column.to_string();
206
207    let stream = reader
208        .then(move |batch_res| {
209            let path = path.clone();
210            let index_column = index_column.clone();
211            let index_field = index_field.clone();
212
213            async move {
214                let batch = batch_res.map_err(|source| TableError::ParquetRead {
215                    path: path.clone(),
216                    source,
217                })?;
218
219                let filtered = match (bounds, index_field.data_type()) {
220                    (ScanBounds::Timestamp(start, end), DataType::Timestamp(unit, timezone)) => {
221                        macro_rules! filter_timestamp {
222                            ($array_ty:ty) => {{
223                                let start = Scalar::new(
224                                    <$array_ty>::from(vec![start])
225                                        .with_timezone_opt(timezone.clone()),
226                                );
227                                let end = Scalar::new(
228                                    <$array_ty>::from(vec![end])
229                                        .with_timezone_opt(timezone.clone()),
230                                );
231                                filter_index_batch(
232                                    batch,
233                                    index_idx,
234                                    &start,
235                                    &end,
236                                    &path,
237                                    &index_column,
238                                )
239                            }};
240                        }
241                        match unit {
242                            TimeUnit::Second => filter_timestamp!(TimestampSecondArray)?,
243                            TimeUnit::Millisecond => filter_timestamp!(TimestampMillisecondArray)?,
244                            TimeUnit::Microsecond => filter_timestamp!(TimestampMicrosecondArray)?,
245                            TimeUnit::Nanosecond => filter_timestamp!(TimestampNanosecondArray)?,
246                        }
247                    }
248                    (ScanBounds::Int64(start, end), DataType::Int64) => {
249                        let start = Scalar::new(Int64Array::from(vec![start]));
250                        let end = Scalar::new(Int64Array::from(vec![end]));
251                        filter_index_batch(batch, index_idx, &start, &end, &path, &index_column)?
252                    }
253                    (ScanBounds::UInt64(start, end), DataType::UInt64) => {
254                        let start = Scalar::new(UInt64Array::from(vec![start]));
255                        let end = Scalar::new(UInt64Array::from(vec![end]));
256                        filter_index_batch(batch, index_idx, &start, &end, &path, &index_column)?
257                    }
258                    (_, datatype) => {
259                        return Err(TableError::IndexColumnTypeMismatch {
260                            path,
261                            column: index_column,
262                            expected,
263                            datatype: datatype.clone(),
264                        });
265                    }
266                };
267
268                if filtered.is_none() {
269                    tokio::task::yield_now().await;
270                }
271
272                Ok(filtered)
273            }
274        })
275        .try_filter_map(|batch| future::ready(Ok(batch)));
276
277    Ok(Box::pin(stream))
278}
279
280async fn read_segment_range<S, E>(
281    location: &TableLocation,
282    segment: &SegmentMeta,
283    index_column: &str,
284    start: S,
285    end: E,
286) -> Result<TimeSeriesScan, TableError>
287where
288    S: Into<IndexValue>,
289    E: Into<IndexValue>,
290{
291    let rel_path = Path::new(&segment.path);
292    let reader = storage::open_parquet_reader(location.as_ref(), rel_path)
293        .await
294        .context(StorageSnafu)?;
295
296    segment_range_stream(reader, segment.path.clone(), index_column, start, end).await
297}
298
299struct ScanState {
300    candidates: std::vec::IntoIter<SegmentMeta>,
301    current: Option<TimeSeriesScan>,
302    location: TableLocation,
303    index_column: String,
304    start: IndexValue,
305    end: IndexValue,
306}
307
308impl TimeSeriesTable {
309    /// Scan the time-series table for record batches overlapping `[start, end)`,
310    /// returning a stream of filtered batches from the segments covering that range.
311    ///
312    /// Input rows need not be ordered. The returned batches and rows have no
313    /// ordering guarantee; callers that need ordered results must sort them.
314    pub async fn scan_range<S, E>(&self, start: S, end: E) -> Result<TimeSeriesScan, TableError>
315    where
316        S: Into<IndexValue>,
317        E: Into<IndexValue>,
318    {
319        let start = start.into();
320        let end = end.into();
321        validate_index_range(&self.index.kind, &start, &end).context(InvalidRangeSnafu)?;
322
323        // Pick candidate segments and sort them by index_min, index_max, and path.
324        let candidates =
325            segments_for_range(&self.state, &start, &end).context(InvalidSegmentBoundsSnafu)?;
326
327        let state = ScanState {
328            candidates: candidates.into_iter(),
329            current: None,
330            location: self.location().clone(),
331            index_column: self.index.column.clone(),
332            start,
333            end,
334        };
335
336        // Process one lazily opened segment stream at a time. `try_unfold`
337        // drops its state after the first error, so later segments are not
338        // opened after a terminal scan failure.
339        let stream = futures::stream::try_unfold(state, |mut state| async move {
340            loop {
341                if let Some(current) = state.current.as_mut() {
342                    match current.next().await {
343                        Some(Ok(batch)) => return Ok(Some((batch, state))),
344                        Some(Err(error)) => return Err(error),
345                        None => state.current = None,
346                    }
347                }
348
349                let Some(segment) = state.candidates.next() else {
350                    return Ok(None);
351                };
352                state.current = Some(
353                    read_segment_range(
354                        &state.location,
355                        &segment,
356                        &state.index_column,
357                        state.start.clone(),
358                        state.end.clone(),
359                    )
360                    .await?,
361                );
362            }
363        });
364
365        Ok(Box::pin(stream))
366    }
367}
368
369#[cfg(test)]
370mod tests {
371    use super::*;
372    use crate::storage::TableLocation;
373    use crate::table::test_util::*;
374
375    use crate::metadata::logical_schema::LogicalTimestampUnit;
376    use crate::metadata::segments::{FileFormat, SegmentEntityLayout};
377    use crate::metadata::table_metadata::{IndexKind, IndexSpec, TableMeta};
378
379    use arrow::array::ArrayRef;
380    use arrow::datatypes::{Schema, TimeUnit as ArrowTimeUnit};
381
382    use chrono::{TimeZone, Utc};
383    use futures::{FutureExt, StreamExt, future::BoxFuture};
384    use parquet::arrow::ArrowWriter;
385    use parquet::arrow::arrow_reader::ArrowReaderOptions;
386    use parquet::errors::Result as ParquetResult;
387    use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
388    use parquet::file::properties::WriterProperties;
389
390    use std::fs::File;
391    use std::num::NonZeroU64;
392    use std::ops::Range;
393    use std::sync::Arc;
394    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
395    use tempfile::TempDir;
396
397    #[derive(Default)]
398    struct TrackingStats {
399        read_calls: AtomicUsize,
400        bytes_read: AtomicUsize,
401        dropped: AtomicBool,
402    }
403
404    struct TrackingReader {
405        data: bytes::Bytes,
406        stats: Arc<TrackingStats>,
407        gate: Option<(usize, futures::channel::oneshot::Receiver<()>)>,
408    }
409
410    impl TrackingReader {
411        fn new(data: bytes::Bytes) -> (Self, Arc<TrackingStats>) {
412            let stats = Arc::new(TrackingStats::default());
413            (
414                Self {
415                    data,
416                    stats: Arc::clone(&stats),
417                    gate: None,
418                },
419                stats,
420            )
421        }
422
423        fn with_gate(
424            data: bytes::Bytes,
425            gated_call: usize,
426        ) -> (
427            Self,
428            Arc<TrackingStats>,
429            futures::channel::oneshot::Sender<()>,
430        ) {
431            let (mut reader, stats) = Self::new(data);
432            let (release, gate) = futures::channel::oneshot::channel();
433            reader.gate = Some((gated_call, gate));
434            (reader, stats, release)
435        }
436
437        fn read_ranges(&self, ranges: Vec<Range<u64>>) -> (usize, Vec<bytes::Bytes>) {
438            let call = self.stats.read_calls.fetch_add(1, Ordering::SeqCst) + 1;
439            self.stats.bytes_read.fetch_add(
440                ranges
441                    .iter()
442                    .map(|range| (range.end - range.start) as usize)
443                    .sum::<usize>(),
444                Ordering::SeqCst,
445            );
446            (
447                call,
448                ranges
449                    .into_iter()
450                    .map(|range| self.data.slice(range.start as usize..range.end as usize))
451                    .collect(),
452            )
453        }
454    }
455
456    impl Drop for TrackingReader {
457        fn drop(&mut self) {
458            self.stats.dropped.store(true, Ordering::SeqCst);
459        }
460    }
461
462    impl AsyncFileReader for TrackingReader {
463        fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, ParquetResult<bytes::Bytes>> {
464            let (_, mut ranges) = self.read_ranges(vec![range]);
465            futures::future::ready(Ok(ranges.pop().expect("one requested range"))).boxed()
466        }
467
468        fn get_byte_ranges(
469            &mut self,
470            ranges: Vec<Range<u64>>,
471        ) -> BoxFuture<'_, ParquetResult<Vec<bytes::Bytes>>> {
472            let (call, bytes) = self.read_ranges(ranges);
473            let gate = self
474                .gate
475                .as_ref()
476                .is_some_and(|(gated_call, _)| *gated_call == call)
477                .then(|| self.gate.take().expect("configured gate").1);
478            async move {
479                if let Some(gate) = gate {
480                    let _ = gate.await;
481                }
482                Ok(bytes)
483            }
484            .boxed()
485        }
486
487        fn get_metadata<'a>(
488            &'a mut self,
489            _options: Option<&'a ArrowReaderOptions>,
490        ) -> BoxFuture<'a, ParquetResult<Arc<ParquetMetaData>>> {
491            let metadata = ParquetMetaDataReader::new()
492                .parse_and_finish(&self.data)
493                .map(Arc::new);
494            futures::future::ready(metadata).boxed()
495        }
496    }
497
498    fn indexed_segment(path: &str, min: IndexValue, max: IndexValue) -> SegmentMeta {
499        SegmentMeta {
500            path: path.to_string(),
501            format: FileFormat::Parquet,
502            entity_layout: SegmentEntityLayout::NotApplicable,
503            index_min: min,
504            index_max: max,
505            row_count: 1,
506            file_size: None,
507            coverage_path: None,
508        }
509    }
510
511    fn state_with_segments(segments: Vec<SegmentMeta>) -> TableState {
512        TableState {
513            version: 1,
514            table_meta: make_basic_table_meta(),
515            segments: segments
516                .into_iter()
517                .map(|segment| (segment.path.clone(), segment))
518                .collect(),
519            table_coverage: None,
520        }
521    }
522
523    fn integer_table_meta(kind: IndexKind) -> TableMeta {
524        TableMeta::new_time_series(IndexSpec {
525            column: "ts".to_string(),
526            entity_columns: Vec::new(),
527            kind,
528        })
529    }
530
531    fn write_index_parquet(path: &Path, values: ArrayRef) -> TestResult {
532        if let Some(parent) = path.parent() {
533            std::fs::create_dir_all(parent)?;
534        }
535        let schema = Arc::new(Schema::new(vec![Field::new(
536            "ts",
537            values.data_type().clone(),
538            values.null_count() > 0,
539        )]));
540        let batch = RecordBatch::try_new(Arc::clone(&schema), vec![values])?;
541        let properties = WriterProperties::builder()
542            .set_max_row_group_size(1)
543            .build();
544        let mut writer = ArrowWriter::try_new(File::create(path)?, schema, Some(properties))?;
545        writer.write(&batch)?;
546        writer.close()?;
547        Ok(())
548    }
549
550    async fn collect_i64_index(
551        table: &TimeSeriesTable,
552        start: i64,
553        end: i64,
554    ) -> Result<Vec<i64>, TableError> {
555        let mut stream = table.scan_range(start, end).await?;
556        let mut values = Vec::new();
557        while let Some(batch) = stream.next().await.transpose()? {
558            assert_eq!(batch.schema().field(0).data_type(), &DataType::Int64);
559            values.extend(
560                batch
561                    .column(0)
562                    .as_any()
563                    .downcast_ref::<Int64Array>()
564                    .expect("int64 index")
565                    .iter()
566                    .flatten(),
567            );
568        }
569        Ok(values)
570    }
571
572    async fn collect_u64_index(
573        table: &TimeSeriesTable,
574        start: u64,
575        end: u64,
576    ) -> Result<Vec<u64>, TableError> {
577        let mut stream = table.scan_range(start, end).await?;
578        let mut values = Vec::new();
579        while let Some(batch) = stream.next().await.transpose()? {
580            assert_eq!(batch.schema().field(0).data_type(), &DataType::UInt64);
581            values.extend(
582                batch
583                    .column(0)
584                    .as_any()
585                    .downcast_ref::<UInt64Array>()
586                    .expect("uint64 index")
587                    .iter()
588                    .flatten(),
589            );
590        }
591        Ok(values)
592    }
593
594    #[test]
595    fn integer_candidates_preserve_half_open_order() -> Result<(), IndexValueError> {
596        let signed = state_with_segments(vec![
597            indexed_segment("before", (-10i64).into(), (-1i64).into()),
598            indexed_segment("touch-start", (-5i64).into(), 0i64.into()),
599            indexed_segment("same-b", 1i64.into(), 9i64.into()),
600            indexed_segment("same-a", 1i64.into(), 9i64.into()),
601            indexed_segment("at-end", 10i64.into(), 20i64.into()),
602        ]);
603        let signed_paths = segments_for_range(&signed, &0i64.into(), &10i64.into())?
604            .into_iter()
605            .map(|segment| segment.path)
606            .collect::<Vec<_>>();
607        assert_eq!(signed_paths, ["touch-start", "same-a", "same-b"]);
608
609        let start = i64::MAX as u64 + 1;
610        let unsigned = state_with_segments(vec![
611            indexed_segment("below", 0u64.into(), (start - 1).into()),
612            indexed_segment("touch-start", (start - 1).into(), start.into()),
613            indexed_segment("inside", (start + 1).into(), (u64::MAX - 1).into()),
614            indexed_segment("at-end", u64::MAX.into(), u64::MAX.into()),
615        ]);
616        let unsigned_paths = segments_for_range(&unsigned, &start.into(), &u64::MAX.into())?
617            .into_iter()
618            .map(|segment| segment.path)
619            .collect::<Vec<_>>();
620        assert_eq!(unsigned_paths, ["touch-start", "inside"]);
621        Ok(())
622    }
623
624    #[test]
625    fn candidate_selection_rejects_invalid_persisted_bounds() {
626        let inverted =
627            state_with_segments(vec![indexed_segment("inverted", 2i64.into(), 1i64.into())]);
628        assert!(matches!(
629            segments_for_range(&inverted, &0i64.into(), &3i64.into()),
630            Err(IndexValueError::InvalidBounds { .. })
631        ));
632
633        let signed = state_with_segments(vec![indexed_segment("signed", 0i64.into(), 1i64.into())]);
634        assert!(matches!(
635            segments_for_range(&signed, &0u64.into(), &2u64.into()),
636            Err(IndexValueError::DomainMismatch { .. })
637        ));
638    }
639
640    #[tokio::test]
641    async fn scan_range_filters_signed_integer_boundaries_and_nulls() -> TestResult {
642        let tmp = TempDir::new()?;
643        let kind = IndexKind::Int64 {
644            bucket_width: NonZeroU64::new(1).unwrap(),
645        };
646        let mut table =
647            TimeSeriesTable::create(TableLocation::local(tmp.path()), integer_table_meta(kind))
648                .await?;
649        let rel = "data/int64-scan.parquet";
650        write_index_parquet(
651            &tmp.path().join(rel),
652            Arc::new(Int64Array::from(vec![
653                Some(i64::MIN),
654                Some(-2),
655                None,
656                Some(-1),
657                Some(0),
658                Some(1),
659                Some(i64::MAX - 1),
660                Some(i64::MAX),
661            ])),
662        )?;
663        table.append_parquet_segment(rel).await?;
664
665        assert_eq!(collect_i64_index(&table, -2, 2).await?, [-2, -1, 0, 1]);
666        assert_eq!(
667            collect_i64_index(&table, i64::MIN, i64::MIN + 1).await?,
668            [i64::MIN]
669        );
670        assert_eq!(
671            collect_i64_index(&table, i64::MAX - 1, i64::MAX).await?,
672            [i64::MAX - 1]
673        );
674        Ok(())
675    }
676
677    #[tokio::test]
678    async fn scan_range_preserves_large_unsigned_values() -> TestResult {
679        let tmp = TempDir::new()?;
680        let kind = IndexKind::UInt64 {
681            bucket_width: NonZeroU64::new(1).unwrap(),
682        };
683        let mut table =
684            TimeSeriesTable::create(TableLocation::local(tmp.path()), integer_table_meta(kind))
685                .await?;
686        let rel = "data/uint64-scan.parquet";
687        let start = i64::MAX as u64 + 1;
688        write_index_parquet(
689            &tmp.path().join(rel),
690            Arc::new(UInt64Array::from(vec![
691                Some(start - 1),
692                Some(start),
693                None,
694                Some(start + 1),
695                Some(u64::MAX - 1),
696                Some(u64::MAX),
697            ])),
698        )?;
699        table.append_parquet_segment(rel).await?;
700
701        assert_eq!(
702            collect_u64_index(&table, start, u64::MAX).await?,
703            [start, start + 1, u64::MAX - 1]
704        );
705        Ok(())
706    }
707
708    #[tokio::test]
709    async fn scan_range_validates_typed_bounds_before_segment_checks() -> TestResult {
710        let tmp = TempDir::new()?;
711        let kind = IndexKind::Int64 {
712            bucket_width: NonZeroU64::new(1).unwrap(),
713        };
714        let mut table =
715            TimeSeriesTable::create(TableLocation::local(tmp.path()), integer_table_meta(kind))
716                .await?;
717        table.state.segments.insert(
718            "missing.parquet".to_string(),
719            indexed_segment("missing.parquet", 2i64.into(), 1i64.into()),
720        );
721
722        let equal = match table.scan_range(0i64, 0i64).await {
723            Err(error) => error,
724            Ok(_) => panic!("equal range must fail"),
725        };
726        let reversed = match table.scan_range(1i64, 0i64).await {
727            Err(error) => error,
728            Ok(_) => panic!("reversed range must fail"),
729        };
730        for error in [equal, reversed] {
731            assert!(matches!(
732                error,
733                TableError::InvalidRange {
734                    source: IndexValueError::InvalidRange { .. }
735                }
736            ));
737        }
738
739        assert!(matches!(
740            table.scan_range(0i64, 1u64).await,
741            Err(TableError::InvalidRange {
742                source: IndexValueError::KindMismatch {
743                    expected: "int64",
744                    actual: "uint64"
745                }
746            })
747        ));
748        let start = Utc.timestamp_opt(0, 0).single().unwrap();
749        let end = Utc.timestamp_opt(1, 0).single().unwrap();
750        assert!(matches!(
751            table.scan_range(start, end).await,
752            Err(TableError::InvalidRange {
753                source: IndexValueError::KindMismatch {
754                    expected: "int64",
755                    actual: "timestamp"
756                }
757            })
758        ));
759        Ok(())
760    }
761
762    #[tokio::test]
763    async fn integer_segment_stream_is_directly_polled() -> TestResult {
764        let tmp = TempDir::new()?;
765        let path = tmp.path().join("int64-stream.parquet");
766        write_index_parquet(
767            &path,
768            Arc::new(Int64Array::from(vec![Some(0), Some(1), Some(-1), None])),
769        )?;
770        let data = bytes::Bytes::from(std::fs::read(path)?);
771        let (reader, stats) = TrackingReader::new(data.clone());
772        let mut stream = segment_range_stream(
773            reader,
774            "data/int64-stream.parquet".to_string(),
775            "ts",
776            0i64,
777            2i64,
778        )
779        .await?;
780        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 0);
781
782        let batch = stream.next().await.transpose()?.expect("filtered batch");
783        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 1);
784        assert_eq!(
785            batch
786                .column(0)
787                .as_any()
788                .downcast_ref::<Int64Array>()
789                .expect("int64 index")
790                .iter()
791                .flatten()
792                .collect::<Vec<_>>(),
793            [0]
794        );
795        drop(stream);
796        assert!(stats.dropped.load(Ordering::SeqCst));
797
798        let (reader, stats, release) = TrackingReader::with_gate(data, 2);
799        let mut stream = segment_range_stream(
800            reader,
801            "data/int64-stream.parquet".to_string(),
802            "ts",
803            0i64,
804            2i64,
805        )
806        .await?;
807        assert_eq!(
808            stream
809                .next()
810                .await
811                .transpose()?
812                .expect("first integer batch")
813                .num_rows(),
814            1
815        );
816        let second = stream.next();
817        futures::pin_mut!(second);
818        assert!(futures::poll!(&mut second).is_pending());
819        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 2);
820        release.send(()).expect("release second row-group read");
821        assert_eq!(
822            second
823                .await
824                .transpose()?
825                .expect("second integer batch")
826                .num_rows(),
827            1
828        );
829        Ok(())
830    }
831
832    #[tokio::test]
833    async fn integer_scan_skips_non_candidates_and_stops_after_error() -> TestResult {
834        let tmp = TempDir::new()?;
835        let kind = IndexKind::Int64 {
836            bucket_width: NonZeroU64::new(1).unwrap(),
837        };
838        let mut table =
839            TimeSeriesTable::create(TableLocation::local(tmp.path()), integer_table_meta(kind))
840                .await?;
841        let wrong_type = "data/wrong-type.parquet";
842        write_index_parquet(
843            &tmp.path().join(wrong_type),
844            Arc::new(UInt64Array::from(vec![1, 2])),
845        )?;
846        for segment in [
847            indexed_segment(
848                "data/non-candidate.parquet",
849                (-10i64).into(),
850                (-5i64).into(),
851            ),
852            indexed_segment(wrong_type, 1i64.into(), 2i64.into()),
853            indexed_segment("data/later.parquet", 3i64.into(), 4i64.into()),
854        ] {
855            table.state.segments.insert(segment.path.clone(), segment);
856        }
857
858        let mut empty = table.scan_range(-20i64, -15i64).await?;
859        assert!(empty.next().await.is_none());
860
861        let mut failed = table.scan_range(0i64, 10i64).await?;
862        assert!(matches!(
863            failed.next().await,
864            Some(Err(TableError::IndexColumnTypeMismatch { .. }))
865        ));
866        assert!(failed.next().await.is_none());
867        Ok(())
868    }
869
870    #[tokio::test]
871    async fn scan_range_reports_registered_and_decoded_index_types() -> TestResult {
872        let tmp = TempDir::new()?;
873        let kind = IndexKind::Int64 {
874            bucket_width: NonZeroU64::new(1).unwrap(),
875        };
876        let mut table =
877            TimeSeriesTable::create(TableLocation::local(tmp.path()), integer_table_meta(kind))
878                .await?;
879        let rel = "data/wrong-index-type.parquet";
880        write_index_parquet(
881            &tmp.path().join(rel),
882            Arc::new(UInt64Array::from(vec![1, 2])),
883        )?;
884        table.state.segments.insert(
885            rel.to_string(),
886            indexed_segment(rel, 1i64.into(), 2i64.into()),
887        );
888
889        let mut stream = table.scan_range(0i64, 3i64).await?;
890        let error = stream
891            .next()
892            .await
893            .expect("segment error")
894            .expect_err("decoded type must match registered index");
895        assert!(matches!(
896            error,
897            TableError::IndexColumnTypeMismatch {
898                path,
899                column,
900                expected: "int64",
901                datatype: DataType::UInt64,
902            } if path == rel && column == "ts"
903        ));
904        assert!(stream.next().await.is_none());
905        Ok(())
906    }
907
908    fn write_multi_row_group_parquet(
909        path: &Path,
910        row_groups: &[&[i64]],
911    ) -> Result<(), Box<dyn std::error::Error>> {
912        let primary_timezone: Arc<str> = Arc::from("UTC");
913        let secondary_timezone: Arc<str> = Arc::from("+00:00");
914        let schema = Arc::new(Schema::new(vec![
915            Field::new(
916                "ts",
917                DataType::Timestamp(ArrowTimeUnit::Millisecond, Some(primary_timezone.clone())),
918                false,
919            ),
920            Field::new(
921                "observed_at",
922                DataType::Timestamp(ArrowTimeUnit::Microsecond, Some(secondary_timezone.clone())),
923                false,
924            ),
925        ]));
926
927        let file = File::create(path)?;
928        let mut writer = ArrowWriter::try_new(
929            file,
930            Arc::clone(&schema),
931            Some(WriterProperties::builder().build()),
932        )?;
933        for values in row_groups {
934            let ts = TimestampMillisecondArray::from(values.to_vec())
935                .with_timezone_opt(Some(primary_timezone.clone()));
936            let observed_at = TimestampMicrosecondArray::from(
937                values.iter().map(|value| value * 1_000).collect::<Vec<_>>(),
938            )
939            .with_timezone_opt(Some(secondary_timezone.clone()));
940            writer.write(&RecordBatch::try_new(
941                Arc::clone(&schema),
942                vec![Arc::new(ts), Arc::new(observed_at)],
943            )?)?;
944            writer.flush()?;
945        }
946        writer.close()?;
947        Ok(())
948    }
949
950    #[test]
951    fn timestamp_bounds_round_up_to_column_precision() {
952        let timestamp = |seconds, nanos| Utc.timestamp_opt(seconds, nanos).single().unwrap();
953        let cases = [
954            (
955                ArrowTimeUnit::Second,
956                timestamp(1, 500_000_000),
957                timestamp(2, 500_000_000),
958                (2, 3),
959            ),
960            (
961                ArrowTimeUnit::Second,
962                timestamp(-2, 500_000_000),
963                timestamp(-1, 500_000_000),
964                (-1, 0),
965            ),
966            (
967                ArrowTimeUnit::Second,
968                timestamp(1, 0),
969                timestamp(2, 0),
970                (1, 2),
971            ),
972            (
973                ArrowTimeUnit::Millisecond,
974                timestamp(1, 500_000),
975                timestamp(2, 500_000),
976                (1_001, 2_001),
977            ),
978            (
979                ArrowTimeUnit::Millisecond,
980                timestamp(-1, 999_500_000),
981                timestamp(0, 500_000),
982                (0, 1),
983            ),
984            (
985                ArrowTimeUnit::Millisecond,
986                timestamp(1, 0),
987                timestamp(2, 0),
988                (1_000, 2_000),
989            ),
990            (
991                ArrowTimeUnit::Microsecond,
992                timestamp(1, 500),
993                timestamp(2, 500),
994                (1_000_001, 2_000_001),
995            ),
996            (
997                ArrowTimeUnit::Microsecond,
998                timestamp(-1, 999_999_500),
999                timestamp(0, 500),
1000                (0, 1),
1001            ),
1002            (
1003                ArrowTimeUnit::Microsecond,
1004                timestamp(1, 0),
1005                timestamp(2, 0),
1006                (1_000_000, 2_000_000),
1007            ),
1008        ];
1009
1010        for (unit, start, end, expected) in cases {
1011            let field = Field::new("ts", DataType::Timestamp(unit, None), false);
1012            assert_eq!(
1013                to_bounds_i64(&field, "data/segment.parquet", "ts", start, end).unwrap(),
1014                expected
1015            );
1016        }
1017    }
1018
1019    #[tokio::test]
1020    async fn segment_stream_reads_on_demand_and_preserves_schema() -> TestResult {
1021        let tmp = TempDir::new()?;
1022        let path = tmp.path().join("multi-row-group.parquet");
1023        write_multi_row_group_parquet(&path, &[&[1_000], &[2_000], &[3_000]])?;
1024        let data = bytes::Bytes::from(std::fs::read(path)?);
1025        let file_size = data.len();
1026
1027        let (reader, stats) = TrackingReader::new(data.clone());
1028        let mut stream = segment_range_stream(
1029            reader,
1030            "data/multi-row-group.parquet".to_string(),
1031            "ts",
1032            Utc.timestamp_millis_opt(0).single().unwrap(),
1033            Utc.timestamp_millis_opt(4_000).single().unwrap(),
1034        )
1035        .await?;
1036
1037        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 0);
1038        let first = stream.next().await.transpose()?.expect("first batch");
1039        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 1);
1040        assert!(stats.bytes_read.load(Ordering::SeqCst) < file_size);
1041        assert_eq!(
1042            first.schema().field(0).data_type(),
1043            &DataType::Timestamp(ArrowTimeUnit::Millisecond, Some(Arc::<str>::from("UTC")))
1044        );
1045        assert_eq!(
1046            first.schema().field(1).data_type(),
1047            &DataType::Timestamp(ArrowTimeUnit::Microsecond, Some(Arc::<str>::from("+00:00")))
1048        );
1049
1050        tokio::task::yield_now().await;
1051        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 1);
1052        drop(stream);
1053        assert!(stats.dropped.load(Ordering::SeqCst));
1054        assert_eq!(stats.read_calls.load(Ordering::SeqCst), 1);
1055
1056        let (reader, gated_stats, release) = TrackingReader::with_gate(data, 2);
1057        let mut stream = segment_range_stream(
1058            reader,
1059            "data/multi-row-group.parquet".to_string(),
1060            "ts",
1061            Utc.timestamp_millis_opt(0).single().unwrap(),
1062            Utc.timestamp_millis_opt(4_000).single().unwrap(),
1063        )
1064        .await?;
1065        let mut timestamps = Vec::new();
1066        let first = stream.next().await.transpose()?.expect("first batch");
1067        timestamps.extend(
1068            first
1069                .column(0)
1070                .as_any()
1071                .downcast_ref::<TimestampMillisecondArray>()
1072                .expect("millisecond timestamp")
1073                .values()
1074                .iter()
1075                .copied(),
1076        );
1077        assert_eq!(gated_stats.read_calls.load(Ordering::SeqCst), 1);
1078
1079        let second = stream.next();
1080        futures::pin_mut!(second);
1081        assert!(futures::poll!(&mut second).is_pending());
1082        assert_eq!(gated_stats.read_calls.load(Ordering::SeqCst), 2);
1083        release.send(()).expect("release second row-group read");
1084        let second = second.await.transpose()?.expect("second batch");
1085        timestamps.extend(
1086            second
1087                .column(0)
1088                .as_any()
1089                .downcast_ref::<TimestampMillisecondArray>()
1090                .expect("millisecond timestamp")
1091                .values()
1092                .iter()
1093                .copied(),
1094        );
1095
1096        while let Some(batch) = stream.next().await.transpose()? {
1097            timestamps.extend(
1098                batch
1099                    .column(0)
1100                    .as_any()
1101                    .downcast_ref::<TimestampMillisecondArray>()
1102                    .expect("millisecond timestamp")
1103                    .values()
1104                    .iter()
1105                    .copied(),
1106            );
1107        }
1108        assert_eq!(timestamps, vec![1_000, 2_000, 3_000]);
1109        Ok(())
1110    }
1111
1112    #[tokio::test(flavor = "current_thread")]
1113    async fn fully_filtered_segment_stream_yields_and_cancels() -> TestResult {
1114        let tmp = TempDir::new()?;
1115        let path = tmp.path().join("filtered-row-groups.parquet");
1116        write_multi_row_group_parquet(&path, &[&[1_000], &[2_000], &[3_000]])?;
1117        let data = bytes::Bytes::from(std::fs::read(path)?);
1118        let (reader, stats, release) = TrackingReader::with_gate(data, 2);
1119
1120        let mut stream = segment_range_stream(
1121            reader,
1122            "data/filtered-row-groups.parquet".to_string(),
1123            "ts",
1124            Utc.timestamp_millis_opt(10_000).single().unwrap(),
1125            Utc.timestamp_millis_opt(20_000).single().unwrap(),
1126        )
1127        .await?;
1128
1129        {
1130            let next = stream.next();
1131            futures::pin_mut!(next);
1132            assert!(futures::poll!(&mut next).is_pending());
1133            assert_eq!(stats.read_calls.load(Ordering::SeqCst), 1);
1134            assert!(futures::poll!(&mut next).is_pending());
1135            assert_eq!(stats.read_calls.load(Ordering::SeqCst), 2);
1136        }
1137
1138        drop(stream);
1139        assert!(stats.dropped.load(Ordering::SeqCst));
1140        assert!(release.send(()).is_err());
1141        Ok(())
1142    }
1143
1144    #[tokio::test]
1145    async fn read_segment_range_errors_when_missing_time_column() -> TestResult {
1146        let tmp = TempDir::new()?;
1147        let location = TableLocation::local(tmp.path());
1148
1149        let rel = "data/no-ts.parquet";
1150        let path = tmp.path().join(rel);
1151        write_parquet_without_time_column(&path, &["A"], &[1.0])?;
1152
1153        let segment = SegmentMeta {
1154            path: rel.to_string(),
1155            format: FileFormat::Parquet,
1156            entity_layout: SegmentEntityLayout::NotApplicable,
1157            index_min: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1158            index_max: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1159            row_count: 1,
1160            file_size: None,
1161            coverage_path: None,
1162        };
1163
1164        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1165        let end = utc_datetime(2024, 1, 1, 0, 1, 0);
1166
1167        let err = match read_segment_range(&location, &segment, "ts", start, end).await {
1168            Err(err) => err,
1169            Ok(_) => panic!("missing ts column should error"),
1170        };
1171
1172        assert!(matches!(err, TableError::MissingIndexColumn { .. }));
1173        assert!(err.to_string().contains(rel));
1174        Ok(())
1175    }
1176
1177    #[tokio::test]
1178    async fn read_segment_range_errors_on_unsupported_time_type() -> TestResult {
1179        let tmp = TempDir::new()?;
1180        let location = TableLocation::local(tmp.path());
1181
1182        let rel = "data/int-ts.parquet";
1183        let path = tmp.path().join(rel);
1184        let ts_vals = [1_000_i64, 2_000];
1185        write_arrow_parquet_int_time(&path, &ts_vals, &["A", "B"], &[1.0, 2.0])?;
1186
1187        let segment = SegmentMeta {
1188            path: rel.to_string(),
1189            format: FileFormat::Parquet,
1190            entity_layout: SegmentEntityLayout::NotApplicable,
1191            index_min: (utc_datetime(2024, 1, 1, 0, 0, 1)).into(),
1192            index_max: (utc_datetime(2024, 1, 1, 0, 0, 2)).into(),
1193            row_count: ts_vals.len() as u64,
1194            file_size: None,
1195            coverage_path: None,
1196        };
1197
1198        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1199        let end = utc_datetime(2024, 1, 1, 0, 1, 0);
1200
1201        let err = match read_segment_range(&location, &segment, "ts", start, end).await {
1202            Err(err) => err,
1203            Ok(_) => panic!("unsupported time type should error"),
1204        };
1205
1206        assert!(matches!(err, TableError::IndexColumnTypeMismatch { .. }));
1207        assert!(err.to_string().contains(rel));
1208        Ok(())
1209    }
1210
1211    #[tokio::test]
1212    async fn read_segment_range_overflow_bounds_nanoseconds() -> TestResult {
1213        let tmp = TempDir::new()?;
1214        let location = TableLocation::local(tmp.path());
1215
1216        let rel = "data/nano-empty.parquet";
1217        let path = tmp.path().join(rel);
1218        write_arrow_parquet_with_unit(&path, ArrowTimeUnit::Nanosecond, &[], &[], &[])?;
1219
1220        let segment = SegmentMeta {
1221            path: rel.to_string(),
1222            format: FileFormat::Parquet,
1223            entity_layout: SegmentEntityLayout::NotApplicable,
1224            index_min: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1225            index_max: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1226            row_count: 0,
1227            file_size: None,
1228            coverage_path: None,
1229        };
1230
1231        let huge = Utc
1232            .timestamp_opt(9_223_372_037, 0)
1233            .single()
1234            .expect("overflow ts");
1235        let end = huge
1236            .checked_add_signed(chrono::Duration::seconds(1))
1237            .unwrap();
1238        let err = match read_segment_range(&location, &segment, "ts", huge, end).await {
1239            Err(err) => err,
1240            Ok(_) => panic!("overflow during bound conversion should error"),
1241        };
1242
1243        assert!(matches!(err, TableError::TimeConversionOverflow { .. }));
1244        Ok(())
1245    }
1246
1247    #[tokio::test]
1248    async fn scan_range_filters_across_segments() -> TestResult {
1249        let tmp = TempDir::new()?;
1250        let location = TableLocation::local(tmp.path());
1251        let meta = make_basic_table_meta();
1252        let mut table = TimeSeriesTable::create(location, meta).await?;
1253
1254        let rel1 = "data/seg-scan-1.parquet";
1255        let path1 = tmp.path().join(rel1);
1256        write_test_parquet(
1257            &path1,
1258            true,
1259            false,
1260            &[
1261                TestRow {
1262                    ts_millis: 1_000,
1263                    symbol: "A",
1264                    price: 10.0,
1265                },
1266                TestRow {
1267                    ts_millis: 2_000,
1268                    symbol: "A",
1269                    price: 20.0,
1270                },
1271            ],
1272        )?;
1273
1274        let rel2 = "data/seg-scan-2.parquet";
1275        let path2 = tmp.path().join(rel2);
1276        write_test_parquet(
1277            &path2,
1278            true,
1279            false,
1280            &[
1281                TestRow {
1282                    ts_millis: 61_000,
1283                    symbol: "A",
1284                    price: 30.0,
1285                },
1286                TestRow {
1287                    ts_millis: 62_000,
1288                    symbol: "A",
1289                    price: 40.0,
1290                },
1291            ],
1292        )?;
1293
1294        table.append_parquet_segment(rel1).await?;
1295        table.append_parquet_segment(rel2).await?;
1296
1297        // Query spans both segments but excludes the last row of the second segment.
1298        let start = Utc.timestamp_millis_opt(1_500).single().expect("valid ts");
1299        let end = Utc.timestamp_millis_opt(61_500).single().expect("valid ts");
1300
1301        let mut rows = collect_scan_rows(&table, start, end).await?;
1302        rows.sort_by_key(|row| row.0);
1303
1304        assert_eq!(
1305            rows,
1306            vec![
1307                (2_000, "A".to_string(), 20.0),
1308                (61_000, "A".to_string(), 30.0),
1309            ]
1310        );
1311
1312        Ok(())
1313    }
1314
1315    #[tokio::test]
1316    async fn scan_range_exclusive_end_and_empty() -> TestResult {
1317        let tmp = TempDir::new()?;
1318        let location = TableLocation::local(tmp.path());
1319        let meta = make_basic_table_meta();
1320        let mut table = TimeSeriesTable::create(location, meta).await?;
1321
1322        let rel = "data/seg-boundary.parquet";
1323        let path = tmp.path().join(rel);
1324        write_test_parquet(
1325            &path,
1326            true,
1327            false,
1328            &[
1329                TestRow {
1330                    ts_millis: 1_000,
1331                    symbol: "A",
1332                    price: 10.0,
1333                },
1334                TestRow {
1335                    ts_millis: 2_000,
1336                    symbol: "A",
1337                    price: 20.0,
1338                },
1339            ],
1340        )?;
1341
1342        table.append_parquet_segment(rel).await?;
1343
1344        let start = Utc.timestamp_millis_opt(1_000).single().expect("valid ts");
1345        let end = Utc.timestamp_millis_opt(2_000).single().expect("valid ts");
1346        let rows = collect_scan_rows(&table, start, end).await?;
1347        assert_eq!(rows, vec![(1_000, "A".to_string(), 10.0)]);
1348
1349        let empty_start = Utc.timestamp_millis_opt(5_000).single().expect("valid ts");
1350        let empty_end = Utc.timestamp_millis_opt(6_000).single().expect("valid ts");
1351        let rows = collect_scan_rows(&table, empty_start, empty_end).await?;
1352        assert!(rows.is_empty());
1353
1354        Ok(())
1355    }
1356
1357    #[tokio::test]
1358    async fn scan_range_rejects_invalid_range() -> TestResult {
1359        let tmp = TempDir::new()?;
1360        let location = TableLocation::local(tmp.path());
1361        let meta = make_basic_table_meta();
1362        let table = TimeSeriesTable::create(location, meta).await?;
1363
1364        let start = Utc.timestamp_millis_opt(1_000).single().expect("valid ts");
1365        let end = start;
1366
1367        let result = table.scan_range(start, end).await;
1368
1369        assert!(matches!(result, Err(TableError::InvalidRange { .. })));
1370        Ok(())
1371    }
1372
1373    #[tokio::test]
1374    async fn scan_range_supports_second_unit() -> TestResult {
1375        let tmp = TempDir::new()?;
1376        let location = TableLocation::local(tmp.path());
1377        let meta = make_basic_table_meta();
1378        let mut table = TimeSeriesTable::create(location, meta).await?;
1379
1380        let rel = "data/seg-seconds.parquet";
1381        let path = tmp.path().join(rel);
1382        write_arrow_parquet_with_unit(
1383            &path,
1384            ArrowTimeUnit::Second,
1385            &[Some(1), Some(2), Some(3)],
1386            &["A", "A", "A"],
1387            &[1.0, 2.0, 3.0],
1388        )?;
1389        let segment = SegmentMeta {
1390            path: rel.to_string(),
1391            format: FileFormat::Parquet,
1392            entity_layout: SegmentEntityLayout::NotApplicable,
1393            index_min: (Utc.timestamp_opt(1, 0).single().unwrap()).into(),
1394            index_max: (Utc.timestamp_opt(3, 0).single().unwrap()).into(),
1395            row_count: 3,
1396            file_size: None,
1397            coverage_path: None,
1398        };
1399        table.state.segments.insert(segment.path.clone(), segment);
1400
1401        let start = Utc.timestamp_millis_opt(1_500).single().unwrap();
1402        let end = Utc.timestamp_millis_opt(2_500).single().unwrap();
1403        let rows = collect_scan_rows(&table, start, end).await?;
1404
1405        assert_eq!(rows, vec![(2, "A".to_string(), 2.0)]);
1406        Ok(())
1407    }
1408
1409    #[tokio::test]
1410    async fn scan_range_supports_microsecond_unit() -> TestResult {
1411        let tmp = TempDir::new()?;
1412        let location = TableLocation::local(tmp.path());
1413        let meta = make_table_meta_with_unit(LogicalTimestampUnit::Micros);
1414        let mut table = TimeSeriesTable::create(location, meta).await?;
1415
1416        let rel = "data/seg-micros.parquet";
1417        let path = tmp.path().join(rel);
1418        write_arrow_parquet_with_unit(
1419            &path,
1420            ArrowTimeUnit::Microsecond,
1421            &[Some(1_000_000), Some(2_000_000), Some(3_000_000)],
1422            &["A", "A", "A"],
1423            &[1.0, 2.0, 3.0],
1424        )?;
1425
1426        table.append_parquet_segment(rel).await?;
1427
1428        let start = Utc
1429            .timestamp_opt(1, 500_000_000)
1430            .single()
1431            .expect("valid start");
1432        let end = Utc
1433            .timestamp_opt(2, 500_000_000)
1434            .single()
1435            .expect("valid end");
1436        let rows = collect_scan_rows(&table, start, end).await?;
1437
1438        assert_eq!(rows, vec![(2_000_000, "A".to_string(), 2.0)]);
1439        Ok(())
1440    }
1441
1442    #[tokio::test]
1443    async fn scan_range_supports_nanosecond_unit() -> TestResult {
1444        let tmp = TempDir::new()?;
1445        let location = TableLocation::local(tmp.path());
1446        let meta = make_table_meta_with_unit(LogicalTimestampUnit::Nanos);
1447        let mut table = TimeSeriesTable::create(location, meta).await?;
1448
1449        let rel = "data/seg-nanos.parquet";
1450        let path = tmp.path().join(rel);
1451        write_arrow_parquet_with_unit(
1452            &path,
1453            ArrowTimeUnit::Nanosecond,
1454            &[
1455                Some(1_000_000_000),
1456                Some(1_500_000_000),
1457                Some(2_000_000_000),
1458            ],
1459            &["A", "A", "A"],
1460            &[1.0, 2.0, 3.0],
1461        )?;
1462
1463        table.append_parquet_segment(rel).await?;
1464
1465        let start = Utc
1466            .timestamp_opt(1, 250_000_000)
1467            .single()
1468            .expect("valid start");
1469        let end = Utc
1470            .timestamp_opt(1, 750_000_000)
1471            .single()
1472            .expect("valid end");
1473        let rows = collect_scan_rows(&table, start, end).await?;
1474
1475        assert_eq!(rows, vec![(1_500_000_000, "A".to_string(), 2.0)]);
1476        Ok(())
1477    }
1478
1479    #[tokio::test]
1480    async fn scan_range_filters_null_timestamps() -> TestResult {
1481        let tmp = TempDir::new()?;
1482        let location = TableLocation::local(tmp.path());
1483        let meta = make_table_meta_with_unit(LogicalTimestampUnit::Millis);
1484        let mut table = TimeSeriesTable::create(location, meta).await?;
1485
1486        let rel = "data/seg-null-ts.parquet";
1487        let path = tmp.path().join(rel);
1488        write_arrow_parquet_with_unit(
1489            &path,
1490            ArrowTimeUnit::Millisecond,
1491            &[Some(1_000), None, Some(2_000)],
1492            &["A", "A", "A"],
1493            &[1.0, 2.0, 3.0],
1494        )?;
1495
1496        table.append_parquet_segment(rel).await?;
1497
1498        let start = Utc.timestamp_millis_opt(500).single().unwrap();
1499        let end = Utc.timestamp_millis_opt(2_500).single().unwrap();
1500        let rows = collect_scan_rows(&table, start, end).await?;
1501
1502        assert_eq!(
1503            rows,
1504            vec![(1_000, "A".to_string(), 1.0), (2_000, "A".to_string(), 3.0)]
1505        );
1506        Ok(())
1507    }
1508
1509    #[tokio::test]
1510    async fn scan_range_empty_when_no_segments() -> TestResult {
1511        let tmp = TempDir::new()?;
1512        let location = TableLocation::local(tmp.path());
1513        let meta = make_basic_table_meta();
1514        let table = TimeSeriesTable::create(location, meta).await?;
1515
1516        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1517        let end = utc_datetime(2024, 1, 1, 0, 1, 0);
1518
1519        let mut stream = table.scan_range(start, end).await?;
1520        assert!(stream.next().await.is_none());
1521        Ok(())
1522    }
1523
1524    #[tokio::test]
1525    async fn scan_range_empty_for_zero_row_segment() -> TestResult {
1526        let tmp = TempDir::new()?;
1527        let location = TableLocation::local(tmp.path());
1528        let meta = make_basic_table_meta();
1529        let mut table = TimeSeriesTable::create(location, meta).await?;
1530
1531        let rel = "data/seg-empty.parquet";
1532        let path = tmp.path().join(rel);
1533        write_arrow_parquet_with_unit(&path, ArrowTimeUnit::Millisecond, &[], &[], &[])?;
1534
1535        let segment = SegmentMeta {
1536            path: rel.to_string(),
1537            format: FileFormat::Parquet,
1538            entity_layout: SegmentEntityLayout::NotApplicable,
1539            index_min: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1540            index_max: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1541            row_count: 0,
1542            file_size: None,
1543            coverage_path: None,
1544        };
1545
1546        table.state.segments.insert(segment.path.clone(), segment);
1547
1548        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1549        let end = utc_datetime(2024, 1, 1, 0, 1, 0);
1550
1551        let mut stream = table.scan_range(start, end).await?;
1552        assert!(stream.next().await.is_none());
1553        Ok(())
1554    }
1555
1556    #[tokio::test]
1557    async fn scan_range_all_null_time_filtered_out() -> TestResult {
1558        let tmp = TempDir::new()?;
1559        let location = TableLocation::local(tmp.path());
1560        let meta = make_table_meta_with_unit(LogicalTimestampUnit::Millis);
1561        let mut table = TimeSeriesTable::create(location, meta).await?;
1562
1563        let rel = "data/seg-null-only.parquet";
1564        let path = tmp.path().join(rel);
1565        write_arrow_parquet_with_unit(
1566            &path,
1567            ArrowTimeUnit::Millisecond,
1568            &[None, None],
1569            &["A", "B"],
1570            &[1.0, 2.0],
1571        )?;
1572
1573        let segment = SegmentMeta {
1574            path: rel.to_string(),
1575            format: FileFormat::Parquet,
1576            entity_layout: SegmentEntityLayout::NotApplicable,
1577            index_min: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1578            index_max: (utc_datetime(2024, 1, 1, 0, 0, 1)).into(),
1579            row_count: 2,
1580            file_size: None,
1581            coverage_path: None,
1582        };
1583
1584        table.state.segments.insert(segment.path.clone(), segment);
1585
1586        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1587        let end = utc_datetime(2024, 1, 1, 0, 0, 5);
1588
1589        let mut stream = table.scan_range(start, end).await?;
1590        assert!(stream.next().await.is_none());
1591        Ok(())
1592    }
1593
1594    #[tokio::test]
1595    async fn scan_range_errors_on_missing_time_column_in_segment() -> TestResult {
1596        let tmp = TempDir::new()?;
1597        let location = TableLocation::local(tmp.path());
1598        let meta = make_basic_table_meta();
1599        let mut table = TimeSeriesTable::create(location, meta).await?;
1600
1601        let rel = "data/seg-scan-no-ts.parquet";
1602        let path = tmp.path().join(rel);
1603        write_parquet_without_time_column(&path, &["A"], &[1.0])?;
1604
1605        let segment = SegmentMeta {
1606            path: rel.to_string(),
1607            format: FileFormat::Parquet,
1608            entity_layout: SegmentEntityLayout::NotApplicable,
1609            index_min: (utc_datetime(2024, 1, 1, 0, 0, 0)).into(),
1610            index_max: (utc_datetime(2024, 1, 1, 0, 1, 0)).into(),
1611            row_count: 1,
1612            file_size: None,
1613            coverage_path: None,
1614        };
1615
1616        table.state.segments.insert(segment.path.clone(), segment);
1617        let unopened = SegmentMeta {
1618            path: "data/should-not-open.parquet".to_string(),
1619            format: FileFormat::Parquet,
1620            entity_layout: SegmentEntityLayout::NotApplicable,
1621            index_min: (utc_datetime(2024, 1, 1, 0, 1, 30)).into(),
1622            index_max: (utc_datetime(2024, 1, 1, 0, 1, 31)).into(),
1623            row_count: 1,
1624            file_size: None,
1625            coverage_path: None,
1626        };
1627        table.state.segments.insert(unopened.path.clone(), unopened);
1628
1629        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1630        let end = utc_datetime(2024, 1, 1, 0, 2, 0);
1631
1632        let mut stream = table.scan_range(start, end).await?;
1633        let err = stream.next().await.expect("expected error from scan");
1634
1635        assert!(matches!(err, Err(TableError::MissingIndexColumn { .. })));
1636        assert!(stream.next().await.is_none());
1637        Ok(())
1638    }
1639
1640    #[tokio::test]
1641    async fn scan_range_errors_on_unsupported_time_type_segment() -> TestResult {
1642        let tmp = TempDir::new()?;
1643        let location = TableLocation::local(tmp.path());
1644        let meta = make_basic_table_meta();
1645        let mut table = TimeSeriesTable::create(location, meta).await?;
1646
1647        let rel = "data/seg-scan-int-ts.parquet";
1648        let path = tmp.path().join(rel);
1649        write_arrow_parquet_int_time(&path, &[1_000], &["A"], &[1.0])?;
1650
1651        let segment = SegmentMeta {
1652            path: rel.to_string(),
1653            format: FileFormat::Parquet,
1654            entity_layout: SegmentEntityLayout::NotApplicable,
1655            index_min: (utc_datetime(2024, 1, 1, 0, 0, 1)).into(),
1656            index_max: (utc_datetime(2024, 1, 1, 0, 0, 1)).into(),
1657            row_count: 1,
1658            file_size: None,
1659            coverage_path: None,
1660        };
1661
1662        table.state.segments.insert(segment.path.clone(), segment);
1663
1664        let start = utc_datetime(2024, 1, 1, 0, 0, 0);
1665        let end = utc_datetime(2024, 1, 1, 0, 1, 0);
1666
1667        let mut stream = table.scan_range(start, end).await?;
1668        let err = stream.next().await.expect("expected error from scan");
1669
1670        assert!(matches!(
1671            err,
1672            Err(TableError::IndexColumnTypeMismatch { .. })
1673        ));
1674        Ok(())
1675    }
1676
1677    #[tokio::test]
1678    async fn scan_range_reads_segments_independent_of_append_order() -> TestResult {
1679        let tmp = TempDir::new()?;
1680        let location = TableLocation::local(tmp.path());
1681        let meta = make_basic_table_meta();
1682        let mut table = TimeSeriesTable::create(location, meta).await?;
1683
1684        let rel_b = "data/seg-overlap-b.parquet";
1685        let path_b = tmp.path().join(rel_b);
1686        write_test_parquet(
1687            &path_b,
1688            true,
1689            false,
1690            &[TestRow {
1691                ts_millis: 120_000,
1692                symbol: "A",
1693                price: 2.0,
1694            }],
1695        )?;
1696
1697        let rel_a = "data/seg-overlap-a.parquet";
1698        let path_a = tmp.path().join(rel_a);
1699        write_test_parquet(
1700            &path_a,
1701            true,
1702            false,
1703            &[TestRow {
1704                ts_millis: 60_000,
1705                symbol: "A",
1706                price: 1.0,
1707            }],
1708        )?;
1709
1710        // Append in reverse index order to exercise segment discovery.
1711        table.append_parquet_segment(rel_b).await?;
1712        table.append_parquet_segment(rel_a).await?;
1713
1714        let start = Utc.timestamp_millis_opt(50_000).single().unwrap();
1715        let end = Utc.timestamp_millis_opt(150_000).single().unwrap();
1716        let mut rows = collect_scan_rows(&table, start, end).await?;
1717        rows.sort_by_key(|row| row.0);
1718
1719        assert_eq!(
1720            rows,
1721            vec![
1722                (60_000, "A".to_string(), 1.0),
1723                (120_000, "A".to_string(), 2.0)
1724            ]
1725        );
1726        Ok(())
1727    }
1728
1729    #[tokio::test]
1730    async fn scan_range_skips_non_overlapping_segments() -> TestResult {
1731        let tmp = TempDir::new()?;
1732        let location = TableLocation::local(tmp.path());
1733        let meta = make_basic_table_meta();
1734        let mut table = TimeSeriesTable::create(location, meta).await?;
1735
1736        let rel1 = "data/seg-early.parquet";
1737        let path1 = tmp.path().join(rel1);
1738        write_test_parquet(
1739            &path1,
1740            true,
1741            false,
1742            &[TestRow {
1743                ts_millis: 1_000,
1744                symbol: "A",
1745                price: 1.0,
1746            }],
1747        )?;
1748
1749        let rel2 = "data/seg-late.parquet";
1750        let path2 = tmp.path().join(rel2);
1751        write_test_parquet(
1752            &path2,
1753            true,
1754            false,
1755            &[TestRow {
1756                ts_millis: 70_000,
1757                symbol: "A",
1758                price: 9.0,
1759            }],
1760        )?;
1761
1762        table.append_parquet_segment(rel1).await?;
1763        table.append_parquet_segment(rel2).await?;
1764
1765        let start = Utc.timestamp_millis_opt(1_500).single().unwrap();
1766        let end = Utc.timestamp_millis_opt(2_000).single().unwrap();
1767        let rows = collect_scan_rows(&table, start, end).await?;
1768
1769        assert_eq!(rows, Vec::new());
1770        Ok(())
1771    }
1772
1773    #[tokio::test]
1774    async fn scan_range_reports_missing_segment_path() -> TestResult {
1775        let tmp = TempDir::new()?;
1776        let location = TableLocation::local(tmp.path());
1777        let meta = make_basic_table_meta();
1778        let mut table = TimeSeriesTable::create(location, meta).await?;
1779        let rel = "data/missing.parquet";
1780        let segment = SegmentMeta {
1781            path: rel.to_string(),
1782            format: FileFormat::Parquet,
1783            entity_layout: SegmentEntityLayout::NotApplicable,
1784            index_min: (Utc.timestamp_millis_opt(1_000).single().unwrap()).into(),
1785            index_max: (Utc.timestamp_millis_opt(2_000).single().unwrap()).into(),
1786            row_count: 1,
1787            file_size: None,
1788            coverage_path: None,
1789        };
1790        table.state.segments.insert(segment.path.clone(), segment);
1791
1792        let mut stream = table
1793            .scan_range(
1794                Utc.timestamp_millis_opt(0).single().unwrap(),
1795                Utc.timestamp_millis_opt(3_000).single().unwrap(),
1796            )
1797            .await?;
1798        let error = stream
1799            .next()
1800            .await
1801            .expect("missing segment error")
1802            .expect_err("missing segment should fail");
1803
1804        assert!(matches!(error, TableError::Storage { .. }));
1805        assert!(error.to_string().contains(rel));
1806        assert!(stream.next().await.is_none());
1807        Ok(())
1808    }
1809
1810    #[tokio::test]
1811    async fn scan_range_propagates_parquet_read_error() -> TestResult {
1812        let tmp = TempDir::new()?;
1813        let location = TableLocation::local(tmp.path());
1814        let meta = make_basic_table_meta();
1815        let mut table = TimeSeriesTable::create(location.clone(), meta).await?;
1816
1817        let rel = "data/seg-corrupt.parquet";
1818        let path = tmp.path().join(rel);
1819        write_test_parquet(
1820            &path,
1821            true,
1822            false,
1823            &[TestRow {
1824                ts_millis: 1_000,
1825                symbol: "A",
1826                price: 1.0,
1827            }],
1828        )?;
1829
1830        table.append_parquet_segment(rel).await?;
1831
1832        // Corrupt the file after append so scan encounters a read failure.
1833        let f = std::fs::OpenOptions::new().write(true).open(&path)?;
1834        f.set_len(4)?;
1835
1836        let start = Utc.timestamp_millis_opt(0).single().unwrap();
1837        let end = Utc.timestamp_millis_opt(2_000).single().unwrap();
1838
1839        let mut stream = table.scan_range(start, end).await?;
1840        let err = stream
1841            .next()
1842            .await
1843            .expect("first item should be error")
1844            .expect_err("corrupt segment should fail");
1845
1846        assert!(matches!(err, TableError::ParquetRead { .. }));
1847        assert!(err.to_string().contains(rel));
1848        Ok(())
1849    }
1850}