1use std::path::Path;
14
15use arrow::datatypes::{DataType, TimeUnit};
16use arrow_array::{
17 Array, Int64Array, TimestampMicrosecondArray, TimestampMillisecondArray,
18 TimestampNanosecondArray, TimestampSecondArray, UInt64Array,
19};
20use chrono::{TimeZone, Utc};
21use futures::{Stream, StreamExt};
22use parquet::{
23 arrow::{
24 ProjectionMask,
25 arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions},
26 async_reader::ParquetRecordBatchStreamBuilder,
27 },
28 errors::ParquetError,
29};
30use roaring::RoaringTreemap;
31use snafu::{Backtrace, Snafu};
32use tokio::task::JoinSet;
33
34use crate::{
35 coverage::index_interval::{
36 IndexInterval, IndexIntervalMappingError, index_interval_for_id,
37 index_interval_id_for_value,
38 },
39 coverage::{Coverage, EntityIdentity, EntityIdentityError, IndexIntervalId},
40 metadata::{
41 index::{IndexKind, IndexSpec, IndexValue},
42 segments::ParquetIndexColumnError,
43 },
44 storage::{StorageError, TableLocation, open_parquet_reader},
45};
46
47use super::schema::validate_parquet_index;
48use super::{INSPECTION_BATCH_SIZE, resolve_rg_settings};
49
50#[derive(Debug, Snafu)]
62#[non_exhaustive]
63pub enum SegmentCoverageError {
64 #[snafu(display("Storage error reading parquet file {path}: {source}"))]
68 Storage {
69 path: String,
71 #[snafu(source, backtrace)]
73 source: StorageError,
74 },
75
76 #[snafu(display("Parquet read error for {path}: {source}"))]
81 ParquetRead {
82 path: String,
84 #[snafu(source)]
86 source: ParquetError,
87 backtrace: Backtrace,
89 },
90
91 #[snafu(display("Row-group scan task failed for segment at {path}: {source}"))]
93 RowGroupTask {
94 path: String,
96 source: tokio::task::JoinError,
98 backtrace: Backtrace,
100 },
101
102 #[snafu(transparent)]
104 OrderedIndexColumn {
105 source: ParquetIndexColumnError,
107 },
108
109 #[snafu(display(
111 "Invalid {expected_domain} value for ordered-index column {column} in segment at {path}: {detail}"
112 ))]
113 IndexValue {
114 path: String,
116 column: String,
118 expected_domain: &'static str,
120 detail: String,
122 },
123
124 #[snafu(display("Index interval mapping failed for segment {path}: {source}"))]
126 IndexIntervalMapping {
127 path: String,
129 source: IndexIntervalMappingError,
131 },
132
133 #[snafu(display(
135 "Duplicate ordered-index interval {example_index_interval} in segment {path}"
136 ))]
137 DuplicateIndexInterval {
138 path: String,
140 example_identity: Option<EntityIdentity>,
142 example_index_interval: IndexInterval,
144 },
145
146 #[snafu(display("Entity column not found in {path}: {column}"))]
148 EntityColumnNotFound {
149 path: String,
151 column: String,
153 },
154
155 #[snafu(display("Unsupported entity column type in {path}: {column} has {datatype}"))]
157 EntityColumnUnsupportedType {
158 path: String,
160 column: String,
162 datatype: String,
164 },
165
166 #[snafu(display("Entity column contains nulls in {path}: {column}"))]
168 EntityColumnHasNull {
169 path: String,
171 column: String,
173 },
174
175 #[snafu(display("Entity column has no values (empty segment) in {path}: {column}"))]
177 EntityColumnEmpty {
178 path: String,
180 column: String,
182 },
183
184 #[snafu(display("Invalid entity identity in segment {path}: {source}"))]
186 EntityIdentity {
187 path: String,
189 source: EntityIdentityError,
191 },
192}
193
194pub(super) fn arrow_index_error(
195 path: &str,
196 index: &IndexSpec,
197 observed_type: String,
198) -> SegmentCoverageError {
199 SegmentCoverageError::OrderedIndexColumn {
200 source: ParquetIndexColumnError {
201 path: path.to_string(),
202 column: index.column.clone(),
203 expected_domain: index.kind.name(),
204 observed_type,
205 },
206 }
207}
208
209pub(super) fn timestamp_value(
210 path: &str,
211 index: &IndexSpec,
212 unit: TimeUnit,
213 raw: i64,
214) -> Result<IndexValue, SegmentCoverageError> {
215 let value = match unit {
216 TimeUnit::Second => Utc.timestamp_opt(raw, 0),
217 TimeUnit::Millisecond => Utc.timestamp_millis_opt(raw),
218 TimeUnit::Microsecond => Utc.timestamp_micros(raw),
219 TimeUnit::Nanosecond => {
220 let seconds = raw.div_euclid(1_000_000_000);
221 let nanos = raw.rem_euclid(1_000_000_000) as u32;
222 Utc.timestamp_opt(seconds, nanos)
223 }
224 };
225 value
226 .single()
227 .map(IndexValue::Timestamp)
228 .ok_or_else(|| SegmentCoverageError::IndexValue {
229 path: path.to_string(),
230 column: index.column.clone(),
231 expected_domain: index.kind.name(),
232 detail: format!("timestamp value {raw} is out of range for {unit:?}"),
233 })
234}
235
236pub(super) fn map_and_insert_index_interval_id(
237 bitmap: &mut RoaringTreemap,
238 path: &str,
239 index: &IndexSpec,
240 value: IndexValue,
241) -> Result<(IndexIntervalId, bool), SegmentCoverageError> {
242 let index_interval_id = index_interval_id_for_value(&index.kind, &value).map_err(|source| {
243 SegmentCoverageError::IndexIntervalMapping {
244 path: path.to_string(),
245 source,
246 }
247 })?;
248 Ok((index_interval_id, bitmap.insert(index_interval_id)))
249}
250
251pub(super) fn duplicate_index_interval_error(
252 path: &str,
253 index: &IndexSpec,
254 identity: Option<&EntityIdentity>,
255 index_interval_id: IndexIntervalId,
256) -> SegmentCoverageError {
257 match index_interval_for_id(&index.kind, index_interval_id) {
258 Ok(example_index_interval) => SegmentCoverageError::DuplicateIndexInterval {
259 path: path.to_string(),
260 example_identity: identity.cloned(),
261 example_index_interval,
262 },
263 Err(source) => SegmentCoverageError::IndexIntervalMapping {
264 path: path.to_string(),
265 source,
266 },
267 }
268}
269
270fn add_array_index_interval_ids<T, F>(
271 bitmap: &mut RoaringTreemap,
272 path: &str,
273 index: &IndexSpec,
274 array: &arrow_array::PrimitiveArray<T>,
275 mut to_value: F,
276) -> Result<(), SegmentCoverageError>
277where
278 T: arrow_array::types::ArrowPrimitiveType,
279 F: FnMut(T::Native) -> Result<IndexValue, SegmentCoverageError>,
280{
281 if array.null_count() == 0 {
282 for &raw in array.values() {
283 let (index_interval_id, inserted) =
284 map_and_insert_index_interval_id(bitmap, path, index, to_value(raw)?)?;
285 if !inserted {
286 return Err(duplicate_index_interval_error(
287 path,
288 index,
289 None,
290 index_interval_id,
291 ));
292 }
293 }
294 } else {
295 for raw in array.iter().flatten() {
296 let (index_interval_id, inserted) =
297 map_and_insert_index_interval_id(bitmap, path, index, to_value(raw)?)?;
298 if !inserted {
299 return Err(duplicate_index_interval_error(
300 path,
301 index,
302 None,
303 index_interval_id,
304 ));
305 }
306 }
307 }
308 Ok(())
309}
310
311async fn compute_coverage_bitmap_from_stream(
312 mut reader: impl Stream<
313 Item = Result<arrow::record_batch::RecordBatch, parquet::errors::ParquetError>,
314 > + Unpin,
315 path_str: &str,
316 index: &IndexSpec,
317) -> Result<RoaringTreemap, SegmentCoverageError> {
318 let mut bitmap = RoaringTreemap::new();
319
320 while let Some(batch_res) = reader.next().await {
321 let batch = batch_res.map_err(|source| SegmentCoverageError::ParquetRead {
322 path: path_str.to_string(),
323 source,
324 backtrace: Backtrace::capture(),
325 })?;
326
327 let col = batch.column(0);
328
329 match (&index.kind, col.data_type()) {
330 (IndexKind::Timestamp { .. }, DataType::Timestamp(unit, _)) => {
331 macro_rules! process_timestamp_array {
332 ($array_type:ty) => {{
333 let array =
334 col.as_any().downcast_ref::<$array_type>().ok_or_else(|| {
335 arrow_index_error(
336 path_str,
337 index,
338 format!("Arrow {}", col.data_type()),
339 )
340 })?;
341 add_array_index_interval_ids(&mut bitmap, path_str, index, array, |raw| {
342 timestamp_value(path_str, index, unit.clone(), raw)
343 })?;
344 }};
345 }
346 match unit {
347 TimeUnit::Second => process_timestamp_array!(TimestampSecondArray),
348 TimeUnit::Millisecond => {
349 process_timestamp_array!(TimestampMillisecondArray)
350 }
351 TimeUnit::Microsecond => {
352 process_timestamp_array!(TimestampMicrosecondArray)
353 }
354 TimeUnit::Nanosecond => process_timestamp_array!(TimestampNanosecondArray),
355 }
356 }
357 (IndexKind::Int64 { .. }, DataType::Int64) => {
358 let array = col.as_any().downcast_ref::<Int64Array>().ok_or_else(|| {
359 arrow_index_error(path_str, index, format!("Arrow {}", col.data_type()))
360 })?;
361 add_array_index_interval_ids(&mut bitmap, path_str, index, array, |raw| {
362 Ok(IndexValue::Int64(raw))
363 })?;
364 }
365 (IndexKind::UInt64 { .. }, DataType::UInt64) => {
366 let array = col.as_any().downcast_ref::<UInt64Array>().ok_or_else(|| {
367 arrow_index_error(path_str, index, format!("Arrow {}", col.data_type()))
368 })?;
369 add_array_index_interval_ids(&mut bitmap, path_str, index, array, |raw| {
370 Ok(IndexValue::UInt64(raw))
371 })?;
372 }
373 other => {
374 return Err(arrow_index_error(
375 path_str,
376 index,
377 format!("Arrow {other:?}"),
378 ));
379 }
380 }
381
382 tokio::task::yield_now().await;
383 }
384
385 Ok(bitmap)
386}
387
388pub async fn compute_segment_coverage(
407 location: &TableLocation,
408 rel_path: &Path,
409 index: &IndexSpec,
410) -> Result<Coverage, SegmentCoverageError> {
411 let path = rel_path.display().to_string();
412 let mut file = open_parquet_reader(location.as_ref(), rel_path)
413 .await
414 .map_err(|source| SegmentCoverageError::Storage {
415 path: path.clone(),
416 source,
417 })?;
418 let metadata = ArrowReaderMetadata::load_async(&mut file, ArrowReaderOptions::default())
419 .await
420 .map_err(|source| SegmentCoverageError::ParquetRead {
421 path: path.clone(),
422 source,
423 backtrace: Backtrace::capture(),
424 })?;
425 validate_parquet_index(&path, metadata.parquet_schema(), index)
426 .map_err(|source| SegmentCoverageError::OrderedIndexColumn { source })?;
427 drop(file);
428
429 let mask = ProjectionMask::columns(metadata.parquet_schema(), [index.column.as_str()]);
430 let row_groups = metadata.metadata().num_row_groups();
431 let (max_tasks, row_groups_per_task) = resolve_rg_settings(row_groups);
432 let row_groups = (0..row_groups).collect::<Vec<_>>();
433 let chunks = row_groups
434 .chunks(row_groups_per_task)
435 .map(<[usize]>::to_vec)
436 .collect::<Vec<_>>();
437 debug_assert!(chunks.len() <= max_tasks);
438
439 let mut tasks = JoinSet::new();
440 for chunk in chunks {
441 let location = location.clone();
442 let rel_path = rel_path.to_path_buf();
443 let path = path.clone();
444 let index = index.clone();
445 let metadata = metadata.clone();
446 let mask = mask.clone();
447
448 tasks.spawn(async move {
449 let file = open_parquet_reader(location.as_ref(), &rel_path)
450 .await
451 .map_err(|source| SegmentCoverageError::Storage {
452 path: path.clone(),
453 source,
454 })?;
455 let reader = ParquetRecordBatchStreamBuilder::new_with_metadata(file, metadata)
456 .with_projection(mask)
457 .with_row_groups(chunk)
458 .with_batch_size(INSPECTION_BATCH_SIZE)
459 .build()
460 .map_err(|source| SegmentCoverageError::ParquetRead {
461 path: path.clone(),
462 source,
463 backtrace: Backtrace::capture(),
464 })?;
465 compute_coverage_bitmap_from_stream(reader, &path, &index).await
466 });
467 }
468
469 let mut merged = RoaringTreemap::new();
470 while let Some(result) = tasks.join_next().await {
471 let bitmap = result.map_err(|source| SegmentCoverageError::RowGroupTask {
472 path: path.clone(),
473 source,
474 backtrace: Backtrace::capture(),
475 })??;
476 if !merged.is_disjoint(&bitmap)
477 && let Some(duplicate) = (&merged & &bitmap).min()
478 {
479 return Err(duplicate_index_interval_error(
480 &path, index, None, duplicate,
481 ));
482 }
483 merged |= bitmap;
484 }
485
486 Ok(Coverage::from_treemap(merged))
487}
488
489#[cfg(test)]
490mod tests {
491 use super::*;
492 use std::{fs::File, io::SeekFrom, num::NonZeroU64, sync::Arc};
493
494 use crate::metadata::index::TimeIndexGranularity;
495 use arrow::{
496 datatypes::{Field, Schema},
497 record_batch::RecordBatch,
498 };
499 use arrow_array::builder::{
500 BinaryBuilder, Int32Builder, StringBuilder, TimestampMillisecondBuilder,
501 };
502 use parquet::arrow::ArrowWriter;
503 use parquet::{
504 basic::Compression,
505 file::{
506 properties::WriterProperties,
507 reader::{FileReader, SerializedFileReader},
508 },
509 };
510 use tempfile::TempDir;
511 use tokio::io::{AsyncSeekExt, AsyncWriteExt};
512
513 type TestResult = Result<(), Box<dyn std::error::Error>>;
514 const EPOCH_INDEX_INTERVAL_ID: u64 = 0x8000_0000_0000_0000;
515
516 fn timestamp_index(column: &str, index_granularity: TimeIndexGranularity) -> IndexSpec {
517 IndexSpec {
518 column: column.to_string(),
519 entity_columns: Vec::new(),
520 kind: IndexKind::Timestamp {
521 index_granularity,
522 timezone: None,
523 },
524 }
525 }
526
527 fn int64_index(column: &str, index_granularity: u64) -> IndexSpec {
528 IndexSpec {
529 column: column.to_string(),
530 entity_columns: Vec::new(),
531 kind: IndexKind::Int64 {
532 index_granularity: NonZeroU64::new(index_granularity)
533 .expect("nonzero test granularity"),
534 },
535 }
536 }
537
538 fn uint64_index(column: &str, index_granularity: u64) -> IndexSpec {
539 IndexSpec {
540 column: column.to_string(),
541 entity_columns: Vec::new(),
542 kind: IndexKind::UInt64 {
543 index_granularity: NonZeroU64::new(index_granularity)
544 .expect("nonzero test granularity"),
545 },
546 }
547 }
548
549 fn write_parquet_batch(
550 path: &Path,
551 schema: Schema,
552 columns: Vec<Arc<dyn Array>>,
553 ) -> TestResult {
554 let schema = Arc::new(schema);
555 let batch = RecordBatch::try_new(Arc::clone(&schema), columns)?;
556 write_parquet_batches(
557 path,
558 schema,
559 vec![batch],
560 WriterProperties::builder().build(),
561 )
562 }
563
564 fn write_parquet_batches(
565 path: &Path,
566 schema: Arc<Schema>,
567 batches: Vec<RecordBatch>,
568 props: WriterProperties,
569 ) -> TestResult {
570 if let Some(parent) = path.parent() {
571 std::fs::create_dir_all(parent)?;
572 }
573
574 let mut writer = ArrowWriter::try_new(File::create(path)?, schema, Some(props))?;
575 for batch in batches {
576 writer.write(&batch)?;
577 writer.flush()?;
578 }
579 writer.close()?;
580 Ok(())
581 }
582
583 fn write_parquet_with_timestamps(path: &Path, ts_values: &[Option<i64>]) -> TestResult {
584 let schema = Schema::new(vec![
585 Field::new("ts", DataType::Timestamp(TimeUnit::Millisecond, None), true),
586 Field::new("val", DataType::Int32, false),
587 ]);
588
589 let mut ts_builder = TimestampMillisecondBuilder::with_capacity(ts_values.len());
590 for v in ts_values {
591 match v {
592 Some(ts) => ts_builder.append_value(*ts),
593 None => ts_builder.append_null(),
594 }
595 }
596 let ts_array = Arc::new(ts_builder.finish()) as Arc<dyn Array>;
597
598 let mut val_builder = Int32Builder::with_capacity(ts_values.len());
599 for i in 0..ts_values.len() {
600 val_builder.append_value(i as i32);
601 }
602 let val_array = Arc::new(val_builder.finish()) as Arc<dyn Array>;
603
604 write_parquet_batch(path, schema, vec![ts_array, val_array])
605 }
606
607 fn timestamp_batch(schema: Arc<Schema>, values: &[Option<i64>]) -> RecordBatch {
608 let timestamps = Arc::new(TimestampMillisecondArray::from(values.to_vec()));
609 RecordBatch::try_new(schema, vec![timestamps]).expect("timestamp batch")
610 }
611
612 fn int64_batch(schema: Arc<Schema>, values: &[Option<i64>]) -> RecordBatch {
613 RecordBatch::try_new(schema, vec![Arc::new(Int64Array::from(values.to_vec()))])
614 .expect("int64 batch")
615 }
616
617 fn uint64_batch(schema: Arc<Schema>, values: &[Option<u64>]) -> RecordBatch {
618 RecordBatch::try_new(schema, vec![Arc::new(UInt64Array::from(values.to_vec()))])
619 .expect("uint64 batch")
620 }
621
622 fn expected_interval_ids(
623 index: &IndexSpec,
624 values: impl IntoIterator<Item = IndexValue>,
625 ) -> Vec<u64> {
626 let mut interval_ids = values
627 .into_iter()
628 .map(|value| {
629 index_interval_id_for_value(&index.kind, &value).expect("valid test index value")
630 })
631 .collect::<Vec<_>>();
632 interval_ids.sort_unstable();
633 interval_ids.dedup();
634 interval_ids
635 }
636
637 #[tokio::test]
638 async fn compute_coverage_supports_nulls_and_multiple_specs() -> TestResult {
639 let tmp = TempDir::new()?;
640 let rel_path = Path::new("data/seg.parquet");
641 let abs_path = tmp.path().join(rel_path);
642
643 let ts_values = vec![Some(1_000), Some(3_600_000), None];
644 write_parquet_with_timestamps(&abs_path, &ts_values)?;
645
646 let location = TableLocation::local(tmp.path());
647
648 let cov_min = compute_segment_coverage(
649 &location,
650 rel_path,
651 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
652 )
653 .await?;
654 let interval_ids_min: Vec<u64> = cov_min.present().iter().collect();
655 assert_eq!(
656 interval_ids_min,
657 vec![EPOCH_INDEX_INTERVAL_ID, EPOCH_INDEX_INTERVAL_ID + 60]
658 );
659
660 let cov_hr = compute_segment_coverage(
661 &location,
662 rel_path,
663 ×tamp_index("ts", TimeIndexGranularity::Hours(1)),
664 )
665 .await?;
666 let interval_ids_hr: Vec<u64> = cov_hr.present().iter().collect();
667 assert_eq!(
668 interval_ids_hr,
669 vec![EPOCH_INDEX_INTERVAL_ID, EPOCH_INDEX_INTERVAL_ID + 1]
670 );
671
672 Ok(())
673 }
674
675 fn assert_implicit_duplicate(error: SegmentCoverageError, expected_path: &str) {
676 match error {
677 SegmentCoverageError::DuplicateIndexInterval {
678 path,
679 example_identity,
680 example_index_interval,
681 } => {
682 assert_eq!(path, expected_path);
683 assert_eq!(example_identity, None);
684 assert_eq!(
685 example_index_interval.to_string(),
686 "[1970-01-01T00:00:00Z, 1970-01-01T00:01:00Z)"
687 );
688 }
689 other => panic!("expected duplicate interval error, got {other:?}"),
690 }
691 }
692
693 #[tokio::test]
694 async fn compute_coverage_rejects_equal_and_distinct_timestamp_duplicates() -> TestResult {
695 let tmp = TempDir::new()?;
696 for (name, values) in [
697 ("equal", [Some(1_000), Some(1_000)]),
698 ("distinct", [Some(1_000), Some(30_000)]),
699 ] {
700 let rel_path = Path::new("data").join(format!("{name}-timestamp-duplicate.parquet"));
701 write_parquet_with_timestamps(&tmp.path().join(&rel_path), &values)?;
702
703 let error = compute_segment_coverage(
704 &TableLocation::local(tmp.path()),
705 &rel_path,
706 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
707 )
708 .await
709 .expect_err("same-worker duplicate must be rejected");
710
711 assert_implicit_duplicate(error, &rel_path.display().to_string());
712 }
713 Ok(())
714 }
715
716 #[tokio::test]
717 async fn compute_coverage_respects_exact_timestamp_boundary() -> TestResult {
718 let tmp = TempDir::new()?;
719 let rel_path = Path::new("data/timestamp-boundary.parquet");
720 write_parquet_with_timestamps(&tmp.path().join(rel_path), &[Some(59_999), Some(60_000)])?;
721
722 let coverage = compute_segment_coverage(
723 &TableLocation::local(tmp.path()),
724 rel_path,
725 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
726 )
727 .await?;
728
729 assert_eq!(
730 coverage.present().iter().collect::<Vec<_>>(),
731 vec![EPOCH_INDEX_INTERVAL_ID, EPOCH_INDEX_INTERVAL_ID + 1]
732 );
733 Ok(())
734 }
735
736 #[tokio::test]
737 async fn compute_coverage_rejects_duplicate_across_parallel_workers() -> TestResult {
738 let tmp = TempDir::new()?;
739 let rel_path = Path::new("data/cross-worker-duplicate.parquet");
740 assert_eq!(resolve_rg_settings(2), (2, 1));
741 let schema = Arc::new(Schema::new(vec![Field::new(
742 "ts",
743 DataType::Timestamp(TimeUnit::Millisecond, None),
744 true,
745 )]));
746 write_parquet_batches(
747 &tmp.path().join(rel_path),
748 Arc::clone(&schema),
749 vec![
750 timestamp_batch(Arc::clone(&schema), &[Some(1_000)]),
751 timestamp_batch(Arc::clone(&schema), &[Some(30_000)]),
752 ],
753 WriterProperties::builder().build(),
754 )?;
755
756 let error = compute_segment_coverage(
757 &TableLocation::local(tmp.path()),
758 rel_path,
759 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
760 )
761 .await
762 .expect_err("cross-worker duplicate must be rejected");
763
764 assert_implicit_duplicate(error, "data/cross-worker-duplicate.parquet");
765 Ok(())
766 }
767
768 #[tokio::test]
769 async fn compute_coverage_merges_multiple_row_groups() -> TestResult {
770 let tmp = TempDir::new()?;
771 let rel_path = Path::new("data/row_groups.parquet");
772 let schema = Arc::new(Schema::new(vec![Field::new(
773 "ts",
774 DataType::Timestamp(TimeUnit::Millisecond, None),
775 true,
776 )]));
777 let batches = vec![
778 timestamp_batch(Arc::clone(&schema), &[Some(1_000), Some(61_000)]),
779 timestamp_batch(Arc::clone(&schema), &[Some(121_000), None]),
780 timestamp_batch(Arc::clone(&schema), &[Some(181_000), Some(241_000)]),
781 ];
782 write_parquet_batches(
783 &tmp.path().join(rel_path),
784 schema,
785 batches,
786 WriterProperties::builder().build(),
787 )?;
788
789 let coverage = compute_segment_coverage(
790 &TableLocation::local(tmp.path()),
791 rel_path,
792 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
793 )
794 .await?;
795 assert_eq!(
796 coverage.present().iter().collect::<Vec<_>>(),
797 vec![
798 EPOCH_INDEX_INTERVAL_ID,
799 EPOCH_INDEX_INTERVAL_ID + 1,
800 EPOCH_INDEX_INTERVAL_ID + 2,
801 EPOCH_INDEX_INTERVAL_ID + 3,
802 EPOCH_INDEX_INTERVAL_ID + 4
803 ]
804 );
805 Ok(())
806 }
807
808 #[tokio::test]
809 async fn compute_coverage_supports_integer_indexes_across_row_groups() -> TestResult {
810 let tmp = TempDir::new()?;
811 let location = TableLocation::local(tmp.path());
812
813 let signed_path = Path::new("data/int64-row-groups.parquet");
814 let signed_schema = Arc::new(Schema::new(vec![Field::new(
815 "index",
816 DataType::Int64,
817 true,
818 )]));
819 let signed_values = [i64::MIN, -11, -1, 0, 10, i64::MAX];
820 write_parquet_batches(
821 &tmp.path().join(signed_path),
822 Arc::clone(&signed_schema),
823 vec![
824 int64_batch(Arc::clone(&signed_schema), &[Some(i64::MIN), Some(-11)]),
825 int64_batch(Arc::clone(&signed_schema), &[None, Some(-1), Some(0)]),
826 int64_batch(Arc::clone(&signed_schema), &[Some(10), Some(i64::MAX)]),
827 ],
828 WriterProperties::builder().build(),
829 )?;
830 let signed_index = int64_index("index", 10);
831 let signed = compute_segment_coverage(&location, signed_path, &signed_index).await?;
832 assert_eq!(
833 signed.present().iter().collect::<Vec<_>>(),
834 expected_interval_ids(
835 &signed_index,
836 signed_values.into_iter().map(IndexValue::Int64)
837 )
838 );
839
840 let unsigned_path = Path::new("data/uint64-row-groups.parquet");
841 let unsigned_schema = Arc::new(Schema::new(vec![Field::new(
842 "index",
843 DataType::UInt64,
844 true,
845 )]));
846 let unsigned_values = [0, 10, i64::MAX as u64 + 1, u64::MAX];
847 write_parquet_batches(
848 &tmp.path().join(unsigned_path),
849 Arc::clone(&unsigned_schema),
850 vec![
851 uint64_batch(Arc::clone(&unsigned_schema), &[Some(0)]),
852 uint64_batch(Arc::clone(&unsigned_schema), &[None, Some(10)]),
853 uint64_batch(
854 Arc::clone(&unsigned_schema),
855 &[Some(i64::MAX as u64 + 1), Some(u64::MAX)],
856 ),
857 ],
858 WriterProperties::builder().build(),
859 )?;
860 let unsigned_index = uint64_index("index", 10);
861 let unsigned = compute_segment_coverage(&location, unsigned_path, &unsigned_index).await?;
862 assert_eq!(
863 unsigned.present().iter().collect::<Vec<_>>(),
864 expected_interval_ids(
865 &unsigned_index,
866 unsigned_values.into_iter().map(IndexValue::UInt64)
867 )
868 );
869 Ok(())
870 }
871
872 #[tokio::test]
873 async fn compute_coverage_rejects_integer_duplicates_at_domain_boundaries() -> TestResult {
874 let tmp = TempDir::new()?;
875 let location = TableLocation::local(tmp.path());
876
877 for (name, values, expected_range) in [
878 ("negative", [-10, -1], "[-10, 0)"),
879 ("zero", [0, 9], "[0, 10)"),
880 (
881 "maximum",
882 [i64::MAX - 7, i64::MAX],
883 "[9223372036854775800, 9223372036854775807]",
884 ),
885 ] {
886 let rel_path = Path::new("data").join(format!("int64-{name}-duplicate.parquet"));
887 write_parquet_batch(
888 &tmp.path().join(&rel_path),
889 Schema::new(vec![Field::new("index", DataType::Int64, false)]),
890 vec![Arc::new(Int64Array::from(values.to_vec()))],
891 )?;
892
893 let error = compute_segment_coverage(&location, &rel_path, &int64_index("index", 10))
894 .await
895 .expect_err("signed duplicate must be rejected");
896 assert!(matches!(
897 error,
898 SegmentCoverageError::DuplicateIndexInterval {
899 example_identity: None,
900 example_index_interval,
901 ..
902 } if example_index_interval.to_string() == expected_range
903 ));
904 }
905
906 for (name, values, expected_range) in [
907 ("boundary", [10, 11], "[10, 20)"),
908 (
909 "maximum",
910 [u64::MAX - 5, u64::MAX],
911 "[18446744073709551610, 18446744073709551615]",
912 ),
913 ] {
914 let rel_path = Path::new("data").join(format!("uint64-{name}-duplicate.parquet"));
915 write_parquet_batch(
916 &tmp.path().join(&rel_path),
917 Schema::new(vec![Field::new("index", DataType::UInt64, false)]),
918 vec![Arc::new(UInt64Array::from(values.to_vec()))],
919 )?;
920
921 let error = compute_segment_coverage(&location, &rel_path, &uint64_index("index", 10))
922 .await
923 .expect_err("unsigned duplicate must be rejected");
924 assert!(matches!(
925 error,
926 SegmentCoverageError::DuplicateIndexInterval {
927 example_identity: None,
928 example_index_interval,
929 ..
930 } if example_index_interval.to_string() == expected_range
931 ));
932 }
933 Ok(())
934 }
935
936 #[tokio::test]
937 async fn compute_coverage_scans_multiple_bounded_batches() -> TestResult {
938 let tmp = TempDir::new()?;
939 let rel_path = Path::new("data/batches.parquet");
940 let row_count = INSPECTION_BATCH_SIZE * 2 + 17;
941 let values = (0..row_count)
942 .map(|value| Some(value as i64 * 1_000))
943 .collect::<Vec<_>>();
944 write_parquet_with_timestamps(&tmp.path().join(rel_path), &values)?;
945
946 let coverage = compute_segment_coverage(
947 &TableLocation::local(tmp.path()),
948 rel_path,
949 ×tamp_index("ts", TimeIndexGranularity::Seconds(1)),
950 )
951 .await?;
952 assert_eq!(coverage.cardinality(), row_count as u64);
953 assert_eq!(coverage.present().min(), Some(EPOCH_INDEX_INTERVAL_ID));
954 assert_eq!(
955 coverage.present().max(),
956 Some(EPOCH_INDEX_INTERVAL_ID + row_count as u64 - 1)
957 );
958 Ok(())
959 }
960
961 #[tokio::test]
962 async fn compute_coverage_rejects_duplicate_across_decoder_batches() -> TestResult {
963 let tmp = TempDir::new()?;
964 let rel_path = Path::new("data/decoder-batch-duplicate.parquet");
965 let mut values = (0..INSPECTION_BATCH_SIZE)
966 .map(|value| Some(value as i64 * 60_000))
967 .collect::<Vec<_>>();
968 values[INSPECTION_BATCH_SIZE / 2] = None;
969 values.push(Some(30_000));
970 write_parquet_with_timestamps(&tmp.path().join(rel_path), &values)?;
971
972 let error = compute_segment_coverage(
973 &TableLocation::local(tmp.path()),
974 rel_path,
975 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
976 )
977 .await
978 .expect_err("duplicate split across decoder batches must be rejected");
979
980 assert_implicit_duplicate(error, "data/decoder-batch-duplicate.parquet");
981 Ok(())
982 }
983
984 #[tokio::test]
985 async fn compute_coverage_supports_every_parquet_timestamp_unit() -> TestResult {
986 let tmp = TempDir::new()?;
987 let cases: Vec<(&str, DataType, Arc<dyn Array>)> = vec![
988 (
989 "milliseconds.parquet",
990 DataType::Timestamp(TimeUnit::Millisecond, None),
991 Arc::new(TimestampMillisecondArray::from(vec![
992 Some(1_000),
993 Some(60_000),
994 ])),
995 ),
996 (
997 "microseconds.parquet",
998 DataType::Timestamp(TimeUnit::Microsecond, None),
999 Arc::new(TimestampMicrosecondArray::from(vec![
1000 Some(1_000_000),
1001 Some(60_000_000),
1002 ])),
1003 ),
1004 (
1005 "nanoseconds.parquet",
1006 DataType::Timestamp(TimeUnit::Nanosecond, None),
1007 Arc::new(TimestampNanosecondArray::from(vec![
1008 Some(1_000_000_000),
1009 Some(60_000_000_000),
1010 ])),
1011 ),
1012 ];
1013
1014 for (file_name, data_type, array) in cases {
1015 let rel_path = Path::new("data").join(file_name);
1016 write_parquet_batch(
1017 &tmp.path().join(&rel_path),
1018 Schema::new(vec![Field::new("ts", data_type, true)]),
1019 vec![array],
1020 )?;
1021 let coverage = compute_segment_coverage(
1022 &TableLocation::local(tmp.path()),
1023 &rel_path,
1024 ×tamp_index("ts", TimeIndexGranularity::Seconds(1)),
1025 )
1026 .await?;
1027 assert_eq!(
1028 coverage.present().iter().collect::<Vec<_>>(),
1029 vec![EPOCH_INDEX_INTERVAL_ID + 1, EPOCH_INDEX_INTERVAL_ID + 60]
1030 );
1031 }
1032 Ok(())
1033 }
1034
1035 #[tokio::test]
1036 async fn compute_coverage_returns_empty_for_empty_and_all_null_files() -> TestResult {
1037 let tmp = TempDir::new()?;
1038 for (file_name, values) in [
1039 ("empty.parquet", Vec::new()),
1040 ("all_null.parquet", vec![None, None, None]),
1041 ] {
1042 let rel_path = Path::new("data").join(file_name);
1043 write_parquet_with_timestamps(&tmp.path().join(&rel_path), &values)?;
1044 let coverage = compute_segment_coverage(
1045 &TableLocation::local(tmp.path()),
1046 &rel_path,
1047 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
1048 )
1049 .await?;
1050 assert!(coverage.present().is_empty());
1051 }
1052 Ok(())
1053 }
1054
1055 #[tokio::test]
1056 async fn compute_coverage_ignores_large_unprojected_payload() -> TestResult {
1057 let tmp = TempDir::new()?;
1058 let rel_path = Path::new("data/payload.parquet");
1059 let abs_path = tmp.path().join(rel_path);
1060 let schema = Arc::new(Schema::new(vec![
1061 Field::new(
1062 "ts",
1063 DataType::Timestamp(TimeUnit::Millisecond, None),
1064 false,
1065 ),
1066 Field::new("payload", DataType::Binary, false),
1067 ]));
1068 let timestamps = Arc::new(TimestampMillisecondArray::from(vec![
1069 1_000, 61_000, 121_000, 181_000,
1070 ]));
1071 let payload = vec![0xA5; 1024 * 1024];
1072 let mut payloads = BinaryBuilder::with_capacity(4, 4 * payload.len());
1073 for _ in 0..4 {
1074 payloads.append_value(&payload);
1075 }
1076 let batch = RecordBatch::try_new(
1077 Arc::clone(&schema),
1078 vec![timestamps, Arc::new(payloads.finish())],
1079 )?;
1080 let props = WriterProperties::builder()
1081 .set_compression(Compression::UNCOMPRESSED)
1082 .set_dictionary_enabled(false)
1083 .build();
1084 write_parquet_batches(&abs_path, schema, vec![batch], props)?;
1085
1086 let reader = SerializedFileReader::new(File::open(&abs_path)?)?;
1087 let payload_page = reader.metadata().row_group(0).column(1).data_page_offset() as u64;
1088 drop(reader);
1089 let mut file = tokio::fs::OpenOptions::new()
1090 .read(true)
1091 .write(true)
1092 .open(&abs_path)
1093 .await?;
1094 file.seek(SeekFrom::Start(payload_page)).await?;
1095 file.write_all(&[0xFF; 32]).await?;
1096 file.flush().await?;
1097 drop(file);
1098
1099 assert!(tokio::fs::metadata(&abs_path).await?.len() > 4 * 1024 * 1024);
1100 let coverage = compute_segment_coverage(
1101 &TableLocation::local(tmp.path()),
1102 rel_path,
1103 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
1104 )
1105 .await?;
1106 assert_eq!(
1107 coverage.present().iter().collect::<Vec<_>>(),
1108 vec![
1109 EPOCH_INDEX_INTERVAL_ID,
1110 EPOCH_INDEX_INTERVAL_ID + 1,
1111 EPOCH_INDEX_INTERVAL_ID + 2,
1112 EPOCH_INDEX_INTERVAL_ID + 3
1113 ]
1114 );
1115 Ok(())
1116 }
1117
1118 #[tokio::test]
1119 async fn compute_coverage_errors_on_missing_time_column() -> TestResult {
1120 let tmp = TempDir::new()?;
1121 let rel_path = Path::new("data/seg.parquet");
1122 let abs_path = tmp.path().join(rel_path);
1123 write_parquet_with_timestamps(&abs_path, &[Some(1_000)])?;
1124
1125 let location = TableLocation::local(tmp.path());
1126 let err = compute_segment_coverage(
1127 &location,
1128 rel_path,
1129 ×tamp_index("missing_ts", TimeIndexGranularity::Minutes(1)),
1130 )
1131 .await
1132 .expect_err("expected missing column error");
1133
1134 assert!(matches!(
1135 err,
1136 SegmentCoverageError::OrderedIndexColumn {
1137 source: ParquetIndexColumnError {
1138 ref column,
1139 expected_domain: "timestamp",
1140 ref observed_type,
1141 ..
1142 }
1143 } if column == "missing_ts" && observed_type == "missing"
1144 ));
1145 Ok(())
1146 }
1147
1148 #[tokio::test]
1149 async fn compute_coverage_rejects_unsupported_time_type() -> TestResult {
1150 let tmp = TempDir::new()?;
1151 let rel_path = Path::new("data/string_ts.parquet");
1152 let abs_path = tmp.path().join(rel_path);
1153
1154 let schema = Schema::new(vec![
1155 Field::new("ts", DataType::Utf8, false),
1156 Field::new("val", DataType::Int32, false),
1157 ]);
1158 let mut ts_builder = StringBuilder::with_capacity(2, 8);
1159 ts_builder.append_value("a");
1160 ts_builder.append_value("b");
1161 let ts_array = Arc::new(ts_builder.finish()) as Arc<dyn Array>;
1162
1163 let mut val_builder = Int32Builder::with_capacity(2);
1164 val_builder.append_value(1);
1165 val_builder.append_value(2);
1166 let val_array = Arc::new(val_builder.finish()) as Arc<dyn Array>;
1167
1168 write_parquet_batch(&abs_path, schema, vec![ts_array, val_array])?;
1169
1170 let location = TableLocation::local(tmp.path());
1171 let err = compute_segment_coverage(
1172 &location,
1173 rel_path,
1174 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
1175 )
1176 .await
1177 .expect_err("expected unsupported arrow type");
1178
1179 assert!(matches!(
1180 err,
1181 SegmentCoverageError::OrderedIndexColumn {
1182 source: ParquetIndexColumnError {
1183 expected_domain: "timestamp",
1184 ref observed_type,
1185 ..
1186 }
1187 } if observed_type.contains("BYTE_ARRAY")
1188 ));
1189 Ok(())
1190 }
1191
1192 #[tokio::test]
1193 async fn compute_coverage_rejects_signed_unsigned_mismatch() -> TestResult {
1194 let tmp = TempDir::new()?;
1195 let rel_path = Path::new("data/signed.parquet");
1196 write_parquet_batch(
1197 &tmp.path().join(rel_path),
1198 Schema::new(vec![Field::new("index", DataType::Int64, false)]),
1199 vec![Arc::new(Int64Array::from(vec![1]))],
1200 )?;
1201
1202 let error = compute_segment_coverage(
1203 &TableLocation::local(tmp.path()),
1204 rel_path,
1205 &uint64_index("index", 1),
1206 )
1207 .await
1208 .expect_err("signed column must not be read as uint64");
1209
1210 assert!(matches!(
1211 error,
1212 SegmentCoverageError::OrderedIndexColumn {
1213 source: ParquetIndexColumnError {
1214 expected_domain: "uint64",
1215 observed_type,
1216 ..
1217 }
1218 } if observed_type.contains("logical=None")
1219 ));
1220 Ok(())
1221 }
1222
1223 #[tokio::test]
1224 async fn compute_coverage_supports_interval_ids_above_u32() -> TestResult {
1225 let tmp = TempDir::new()?;
1226 let rel_path = Path::new("data/overflow.parquet");
1227 let abs_path = tmp.path().join(rel_path);
1228 let overflow_ms = ((u32::MAX as i64) + 1) * 1_000;
1229 write_parquet_with_timestamps(&abs_path, &[Some(overflow_ms)])?;
1230
1231 let location = TableLocation::local(tmp.path());
1232 let coverage = compute_segment_coverage(
1233 &location,
1234 rel_path,
1235 ×tamp_index("ts", TimeIndexGranularity::Seconds(1)),
1236 )
1237 .await?;
1238
1239 assert!(
1240 coverage
1241 .present()
1242 .contains(0x8000_0000_0000_0000 + u64::from(u32::MAX) + 1)
1243 );
1244 Ok(())
1245 }
1246
1247 #[tokio::test]
1248 async fn compute_coverage_bubbles_up_storage_errors() -> TestResult {
1249 let tmp = TempDir::new()?;
1250 let rel_path = Path::new("missing/seg.parquet");
1251 let location = TableLocation::local(tmp.path());
1252
1253 let err = compute_segment_coverage(
1254 &location,
1255 rel_path,
1256 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
1257 )
1258 .await
1259 .expect_err("expected storage error");
1260
1261 assert!(matches!(
1262 err,
1263 SegmentCoverageError::Storage {
1264 source: StorageError::NotFound { .. },
1265 ..
1266 }
1267 ));
1268 Ok(())
1269 }
1270
1271 #[tokio::test]
1272 async fn compute_coverage_surfaces_parquet_read_errors() -> TestResult {
1273 let tmp = TempDir::new()?;
1274 let rel_path = Path::new("data/corrupt.parquet");
1275 let abs_path = tmp.path().join(rel_path);
1276 if let Some(parent) = abs_path.parent() {
1277 std::fs::create_dir_all(parent)?;
1278 }
1279 std::fs::write(&abs_path, b"not a parquet file")?;
1280
1281 let location = TableLocation::local(tmp.path());
1282 let err = compute_segment_coverage(
1283 &location,
1284 rel_path,
1285 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
1286 )
1287 .await
1288 .expect_err("expected parquet read error");
1289
1290 assert!(matches!(err, SegmentCoverageError::ParquetRead { .. }));
1291 Ok(())
1292 }
1293
1294 #[tokio::test]
1295 async fn compute_coverage_surfaces_projected_column_corruption() -> TestResult {
1296 let tmp = TempDir::new()?;
1297 let rel_path = Path::new("data/corrupt_timestamp.parquet");
1298 let abs_path = tmp.path().join(rel_path);
1299 let schema = Arc::new(Schema::new(vec![Field::new(
1300 "ts",
1301 DataType::Timestamp(TimeUnit::Millisecond, None),
1302 false,
1303 )]));
1304 let batch = RecordBatch::try_new(
1305 Arc::clone(&schema),
1306 vec![Arc::new(TimestampMillisecondArray::from(vec![
1307 1_000, 2_000,
1308 ]))],
1309 )?;
1310 let props = WriterProperties::builder()
1311 .set_compression(Compression::UNCOMPRESSED)
1312 .set_dictionary_enabled(false)
1313 .build();
1314 write_parquet_batches(&abs_path, schema, vec![batch], props)?;
1315
1316 let reader = SerializedFileReader::new(File::open(&abs_path)?)?;
1317 let timestamp_page = reader.metadata().row_group(0).column(0).data_page_offset() as u64;
1318 drop(reader);
1319 let mut file = tokio::fs::OpenOptions::new()
1320 .read(true)
1321 .write(true)
1322 .open(&abs_path)
1323 .await?;
1324 file.seek(SeekFrom::Start(timestamp_page)).await?;
1325 file.write_all(&[0xFF; 16]).await?;
1326 file.flush().await?;
1327 drop(file);
1328
1329 let err = compute_segment_coverage(
1330 &TableLocation::local(tmp.path()),
1331 rel_path,
1332 ×tamp_index("ts", TimeIndexGranularity::Minutes(1)),
1333 )
1334 .await
1335 .unwrap_err();
1336 assert!(matches!(err, SegmentCoverageError::ParquetRead { .. }));
1337 Ok(())
1338 }
1339}