Skip to main content

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