1use 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#[derive(Debug, Snafu)]
50#[snafu(visibility(pub(crate)))]
51#[non_exhaustive]
52pub enum ScanError {
53 #[snafu(display("Invalid scan range: {source}"))]
55 InvalidRange {
56 source: IndexValueError,
58 backtrace: Box<Backtrace>,
60 },
61
62 #[snafu(display("Invalid persisted segment bounds while planning scan: {source}"))]
64 InvalidSegmentBounds {
65 source: IndexValueError,
67 backtrace: Box<Backtrace>,
69 },
70
71 #[snafu(display("Failed to open segment {path} during scan execution: {source}"))]
73 Storage {
74 path: String,
76 #[snafu(source(from(storage::StorageError, Box::new)), backtrace)]
78 source: Box<storage::StorageError>,
79 },
80
81 #[snafu(display(
83 "Parquet error while {operation} for segment {path} during scan execution: {source}"
84 ))]
85 Parquet {
86 path: String,
88 operation: &'static str,
90 #[snafu(source(from(ParquetError, Box::new)))]
92 source: Box<ParquetError>,
93 backtrace: Box<Backtrace>,
95 },
96
97 #[snafu(display(
99 "Arrow error while {operation} for column {column} in segment {path} during scan execution: {source}"
100 ))]
101 Arrow {
102 path: String,
104 column: String,
106 operation: &'static str,
108 #[snafu(source(from(arrow::error::ArrowError, Box::new)))]
110 source: Box<arrow::error::ArrowError>,
111 backtrace: Box<Backtrace>,
113 },
114
115 #[snafu(display(
117 "Missing ordered-index column {column} in segment {path} during scan execution"
118 ))]
119 MissingIndexColumn {
120 path: String,
122 column: String,
124 backtrace: Box<Backtrace>,
126 },
127
128 #[snafu(display(
130 "Ordered-index column {column} in segment {path} has Arrow type {datatype:?}, expected {expected}, during scan execution"
131 ))]
132 IndexColumnTypeMismatch {
133 path: String,
135 column: String,
137 expected: &'static str,
139 datatype: Box<DataType>,
141 backtrace: Box<Backtrace>,
143 },
144
145 #[snafu(display(
147 "Timestamp conversion overflow for column {column} in segment {path} during scan execution (value: {timestamp})"
148 ))]
149 TimeConversionOverflow {
150 path: String,
152 column: String,
154 timestamp: DateTime<Utc>,
156 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
175fn 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, <_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 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 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 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 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 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_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 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}