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