1use std::{
5 borrow::Cow,
6 collections::{BTreeMap, BTreeSet},
7 fmt::Debug,
8 io::Cursor,
9 ops::Range,
10 pin::Pin,
11 sync::Arc,
12};
13
14use arrow_array::RecordBatchReader;
15use arrow_schema::Schema as ArrowSchema;
16use async_trait::async_trait;
17use byteorder::{ByteOrder, LittleEndian, ReadBytesExt};
18use bytes::{Bytes, BytesMut};
19use futures::{Stream, StreamExt, stream::BoxStream};
20use lance_core::deepsize::{Context, DeepSizeOf};
21use lance_encoding::{
22 EncodingsIo,
23 decoder::{
24 ColumnInfo, DecoderConfig, DecoderPlugins, FilterExpression, PageEncoding, ReadBatchTask,
25 RequestedRows, SchedulerDecoderConfig, schedule_and_decode, schedule_and_decode_blocking,
26 },
27 encoder::EncodedBatch,
28};
29use log::debug;
30use object_store::path::Path;
31use prost::Message;
32
33use lance_core::{
34 Error, Result,
35 cache::{CacheKey, CacheKeySchema, KeyBuilder, LanceCache},
36 datatypes::{Field, Schema},
37};
38use lance_encoding::format::pb as pbenc;
39use lance_encoding::format::pb21 as pbenc21;
40use lance_io::{
41 ReadBatchParams,
42 scheduler::FileScheduler,
43 stream::{RecordBatchStream, RecordBatchStreamAdapter},
44};
45
46use crate::{
47 datatypes::{Fields, FieldsWithMeta},
48 format::{MAGIC, pb, pbfile},
49 io::LanceEncodingsIo,
50 version::ConcreteFileVersion,
51 versions,
52};
53
54pub(crate) mod structural;
55
56pub const DEFAULT_READ_CHUNK_SIZE: u64 = 8 * 1024 * 1024;
59
60#[derive(Debug, DeepSizeOf)]
65pub struct BufferDescriptor {
66 pub position: u64,
67 pub size: u64,
68}
69
70impl BufferDescriptor {
71 fn checked_range(&self, buffer_index: usize, file_len: u64) -> Result<Range<u64>> {
72 let end = self.position.checked_add(self.size).ok_or_else(|| {
73 Error::invalid_input_source(
74 format!(
75 "Global buffer {} range overflows: position={}, size={}",
76 buffer_index, self.position, self.size
77 )
78 .into(),
79 )
80 })?;
81 if self.position > file_len {
82 return Err(Error::invalid_input_source(
83 format!(
84 "Global buffer {} position {} is outside file of size {}",
85 buffer_index, self.position, file_len
86 )
87 .into(),
88 ));
89 }
90 if end > file_len {
91 return Err(Error::invalid_input_source(
92 format!(
93 "Global buffer {} range {}..{} is outside file of size {}",
94 buffer_index, self.position, end, file_len
95 )
96 .into(),
97 ));
98 }
99 Ok(self.position..end)
100 }
101}
102
103#[derive(Debug)]
105pub struct FileStatistics {
106 pub columns: Vec<ColumnStatistics>,
108}
109
110#[derive(Debug)]
112pub struct ColumnStatistics {
113 pub num_pages: usize,
115 pub size_bytes: u64,
119}
120
121#[derive(Debug)]
123pub struct CachedFileMetadata {
124 pub file_schema: Arc<Schema>,
126 pub column_metadatas: Vec<pbfile::ColumnMetadata>,
128 pub column_infos: Vec<Arc<ColumnInfo>>,
129 pub num_rows: u64,
131 pub file_buffers: Vec<BufferDescriptor>,
132 pub num_data_bytes: u64,
134 pub num_column_metadata_bytes: u64,
137 pub num_global_buffer_bytes: u64,
139 pub num_footer_bytes: u64,
141 pub major_version: u16,
143 pub minor_version: u16,
145 pub version: ConcreteFileVersion,
146 pub file_size_bytes: u64,
148 pub retained_global_buffers: BTreeMap<u32, Bytes>,
161}
162
163impl CachedFileMetadata {
164 pub fn file_size(&self) -> u64 {
166 self.file_size_bytes
167 }
168}
169
170fn column_metadata_deep_size(column_metadatas: &[pbfile::ColumnMetadata]) -> usize {
171 column_metadatas
172 .iter()
173 .map(|cm| cm.encoded_len() * 4)
174 .sum::<usize>()
175 + std::mem::size_of_val(column_metadatas)
176}
177
178impl DeepSizeOf for CachedFileMetadata {
179 fn deep_size_of_children(&self, context: &mut Context) -> usize {
180 let schema_size = self.file_schema.deep_size_of_children(context);
181
182 let buffers_size: usize = self
183 .file_buffers
184 .iter()
185 .map(|fb| fb.deep_size_of_children(context))
186 .sum();
187
188 let column_metadatas_size = column_metadata_deep_size(self.column_metadatas.as_slice());
194
195 let column_infos_size = self.column_infos.deep_size_of_children(context);
199
200 let retained_buffers_size = self.retained_global_buffers.deep_size_of_children(context);
202
203 schema_size
204 + buffers_size
205 + column_metadatas_size
206 + column_infos_size
207 + retained_buffers_size
208 }
209}
210
211#[derive(Debug, DeepSizeOf)]
217pub struct FileMetadataIndex {
218 pub(crate) file_schema: Arc<Schema>,
219 pub(crate) num_rows: u64,
220 pub(crate) file_buffers: Vec<BufferDescriptor>,
221 pub(crate) column_metadata_offsets: Arc<[(u64, u64)]>,
222 pub(crate) num_columns: u32,
223 pub(crate) version: ConcreteFileVersion,
224 pub(crate) file_size_bytes: u64,
225 pub(crate) retained_global_buffers: BTreeMap<u32, Bytes>,
226}
227
228impl FileMetadataIndex {
229 pub fn file_size(&self) -> u64 {
231 self.file_size_bytes
232 }
233
234 pub fn num_columns(&self) -> u32 {
236 self.num_columns
237 }
238}
239
240#[derive(Debug)]
241struct CachedColumnMetadata {
242 column_metadata: pbfile::ColumnMetadata,
243 column_info: Arc<ColumnInfo>,
244}
245
246impl DeepSizeOf for CachedColumnMetadata {
247 fn deep_size_of_children(&self, context: &mut Context) -> usize {
248 column_metadata_deep_size(std::slice::from_ref(&self.column_metadata))
249 + self.column_info.deep_size_of_children(context)
250 }
251}
252
253#[derive(Debug, Clone)]
254struct ColumnMetadataCacheKey {
255 column_index: u32,
256}
257
258impl CacheKey for ColumnMetadataCacheKey {
259 type ValueType = CachedColumnMetadata;
260
261 fn key(&self) -> Cow<'_, str> {
262 Cow::Owned(format!("column_metadata/{}", self.column_index))
263 }
264
265 fn type_name() -> &'static str {
266 "ColumnMetadata"
267 }
268
269 fn schema() -> CacheKeySchema {
270 CacheKeySchema::new("lance.file.column-metadata-key", 1)
271 }
272
273 fn write_key(&self, builder: &mut KeyBuilder) {
274 builder.write_u32(self.column_index);
275 }
276}
277
278impl CachedFileMetadata {
279 pub fn version(&self) -> ConcreteFileVersion {
280 self.version
281 }
282}
283
284#[derive(Debug, Clone)]
306pub struct ReaderProjection {
307 pub schema: Arc<Schema>,
310 pub column_indices: Vec<u32>,
350}
351
352impl ReaderProjection {
353 pub fn prefers_indexed_metadata(&self, total_columns: usize) -> bool {
359 FileMetadataProvider::projection_matches_indexed_metadata(self)
360 && self.column_indices.len().saturating_mul(4) < total_columns
361 }
362}
363
364#[derive(Clone, Debug)]
366pub struct FileReaderOptions {
367 pub decoder_config: DecoderConfig,
368 pub read_chunk_size: u64,
372 pub batch_size_bytes: Option<u64>,
380}
381
382impl Default for FileReaderOptions {
383 fn default() -> Self {
384 Self {
385 decoder_config: DecoderConfig::default(),
386 read_chunk_size: DEFAULT_READ_CHUNK_SIZE,
387 batch_size_bytes: None,
388 }
389 }
390}
391
392#[derive(Debug, Clone)]
393pub(crate) struct PreparedProjection {
394 pub column_infos: Vec<Arc<ColumnInfo>>,
395 pub decoder_projection: ReaderProjection,
396}
397
398#[derive(Debug, Clone)]
399pub(crate) enum FileMetadataProvider {
400 Full(Arc<CachedFileMetadata>),
401 Indexed(Arc<FileMetadataIndex>),
402}
403
404#[async_trait]
409pub(crate) trait ReadProjection: Debug + Send + Sync {
410 fn validate_indexed(
411 &self,
412 projection: &ReaderProjection,
413 metadata_index: &FileMetadataIndex,
414 ) -> Result<()>;
415
416 fn read_length(&self, prepared: &PreparedProjection) -> Result<u64>;
417
418 async fn prepare(
419 &self,
420 metadata_provider: &FileMetadataProvider,
421 projection: &ReaderProjection,
422 io: &Arc<dyn EncodingsIo>,
423 cache: &Arc<LanceCache>,
424 ) -> Result<(PreparedProjection, u64)>;
425}
426
427#[derive(Debug, Clone)]
428pub(crate) struct DecodeEngine {
429 pub scheduler: Arc<dyn EncodingsIo>,
430 pub base_projection: ReaderProjection,
431 pub metadata_provider: FileMetadataProvider,
432 pub read_projection: Arc<dyn ReadProjection>,
433 pub decoder_plugins: Arc<DecoderPlugins>,
434 pub cache: Arc<LanceCache>,
435 pub options: FileReaderOptions,
436}
437
438#[derive(Debug, Clone)]
446pub struct ProjectedFileReader {
447 core: DecodeEngine,
448}
449
450#[derive(Debug, Clone)]
452pub struct FileReader {
453 pub(crate) core: DecodeEngine,
454 pub(crate) metadata: Arc<CachedFileMetadata>,
455}
456
457pub(crate) fn tasks_to_record_batch_stream(
458 schema: Arc<Schema>,
459 tasks: Pin<Box<dyn Stream<Item = ReadBatchTask> + Send>>,
460 batch_readahead: u32,
461) -> Pin<Box<dyn RecordBatchStream>> {
462 let arrow_schema = Arc::new(ArrowSchema::from(schema.as_ref()));
463 let batches = tasks
464 .map(|task| task.task)
465 .buffered(batch_readahead as usize)
466 .boxed();
467 Box::pin(RecordBatchStreamAdapter::new(arrow_schema, batches))
468}
469
470pub(crate) enum RawFileMetadataOpen {
471 Legacy {
472 major_version: u16,
473 minor_version: u16,
474 },
475 Current {
476 version: ConcreteFileVersion,
477 metadata: RawFileMetadata,
478 },
479}
480
481pub(crate) struct RawFileMetadata {
482 pub file_schema: Arc<Schema>,
483 pub column_metadatas: Vec<pbfile::ColumnMetadata>,
484 pub num_rows: u64,
485 pub file_buffers: Vec<BufferDescriptor>,
486 pub num_data_bytes: u64,
487 pub num_column_metadata_bytes: u64,
488 pub num_global_buffer_bytes: u64,
489 pub num_footer_bytes: u64,
490 pub footer: Footer,
491 pub file_size_bytes: u64,
492 pub retained_global_buffers: BTreeMap<u32, Bytes>,
493}
494
495#[derive(Debug)]
496pub(crate) struct Footer {
497 #[allow(dead_code)]
498 pub column_meta_start: u64,
499 #[allow(dead_code)]
502 pub column_meta_offsets_start: u64,
503 pub global_buff_offsets_start: u64,
504 pub num_global_buffers: u32,
505 pub num_columns: u32,
506 pub major_version: u16,
507 pub minor_version: u16,
508}
509
510const FOOTER_LEN: usize = 40;
511
512fn indexed_projection_column_count(field: &Field) -> Option<usize> {
521 if field.is_blob() || field.is_packed_struct() {
522 return None;
523 }
524
525 if field.children.is_empty() {
526 return Some(1);
527 }
528
529 field.children.iter().try_fold(0usize, |count, child| {
530 count.checked_add(indexed_projection_column_count(child)?)
531 })
532}
533
534pub(crate) fn normalized_column_num_rows(info: &ColumnInfo) -> Result<u64> {
540 info.page_infos.iter().try_fold(0_u64, |rows, page| {
541 let page_rows = match &page.encoding {
542 PageEncoding::Structural(layout) => match &layout.layout {
543 Some(pbenc21::page_layout::Layout::SparseLayout(sparse)) => sparse
544 .structural_layers
545 .first()
546 .and_then(|layer| layer.layer.as_ref())
547 .map_or(page.num_rows, |layer| match layer {
548 pbenc21::sparse_structural_layer::Layer::Validity(layer) => layer.num_slots,
549 pbenc21::sparse_structural_layer::Layer::List(layer) => layer.num_slots,
550 pbenc21::sparse_structural_layer::Layer::FixedSizeList(layer) => {
551 layer.num_slots
552 }
553 }),
554 _ => page.num_rows,
555 },
556 _ => page.num_rows,
557 };
558 rows.checked_add(page_rows)
559 .ok_or_else(|| Error::invalid_input_source("Column row count overflows u64".into()))
560 })
561}
562
563pub(crate) fn verify_uniform_lengths(field_lengths: &[(&str, u64)]) -> Result<u64> {
564 let first = field_lengths.first().map_or(0, |&(_, len)| len);
565 if field_lengths.iter().all(|&(_, len)| len == first) {
566 return Ok(first);
567 }
568 let columns = field_lengths
569 .iter()
570 .map(|(name, len)| format!("{name}={len}"))
571 .collect::<Vec<_>>()
572 .join(", ");
573 Err(Error::invalid_input(format!(
574 "cannot read columns of differing lengths together ({columns}); \
575 read each column (or equal-length group) separately"
576 )))
577}
578
579impl FileReader {
580 pub(crate) fn base_projection(&self) -> &ReaderProjection {
581 &self.core.base_projection
582 }
583
584 pub(crate) fn full_projection(&self, projection: ReaderProjection) -> PreparedProjection {
585 PreparedProjection {
586 column_infos: self.metadata.column_infos.clone(),
587 decoder_projection: projection,
588 }
589 }
590
591 pub(crate) async fn read_prepared_tasks(
592 &self,
593 params: ReadBatchParams,
594 batch_size: u32,
595 prepared: PreparedProjection,
596 read_len: u64,
597 filter: FilterExpression,
598 ) -> Result<Pin<Box<dyn Stream<Item = ReadBatchTask> + Send>>> {
599 self.core
600 .read_prepared_tasks(params, batch_size, prepared, read_len, filter)
601 .await
602 }
603
604 pub fn with_scheduler(&self, scheduler: Arc<dyn EncodingsIo>) -> Self {
605 Self {
606 core: self.core.with_scheduler(scheduler),
607 metadata: self.metadata.clone(),
608 }
609 }
610
611 pub fn with_io_stats(
619 &self,
620 stats: Arc<dyn lance_core::utils::io_stats::IoStatsRecorder>,
621 ) -> Self {
622 match self.core.scheduler.with_io_stats(stats) {
623 Some(scheduler) => self.with_scheduler(scheduler),
624 None => self.clone(),
625 }
626 }
627
628 pub fn num_rows(&self) -> u64 {
629 self.core.num_rows()
630 }
631
632 pub fn column_num_rows(&self, column_index: usize) -> Result<u64> {
641 let column = self
642 .metadata
643 .column_metadatas
644 .get(column_index)
645 .ok_or_else(|| {
646 Error::invalid_input(format!(
647 "column index {} is out of bounds (file has {} columns)",
648 column_index,
649 self.metadata.column_metadatas.len()
650 ))
651 })?;
652 Ok(column.pages.iter().map(|page| page.length).sum())
653 }
654
655 pub fn metadata(&self) -> &Arc<CachedFileMetadata> {
656 &self.metadata
657 }
658
659 fn statistics_from_column_metadata(
660 column_metadatas: &[pbfile::ColumnMetadata],
661 ) -> FileStatistics {
662 let column_stats = column_metadatas
663 .iter()
664 .map(|col_metadata| {
665 let num_pages = col_metadata.pages.len();
666 let size_bytes = col_metadata
667 .pages
668 .iter()
669 .map(|page| page.buffer_sizes.iter().sum::<u64>())
670 .sum::<u64>();
671 ColumnStatistics {
672 num_pages,
673 size_bytes,
674 }
675 })
676 .collect();
677
678 FileStatistics {
679 columns: column_stats,
680 }
681 }
682
683 pub fn file_statistics(&self) -> FileStatistics {
684 Self::statistics_from_column_metadata(&self.metadata().column_metadatas)
685 }
686
687 pub async fn read_global_buffer(&self, index: u32) -> Result<Bytes> {
688 self.core.read_global_buffer(index).await
689 }
690
691 async fn read_tail(scheduler: &FileScheduler) -> Result<(Bytes, u64)> {
692 let file_size = scheduler.reader().size().await? as u64;
693 let begin = if file_size < scheduler.reader().block_size() as u64 {
694 0
695 } else {
696 file_size - scheduler.reader().block_size() as u64
697 };
698 let tail_bytes = scheduler.submit_single(begin..file_size, 0).await?;
699 Ok((tail_bytes, file_size))
700 }
701
702 async fn read_range_from_tail_or_scheduler(
703 tail_bytes: &Bytes,
704 tail_offset: u64,
705 scheduler: &FileScheduler,
706 range: Range<u64>,
707 ) -> Result<Bytes> {
708 let tail_end = tail_offset + tail_bytes.len() as u64;
709 if range.start >= tail_offset && range.end <= tail_end {
710 let rel_start = (range.start - tail_offset) as usize;
711 let rel_end = (range.end - tail_offset) as usize;
712 Ok(tail_bytes.slice(rel_start..rel_end))
713 } else {
714 scheduler.submit_single(range, 0).await
715 }
716 }
717
718 fn retained_global_buffers_from_tail(
719 gbo_table: &[BufferDescriptor],
720 tail_bytes: &Bytes,
721 tail_offset: u64,
722 file_len: u64,
723 ) -> Result<BTreeMap<u32, Bytes>> {
724 let tail_end = tail_offset
725 .checked_add(tail_bytes.len() as u64)
726 .ok_or_else(|| Error::invalid_input_source("Tail byte range overflows".into()))?;
727 let mut retained_buffers = BTreeMap::new();
728 for (index, buffer) in gbo_table.iter().enumerate().skip(1) {
729 let range = buffer.checked_range(index, file_len)?;
730 if range.start >= tail_offset && range.end <= tail_end {
731 let rel_start = (range.start - tail_offset) as usize;
732 let rel_end = (range.end - tail_offset) as usize;
733 let bytes = Bytes::copy_from_slice(&tail_bytes[rel_start..rel_end]);
734 retained_buffers.insert(index as u32, bytes);
735 }
736 }
737 Ok(retained_buffers)
738 }
739
740 fn decode_footer(footer_bytes: &Bytes) -> Result<Footer> {
743 let len = footer_bytes.len();
744 if len < FOOTER_LEN {
745 return Err(Error::invalid_input(format!(
746 "does not have sufficient data, len: {}, bytes: {:?}",
747 len, footer_bytes
748 )));
749 }
750 let mut cursor = Cursor::new(footer_bytes.slice(len - FOOTER_LEN..));
751
752 let column_meta_start = cursor.read_u64::<LittleEndian>()?;
753 let column_meta_offsets_start = cursor.read_u64::<LittleEndian>()?;
754 let global_buff_offsets_start = cursor.read_u64::<LittleEndian>()?;
755 let num_global_buffers = cursor.read_u32::<LittleEndian>()?;
756 let num_columns = cursor.read_u32::<LittleEndian>()?;
757 let major_version = cursor.read_u16::<LittleEndian>()?;
758 let minor_version = cursor.read_u16::<LittleEndian>()?;
759
760 let magic_bytes = footer_bytes.slice(len - 4..);
761 if magic_bytes.as_ref() != MAGIC {
762 return Err(Error::invalid_input(format!(
763 "file does not appear to be a Lance file (invalid magic: {:?})",
764 MAGIC
765 )));
766 }
767 Ok(Footer {
768 column_meta_start,
769 column_meta_offsets_start,
770 global_buff_offsets_start,
771 num_global_buffers,
772 num_columns,
773 major_version,
774 minor_version,
775 })
776 }
777
778 fn current_file_version(footer: &Footer) -> Result<ConcreteFileVersion> {
779 let version =
780 ConcreteFileVersion::from_footer_numbers(footer.major_version, footer.minor_version)?;
781 match version {
782 ConcreteFileVersion::V1 => Err(Error::version_conflict(
783 "Attempt to use the lance v2 reader to read a legacy file".to_string(),
784 footer.major_version,
785 footer.minor_version,
786 )),
787 ConcreteFileVersion::V2_0
788 | ConcreteFileVersion::V2_1
789 | ConcreteFileVersion::V2_2
790 | ConcreteFileVersion::V2_3 => Ok(version),
791 }
792 }
793
794 fn read_all_column_metadata(
796 column_metadata_bytes: Bytes,
797 footer: &Footer,
798 ) -> Result<Vec<pbfile::ColumnMetadata>> {
799 let column_metadata_start = footer.column_meta_start;
800 let cmo_table_size = 16 * footer.num_columns as usize;
802 if column_metadata_bytes.len() < cmo_table_size {
803 return Err(Error::invalid_input(format!(
804 "column metadata region has {} bytes but CMO table needs {} bytes for {} columns",
805 column_metadata_bytes.len(),
806 cmo_table_size,
807 footer.num_columns
808 )));
809 }
810 let cmo_table = column_metadata_bytes.slice(column_metadata_bytes.len() - cmo_table_size..);
811 let column_metadata_offsets = Self::decode_cmo_table(cmo_table, footer)?;
812
813 column_metadata_offsets
814 .iter()
815 .map(|(position, length)| {
816 let normalized_position = (*position - column_metadata_start) as usize;
817 let normalized_end = normalized_position + (*length as usize);
818 Ok(pbfile::ColumnMetadata::decode(
819 &column_metadata_bytes[normalized_position..normalized_end],
820 )?)
821 })
822 .collect::<Result<Vec<_>>>()
823 }
824
825 fn decode_cmo_table(cmo_table: Bytes, footer: &Footer) -> Result<Arc<[(u64, u64)]>> {
826 let expected_size = 16 * footer.num_columns as usize;
827 if cmo_table.len() != expected_size {
828 return Err(Error::invalid_input(format!(
829 "column metadata offset table has {} bytes but expected {} bytes for {} columns",
830 cmo_table.len(),
831 expected_size,
832 footer.num_columns
833 )));
834 }
835
836 let mut offsets = Vec::with_capacity(footer.num_columns as usize);
837 for col_idx in 0..footer.num_columns {
838 let offset = (col_idx * 16) as usize;
839 let position = LittleEndian::read_u64(&cmo_table[offset..offset + 8]);
840 let length = LittleEndian::read_u64(&cmo_table[offset + 8..offset + 16]);
841 let end = position.checked_add(length).ok_or_else(|| {
842 Error::invalid_input(format!(
843 "column metadata range overflows for column index {}, position={}, length={}",
844 col_idx, position, length
845 ))
846 })?;
847 if position < footer.column_meta_start || end > footer.column_meta_offsets_start {
848 return Err(Error::invalid_input(format!(
849 "column metadata range for column index {} is outside metadata region: position={}, length={}, metadata_start={}, cmo_start={}",
850 col_idx,
851 position,
852 length,
853 footer.column_meta_start,
854 footer.column_meta_offsets_start
855 )));
856 }
857 offsets.push((position, length));
858 }
859
860 Ok(Arc::from(offsets))
861 }
862
863 async fn optimistic_tail_read(
864 data: &Bytes,
865 start_pos: u64,
866 scheduler: &FileScheduler,
867 file_len: u64,
868 ) -> Result<Bytes> {
869 let num_bytes_needed = file_len.checked_sub(start_pos).ok_or_else(|| {
870 Error::invalid_input_source(
871 format!(
872 "Tail read position {} is outside file of size {}",
873 start_pos, file_len
874 )
875 .into(),
876 )
877 })? as usize;
878 if data.len() >= num_bytes_needed {
879 Ok(data.slice((data.len() - num_bytes_needed)..))
880 } else {
881 let num_bytes_missing = (num_bytes_needed - data.len()) as u64;
882 let start = file_len - num_bytes_needed as u64;
883 let missing_bytes = scheduler
884 .submit_single(start..start + num_bytes_missing, 0)
885 .await?;
886 let mut combined = BytesMut::with_capacity(data.len() + num_bytes_missing as usize);
887 combined.extend(missing_bytes);
888 combined.extend(data);
889 Ok(combined.freeze())
890 }
891 }
892
893 fn do_decode_gbo_table(gbo_bytes: &Bytes, footer: &Footer) -> Result<Vec<BufferDescriptor>> {
894 let mut global_bufs_cursor = Cursor::new(gbo_bytes);
895
896 let mut global_buffers = Vec::with_capacity(footer.num_global_buffers as usize);
897 for _ in 0..footer.num_global_buffers {
898 let buf_pos = global_bufs_cursor.read_u64::<LittleEndian>()?;
899 let buf_size = global_bufs_cursor.read_u64::<LittleEndian>()?;
900 global_buffers.push(BufferDescriptor {
901 position: buf_pos,
902 size: buf_size,
903 });
904 }
905
906 Ok(global_buffers)
907 }
908
909 fn validate_gbo_table(
910 gbo_table: &[BufferDescriptor],
911 file_len: u64,
912 version: ConcreteFileVersion,
913 ) -> Result<()> {
914 versions::validate_global_buffers(version, gbo_table)?;
915 for (buffer_index, buffer) in gbo_table.iter().enumerate() {
916 buffer.checked_range(buffer_index, file_len)?;
917 }
918 Ok(())
919 }
920
921 async fn decode_gbo_table(
922 tail_bytes: &Bytes,
923 file_len: u64,
924 scheduler: &FileScheduler,
925 footer: &Footer,
926 version: ConcreteFileVersion,
927 ) -> Result<Vec<BufferDescriptor>> {
928 let gbo_bytes = Self::optimistic_tail_read(
931 tail_bytes,
932 footer.global_buff_offsets_start,
933 scheduler,
934 file_len,
935 )
936 .await?;
937 let gbo_table = Self::do_decode_gbo_table(&gbo_bytes, footer)?;
938 Self::validate_gbo_table(&gbo_table, file_len, version)?;
939 Ok(gbo_table)
940 }
941
942 fn decode_schema(schema_bytes: Bytes) -> Result<(u64, lance_core::datatypes::Schema)> {
943 let file_descriptor = pb::FileDescriptor::decode(schema_bytes)?;
944 let pb_schema = file_descriptor.schema.unwrap();
945 let num_rows = file_descriptor.length;
946 let fields_with_meta = FieldsWithMeta {
947 fields: Fields(pb_schema.fields),
948 metadata: pb_schema.metadata,
949 };
950 let schema = Schema::try_from(fields_with_meta)?;
951 Ok((num_rows, schema))
952 }
953
954 pub(crate) async fn read_raw_metadata_for_dispatch(
955 scheduler: &FileScheduler,
956 ) -> Result<RawFileMetadataOpen> {
957 let (tail_bytes, file_len) = Self::read_tail(scheduler).await?;
958 let tail_offset = file_len - tail_bytes.len() as u64;
959 let footer = Self::decode_footer(&tail_bytes)?;
960 let version =
961 ConcreteFileVersion::from_footer_numbers(footer.major_version, footer.minor_version)?;
962 if version == ConcreteFileVersion::V1 {
963 return Ok(RawFileMetadataOpen::Legacy {
964 major_version: footer.major_version,
965 minor_version: footer.minor_version,
966 });
967 }
968
969 let gbo_table =
970 Self::decode_gbo_table(&tail_bytes, file_len, scheduler, &footer, version).await?;
971 if gbo_table.is_empty() {
972 return Err(Error::internal(
973 "File did not contain any global buffers, schema expected".to_string(),
974 ));
975 }
976 let schema_start = gbo_table[0].position;
977 let schema_size = gbo_table[0].size;
978 let num_footer_bytes = file_len.checked_sub(schema_start).ok_or_else(|| {
979 Error::invalid_input_source(
980 format!(
981 "Schema position {} is outside file of size {}",
982 schema_start, file_len
983 )
984 .into(),
985 )
986 })?;
987 let all_metadata_bytes =
988 Self::optimistic_tail_read(&tail_bytes, schema_start, scheduler, file_len).await?;
989 let schema_bytes = all_metadata_bytes.slice(0..schema_size as usize);
990 let (num_rows, schema) = Self::decode_schema(schema_bytes)?;
991
992 let column_metadata_start = (footer.column_meta_start - schema_start) as usize;
993 let column_metadata_end = (footer.global_buff_offsets_start - schema_start) as usize;
994 let column_metadata_bytes =
995 all_metadata_bytes.slice(column_metadata_start..column_metadata_end);
996 let column_metadatas = Self::read_all_column_metadata(column_metadata_bytes, &footer)?;
997
998 let num_global_buffer_bytes = gbo_table.iter().map(|buf| buf.size).sum::<u64>();
999 let num_data_bytes = footer.column_meta_start - num_global_buffer_bytes;
1000 let num_column_metadata_bytes = footer.global_buff_offsets_start - footer.column_meta_start;
1001 let retained_global_buffers = Self::retained_global_buffers_from_tail(
1007 &gbo_table,
1008 &tail_bytes,
1009 tail_offset,
1010 file_len,
1011 )?;
1012
1013 Ok(RawFileMetadataOpen::Current {
1014 version,
1015 metadata: RawFileMetadata {
1016 file_schema: Arc::new(schema),
1017 column_metadatas,
1018 num_rows,
1019 file_buffers: gbo_table,
1020 num_data_bytes,
1021 num_column_metadata_bytes,
1022 num_global_buffer_bytes,
1023 num_footer_bytes,
1024 footer,
1025 file_size_bytes: file_len,
1026 retained_global_buffers,
1027 },
1028 })
1029 }
1030
1031 async fn read_raw_metadata_index_with_known_schema(
1032 scheduler: &FileScheduler,
1033 known_schema: Option<(Arc<Schema>, u64)>,
1034 ) -> Result<FileMetadataIndex> {
1035 let (tail_bytes, file_len) = Self::read_tail(scheduler).await?;
1036 let tail_offset = file_len - tail_bytes.len() as u64;
1037 let footer = Self::decode_footer(&tail_bytes)?;
1038
1039 let file_version = Self::current_file_version(&footer)?;
1040
1041 let gbo_table =
1042 Self::decode_gbo_table(&tail_bytes, file_len, scheduler, &footer, file_version).await?;
1043 if gbo_table.is_empty() {
1044 return Err(Error::internal(
1045 "File did not contain any global buffers, schema expected".to_string(),
1046 ));
1047 }
1048 let (file_schema, num_rows) = match known_schema {
1049 Some((file_schema, num_rows)) => (file_schema, num_rows),
1050 None => {
1051 let schema_buffer = &gbo_table[0];
1052 let schema_range = schema_buffer.checked_range(0, file_len)?;
1053 let schema_bytes = Self::read_range_from_tail_or_scheduler(
1054 &tail_bytes,
1055 tail_offset,
1056 scheduler,
1057 schema_range,
1058 )
1059 .await?;
1060 let (num_rows, schema) = Self::decode_schema(schema_bytes)?;
1061 (Arc::new(schema), num_rows)
1062 }
1063 };
1064
1065 let cmo_table = Self::read_range_from_tail_or_scheduler(
1066 &tail_bytes,
1067 tail_offset,
1068 scheduler,
1069 footer.column_meta_offsets_start..footer.global_buff_offsets_start,
1070 )
1071 .await?;
1072 let column_metadata_offsets = Self::decode_cmo_table(cmo_table, &footer)?;
1073
1074 let retained_global_buffers = Self::retained_global_buffers_from_tail(
1075 &gbo_table,
1076 &tail_bytes,
1077 tail_offset,
1078 file_len,
1079 )?;
1080
1081 Ok(FileMetadataIndex {
1082 file_schema,
1083 num_rows,
1084 file_buffers: gbo_table,
1085 column_metadata_offsets,
1086 num_columns: footer.num_columns,
1087 version: file_version,
1088 file_size_bytes: file_len,
1089 retained_global_buffers,
1090 })
1091 }
1092
1093 pub(crate) async fn read_raw_metadata_index(
1099 scheduler: &FileScheduler,
1100 ) -> Result<FileMetadataIndex> {
1101 Self::read_raw_metadata_index_with_known_schema(scheduler, None).await
1102 }
1103
1104 pub(crate) async fn read_raw_metadata_index_with_schema(
1109 scheduler: &FileScheduler,
1110 file_schema: Arc<Schema>,
1111 num_rows: u64,
1112 ) -> Result<FileMetadataIndex> {
1113 Self::read_raw_metadata_index_with_known_schema(scheduler, Some((file_schema, num_rows)))
1114 .await
1115 }
1116
1117 pub(crate) fn validate_projection(
1118 projection: &ReaderProjection,
1119 metadata: &CachedFileMetadata,
1120 ) -> Result<()> {
1121 if projection.schema.fields.is_empty() {
1122 return Err(Error::invalid_input(
1123 "Attempt to read zero columns from the file, at least one column must be specified"
1124 .to_string(),
1125 ));
1126 }
1127 let mut column_indices_seen = BTreeSet::new();
1128 for column_index in &projection.column_indices {
1129 if !column_indices_seen.insert(*column_index) {
1130 return Err(Error::invalid_input(format!(
1131 "The projection specified the column index {} more than once",
1132 column_index
1133 )));
1134 }
1135 if *column_index >= metadata.column_infos.len() as u32 {
1136 return Err(Error::invalid_input(format!(
1137 "The projection specified the column index {} but there are only {} columns in the file",
1138 column_index,
1139 metadata.column_infos.len()
1140 )));
1141 }
1142 }
1143 Ok(())
1144 }
1145
1146 fn collect_columns_from_projection(
1160 &self,
1161 _projection: &ReaderProjection,
1162 ) -> Result<Vec<Arc<ColumnInfo>>> {
1163 Ok(self.metadata.column_infos.clone())
1164 }
1165
1166 #[allow(clippy::too_many_arguments)]
1167 async fn do_read_range(
1168 column_infos: Vec<Arc<ColumnInfo>>,
1169 io: Arc<dyn EncodingsIo>,
1170 cache: Arc<LanceCache>,
1171 num_rows: u64,
1172 decoder_plugins: Arc<DecoderPlugins>,
1173 range: Range<u64>,
1174 batch_size: u32,
1175 projection: ReaderProjection,
1176 filter: FilterExpression,
1177 decoder_config: DecoderConfig,
1178 batch_size_bytes: Option<u64>,
1179 ) -> Result<BoxStream<'static, ReadBatchTask>> {
1180 debug!(
1181 "Reading range {:?} with batch_size {} from file with {} rows and {} columns into schema with {} columns",
1182 range,
1183 batch_size,
1184 num_rows,
1185 column_infos.len(),
1186 projection.schema.fields.len(),
1187 );
1188
1189 let config = SchedulerDecoderConfig {
1190 batch_size,
1191 cache,
1192 decoder_plugins,
1193 io,
1194 decoder_config,
1195 batch_size_bytes,
1196 };
1197
1198 let requested_rows = RequestedRows::Ranges(vec![range]);
1199
1200 schedule_and_decode(
1201 column_infos,
1202 requested_rows,
1203 filter,
1204 projection.column_indices,
1205 projection.schema,
1206 config,
1207 )
1208 .await
1209 }
1210
1211 #[allow(clippy::too_many_arguments)]
1212 async fn do_take_rows(
1213 column_infos: Vec<Arc<ColumnInfo>>,
1214 io: Arc<dyn EncodingsIo>,
1215 cache: Arc<LanceCache>,
1216 decoder_plugins: Arc<DecoderPlugins>,
1217 indices: Vec<u64>,
1218 batch_size: u32,
1219 projection: ReaderProjection,
1220 filter: FilterExpression,
1221 decoder_config: DecoderConfig,
1222 batch_size_bytes: Option<u64>,
1223 ) -> Result<BoxStream<'static, ReadBatchTask>> {
1224 debug!(
1225 "Taking {} rows spread across range {}..{} with batch_size {} from columns {:?}",
1226 indices.len(),
1227 indices[0],
1228 indices[indices.len() - 1],
1229 batch_size,
1230 column_infos.iter().map(|ci| ci.index).collect::<Vec<_>>()
1231 );
1232
1233 let config = SchedulerDecoderConfig {
1234 batch_size,
1235 cache,
1236 decoder_plugins,
1237 io,
1238 decoder_config,
1239 batch_size_bytes,
1240 };
1241
1242 let requested_rows = RequestedRows::Indices(indices);
1243
1244 schedule_and_decode(
1245 column_infos,
1246 requested_rows,
1247 filter,
1248 projection.column_indices,
1249 projection.schema,
1250 config,
1251 )
1252 .await
1253 }
1254
1255 #[allow(clippy::too_many_arguments)]
1256 async fn do_read_ranges(
1257 column_infos: Vec<Arc<ColumnInfo>>,
1258 io: Arc<dyn EncodingsIo>,
1259 cache: Arc<LanceCache>,
1260 decoder_plugins: Arc<DecoderPlugins>,
1261 ranges: Vec<Range<u64>>,
1262 batch_size: u32,
1263 projection: ReaderProjection,
1264 filter: FilterExpression,
1265 decoder_config: DecoderConfig,
1266 batch_size_bytes: Option<u64>,
1267 ) -> Result<BoxStream<'static, ReadBatchTask>> {
1268 let num_rows = ranges.iter().map(|r| r.end - r.start).sum::<u64>();
1269 debug!(
1270 "Taking {} ranges ({} rows) spread across range {}..{} with batch_size {} from columns {:?}",
1271 ranges.len(),
1272 num_rows,
1273 ranges[0].start,
1274 ranges[ranges.len() - 1].end,
1275 batch_size,
1276 column_infos.iter().map(|ci| ci.index).collect::<Vec<_>>()
1277 );
1278
1279 let config = SchedulerDecoderConfig {
1280 batch_size,
1281 cache,
1282 decoder_plugins,
1283 io,
1284 decoder_config,
1285 batch_size_bytes,
1286 };
1287
1288 let requested_rows = RequestedRows::Ranges(ranges);
1289
1290 schedule_and_decode(
1291 column_infos,
1292 requested_rows,
1293 filter,
1294 projection.column_indices,
1295 projection.schema,
1296 config,
1297 )
1298 .await
1299 }
1300
1301 fn take_rows_blocking(
1302 &self,
1303 indices: Vec<u64>,
1304 batch_size: u32,
1305 projection: ReaderProjection,
1306 filter: FilterExpression,
1307 ) -> Result<Box<dyn RecordBatchReader + Send + 'static>> {
1308 let column_infos = self.collect_columns_from_projection(&projection)?;
1309 debug!(
1310 "Taking {} rows spread across range {}..{} with batch_size {} from columns {:?}",
1311 indices.len(),
1312 indices[0],
1313 indices[indices.len() - 1],
1314 batch_size,
1315 column_infos.iter().map(|ci| ci.index).collect::<Vec<_>>()
1316 );
1317
1318 let config = SchedulerDecoderConfig {
1319 batch_size,
1320 cache: self.core.cache.clone(),
1321 decoder_plugins: self.core.decoder_plugins.clone(),
1322 io: self.core.scheduler.clone(),
1323 decoder_config: self.core.options.decoder_config.clone(),
1324 batch_size_bytes: self.core.options.batch_size_bytes,
1325 };
1326
1327 let requested_rows = RequestedRows::Indices(indices);
1328
1329 schedule_and_decode_blocking(
1330 column_infos,
1331 requested_rows,
1332 filter,
1333 projection.column_indices,
1334 projection.schema,
1335 config,
1336 )
1337 }
1338
1339 fn read_ranges_blocking(
1340 &self,
1341 ranges: Vec<Range<u64>>,
1342 batch_size: u32,
1343 projection: ReaderProjection,
1344 filter: FilterExpression,
1345 ) -> Result<Box<dyn RecordBatchReader + Send + 'static>> {
1346 let column_infos = self.collect_columns_from_projection(&projection)?;
1347 let num_rows = ranges.iter().map(|r| r.end - r.start).sum::<u64>();
1348 debug!(
1349 "Taking {} ranges ({} rows) spread across range {}..{} with batch_size {} from columns {:?}",
1350 ranges.len(),
1351 num_rows,
1352 ranges[0].start,
1353 ranges[ranges.len() - 1].end,
1354 batch_size,
1355 column_infos.iter().map(|ci| ci.index).collect::<Vec<_>>()
1356 );
1357
1358 let config = SchedulerDecoderConfig {
1359 batch_size,
1360 cache: self.core.cache.clone(),
1361 decoder_plugins: self.core.decoder_plugins.clone(),
1362 io: self.core.scheduler.clone(),
1363 decoder_config: self.core.options.decoder_config.clone(),
1364 batch_size_bytes: self.core.options.batch_size_bytes,
1365 };
1366
1367 let requested_rows = RequestedRows::Ranges(ranges);
1368
1369 schedule_and_decode_blocking(
1370 column_infos,
1371 requested_rows,
1372 filter,
1373 projection.column_indices,
1374 projection.schema,
1375 config,
1376 )
1377 }
1378
1379 fn read_range_blocking(
1380 &self,
1381 range: Range<u64>,
1382 batch_size: u32,
1383 projection: ReaderProjection,
1384 filter: FilterExpression,
1385 ) -> Result<Box<dyn RecordBatchReader + Send + 'static>> {
1386 let column_infos = self.collect_columns_from_projection(&projection)?;
1387 let num_rows = self.core.num_rows();
1388
1389 debug!(
1390 "Reading range {:?} with batch_size {} from file with {} rows and {} columns into schema with {} columns",
1391 range,
1392 batch_size,
1393 num_rows,
1394 column_infos.len(),
1395 projection.schema.fields.len(),
1396 );
1397
1398 let config = SchedulerDecoderConfig {
1399 batch_size,
1400 cache: self.core.cache.clone(),
1401 decoder_plugins: self.core.decoder_plugins.clone(),
1402 io: self.core.scheduler.clone(),
1403 decoder_config: self.core.options.decoder_config.clone(),
1404 batch_size_bytes: self.core.options.batch_size_bytes,
1405 };
1406
1407 let requested_rows = RequestedRows::Ranges(vec![range]);
1408
1409 schedule_and_decode_blocking(
1410 column_infos,
1411 requested_rows,
1412 filter,
1413 projection.column_indices,
1414 projection.schema,
1415 config,
1416 )
1417 }
1418
1419 pub(crate) fn read_prepared_blocking(
1420 &self,
1421 params: ReadBatchParams,
1422 batch_size: u32,
1423 prepared: PreparedProjection,
1424 read_len: u64,
1425 filter: FilterExpression,
1426 ) -> Result<Box<dyn RecordBatchReader + Send + 'static>> {
1427 let projection = prepared.decoder_projection;
1428 let verify_bound = |params: &ReadBatchParams, bound: u64, inclusive: bool| {
1429 if bound > read_len || (bound == read_len && inclusive) {
1430 Err(Error::invalid_input(format!(
1431 "cannot read {params:?} from columns with {read_len} rows"
1432 )))
1433 } else {
1434 Ok(())
1435 }
1436 };
1437 match ¶ms {
1438 ReadBatchParams::Indices(indices) => {
1439 for index in indices {
1440 match index {
1441 None => return Err(Error::invalid_input("Null value in indices array")),
1442 Some(index) => verify_bound(¶ms, index as u64, true)?,
1443 }
1444 }
1445 let indices = indices.iter().map(|index| index.unwrap() as u64).collect();
1446 self.take_rows_blocking(indices, batch_size, projection, filter)
1447 }
1448 ReadBatchParams::Range(range) => {
1449 verify_bound(¶ms, range.end as u64, false)?;
1450 self.read_range_blocking(
1451 range.start as u64..range.end as u64,
1452 batch_size,
1453 projection,
1454 filter,
1455 )
1456 }
1457 ReadBatchParams::Ranges(ranges) => {
1458 let mut ranges_u64 = Vec::with_capacity(ranges.len());
1459 for range in ranges.as_ref() {
1460 verify_bound(¶ms, range.end, false)?;
1461 ranges_u64.push(range.start..range.end);
1462 }
1463 self.read_ranges_blocking(ranges_u64, batch_size, projection, filter)
1464 }
1465 ReadBatchParams::RangeFrom(range) => {
1466 verify_bound(¶ms, range.start as u64, true)?;
1467 self.read_range_blocking(
1468 range.start as u64..read_len,
1469 batch_size,
1470 projection,
1471 filter,
1472 )
1473 }
1474 ReadBatchParams::RangeTo(range) => {
1475 verify_bound(¶ms, range.end as u64, false)?;
1476 self.read_range_blocking(0..range.end as u64, batch_size, projection, filter)
1477 }
1478 ReadBatchParams::RangeFull => {
1479 self.read_range_blocking(0..read_len, batch_size, projection, filter)
1480 }
1481 }
1482 }
1483
1484 pub fn schema(&self) -> &Arc<Schema> {
1485 self.core.schema()
1486 }
1487}
1488
1489impl FileMetadataProvider {
1490 pub(crate) fn version(&self) -> ConcreteFileVersion {
1491 match self {
1492 Self::Full(metadata) => metadata.version,
1493 Self::Indexed(metadata_index) => metadata_index.version,
1494 }
1495 }
1496
1497 pub(crate) fn num_rows(&self) -> u64 {
1498 match self {
1499 Self::Full(metadata) => metadata.num_rows,
1500 Self::Indexed(metadata_index) => metadata_index.num_rows,
1501 }
1502 }
1503
1504 pub(crate) fn schema(&self) -> &Arc<Schema> {
1505 match self {
1506 Self::Full(metadata) => &metadata.file_schema,
1507 Self::Indexed(metadata_index) => &metadata_index.file_schema,
1508 }
1509 }
1510
1511 pub(crate) fn file_buffers(&self) -> &Vec<BufferDescriptor> {
1512 match self {
1513 Self::Full(metadata) => &metadata.file_buffers,
1514 Self::Indexed(metadata_index) => &metadata_index.file_buffers,
1515 }
1516 }
1517
1518 pub(crate) fn retained_global_buffers(&self) -> &BTreeMap<u32, Bytes> {
1519 match self {
1520 Self::Full(metadata) => &metadata.retained_global_buffers,
1521 Self::Indexed(metadata_index) => &metadata_index.retained_global_buffers,
1522 }
1523 }
1524
1525 fn file_size(&self) -> u64 {
1526 match self {
1527 Self::Full(metadata) => metadata.file_size_bytes,
1528 Self::Indexed(metadata_index) => metadata_index.file_size_bytes,
1529 }
1530 }
1531
1532 pub(crate) fn file_statistics(&self) -> Option<FileStatistics> {
1533 let metadata = match self {
1534 Self::Full(metadata) => metadata,
1535 Self::Indexed(_) => return None,
1536 };
1537 Some(FileReader::statistics_from_column_metadata(
1538 &metadata.column_metadatas,
1539 ))
1540 }
1541
1542 pub(crate) fn projection_matches_indexed_metadata(projection: &ReaderProjection) -> bool {
1543 if projection.schema.fields.is_empty() {
1544 return false;
1545 }
1546
1547 projection
1548 .schema
1549 .fields
1550 .iter()
1551 .try_fold(0usize, |count, field| {
1552 count.checked_add(indexed_projection_column_count(field)?)
1553 })
1554 == Some(projection.column_indices.len())
1555 }
1556
1557 pub(crate) fn validate_indexed_projection_structure(
1558 projection: &ReaderProjection,
1559 metadata_index: &FileMetadataIndex,
1560 ) -> Result<()> {
1561 if projection.schema.fields.is_empty() {
1562 return Err(Error::invalid_input(
1563 "Attempt to read zero columns from the file, at least one column must be specified"
1564 .to_string(),
1565 ));
1566 }
1567 let mut column_indices_seen = BTreeSet::new();
1568 for column_index in &projection.column_indices {
1569 if !column_indices_seen.insert(*column_index) {
1570 return Err(Error::invalid_input(format!(
1571 "The projection specified the column index {} more than once",
1572 column_index
1573 )));
1574 }
1575 if *column_index >= metadata_index.num_columns {
1576 return Err(Error::invalid_input(format!(
1577 "The projection specified the column index {} but there are only {} columns in the file",
1578 column_index, metadata_index.num_columns
1579 )));
1580 }
1581 }
1582 Ok(())
1583 }
1584
1585 pub(crate) fn indexed_projection_error(
1586 projection: &ReaderProjection,
1587 metadata_index: &FileMetadataIndex,
1588 ) -> Error {
1589 Error::not_supported(format!(
1590 "lazy column metadata loading requires a V2.1+ ordinary structural projection without blob or packed-struct fields whose physical-column count matches the projection; got file version {:?}, {} schema fields, and {} column indices",
1591 metadata_index.version,
1592 projection.schema.fields.len(),
1593 projection.column_indices.len()
1594 ))
1595 }
1596
1597 fn column_metadata_range(
1598 metadata_index: &FileMetadataIndex,
1599 column_index: u32,
1600 ) -> Result<Range<u64>> {
1601 let (position, length) = metadata_index
1602 .column_metadata_offsets
1603 .get(column_index as usize)
1604 .copied()
1605 .ok_or_else(|| {
1606 Error::invalid_input(format!(
1607 "The projection specified the column index {} but there are only {} columns in the file",
1608 column_index, metadata_index.num_columns
1609 ))
1610 })?;
1611 let end = position.checked_add(length).ok_or_else(|| {
1612 Error::invalid_input(format!(
1613 "column metadata range overflows for column index {}, position={}, length={}",
1614 column_index, position, length
1615 ))
1616 })?;
1617 Ok(position..end)
1618 }
1619
1620 pub(crate) async fn load_indexed_column_infos<F>(
1621 metadata_index: &FileMetadataIndex,
1622 io: &Arc<dyn EncodingsIo>,
1623 cache: &Arc<LanceCache>,
1624 column_indices: &[u32],
1625 decode_column: F,
1626 ) -> Result<Vec<Arc<ColumnInfo>>>
1627 where
1628 F: Fn(u32, &pbfile::ColumnMetadata) -> Result<Arc<ColumnInfo>>,
1629 {
1630 let mut column_infos = vec![None; column_indices.len()];
1631 let mut missing_columns = Vec::new();
1632
1633 for (result_index, column_index) in column_indices.iter().copied().enumerate() {
1634 let cache_key = ColumnMetadataCacheKey { column_index };
1635 if let Some(cached) = cache.get_with_key(&cache_key).await {
1636 column_infos[result_index] = Some(cached.column_info.clone());
1637 } else {
1638 let range = Self::column_metadata_range(metadata_index, column_index)?;
1639 missing_columns.push((result_index, column_index, range));
1640 }
1641 }
1642
1643 missing_columns.sort_by_key(|(_, _, range)| range.start);
1644 if !missing_columns.is_empty() {
1645 let ranges = missing_columns
1646 .iter()
1647 .map(|(_, _, range)| range.clone())
1648 .collect::<Vec<_>>();
1649 let metadata_bytes = io.submit_request(ranges, 0).await?;
1650 for ((result_index, column_index, _), bytes) in
1651 missing_columns.into_iter().zip(metadata_bytes)
1652 {
1653 let column_metadata = pbfile::ColumnMetadata::decode(bytes)?;
1654 let column_info = decode_column(column_index, &column_metadata)?;
1655 let cached = Arc::new(CachedColumnMetadata {
1656 column_metadata,
1657 column_info: column_info.clone(),
1658 });
1659 let cache_key = ColumnMetadataCacheKey { column_index };
1660 cache.insert_with_key(&cache_key, cached).await;
1661 column_infos[result_index] = Some(column_info);
1662 }
1663 }
1664
1665 column_infos
1666 .into_iter()
1667 .enumerate()
1668 .map(|(idx, column_info)| {
1669 column_info.ok_or_else(|| {
1670 Error::internal(format!(
1671 "lazy metadata loader did not load requested projection column at position {}",
1672 idx
1673 ))
1674 })
1675 })
1676 .collect()
1677 }
1678}
1679
1680impl DecodeEngine {
1681 pub(crate) fn try_new(
1682 scheduler: Arc<dyn EncodingsIo>,
1683 base_projection: ReaderProjection,
1684 decoder_plugins: Arc<DecoderPlugins>,
1685 metadata_provider: FileMetadataProvider,
1686 read_projection: Arc<dyn ReadProjection>,
1687 cache: Arc<LanceCache>,
1688 options: FileReaderOptions,
1689 ) -> Result<Self> {
1690 Ok(Self {
1691 scheduler,
1692 base_projection,
1693 metadata_provider,
1694 read_projection,
1695 decoder_plugins,
1696 cache,
1697 options,
1698 })
1699 }
1700
1701 pub(crate) fn with_scheduler(&self, scheduler: Arc<dyn EncodingsIo>) -> Self {
1702 Self {
1703 scheduler,
1704 base_projection: self.base_projection.clone(),
1705 metadata_provider: self.metadata_provider.clone(),
1706 read_projection: self.read_projection.clone(),
1707 decoder_plugins: self.decoder_plugins.clone(),
1708 cache: self.cache.clone(),
1709 options: self.options.clone(),
1710 }
1711 }
1712
1713 fn num_rows(&self) -> u64 {
1714 self.metadata_provider.num_rows()
1715 }
1716
1717 fn schema(&self) -> &Arc<Schema> {
1718 self.metadata_provider.schema()
1719 }
1720
1721 pub(crate) async fn read_global_buffer(&self, index: u32) -> Result<Bytes> {
1722 let file_buffers = self.metadata_provider.file_buffers();
1723 let buffer_desc = file_buffers.get(index as usize).ok_or_else(|| {
1724 Error::invalid_input(format!(
1725 "request for global buffer at index {} but there were only {} global buffers in the file",
1726 index,
1727 file_buffers.len()
1728 ))
1729 })?;
1730
1731 if let Some(bytes) = self.metadata_provider.retained_global_buffers().get(&index) {
1732 return Ok(bytes.clone());
1733 }
1734
1735 let bytes = self
1736 .scheduler
1737 .submit_request(
1738 vec![
1739 buffer_desc
1740 .checked_range(index as usize, self.metadata_provider.file_size())?,
1741 ],
1742 0,
1743 )
1744 .await?;
1745 bytes.into_iter().next().ok_or_else(|| {
1746 Error::internal(format!(
1747 "global buffer read for index {} returned no bytes",
1748 index
1749 ))
1750 })
1751 }
1752
1753 async fn read_range(
1754 &self,
1755 range: Range<u64>,
1756 batch_size: u32,
1757 prepared: PreparedProjection,
1758 filter: FilterExpression,
1759 ) -> Result<BoxStream<'static, ReadBatchTask>> {
1760 FileReader::do_read_range(
1761 prepared.column_infos,
1762 self.scheduler.clone(),
1763 self.cache.clone(),
1764 self.num_rows(),
1765 self.decoder_plugins.clone(),
1766 range,
1767 batch_size,
1768 prepared.decoder_projection,
1769 filter,
1770 self.options.decoder_config.clone(),
1771 self.options.batch_size_bytes,
1772 )
1773 .await
1774 }
1775
1776 async fn take_rows(
1777 &self,
1778 indices: Vec<u64>,
1779 batch_size: u32,
1780 prepared: PreparedProjection,
1781 ) -> Result<BoxStream<'static, ReadBatchTask>> {
1782 FileReader::do_take_rows(
1783 prepared.column_infos,
1784 self.scheduler.clone(),
1785 self.cache.clone(),
1786 self.decoder_plugins.clone(),
1787 indices,
1788 batch_size,
1789 prepared.decoder_projection,
1790 FilterExpression::no_filter(),
1791 self.options.decoder_config.clone(),
1792 self.options.batch_size_bytes,
1793 )
1794 .await
1795 }
1796
1797 async fn read_ranges(
1798 &self,
1799 ranges: Vec<Range<u64>>,
1800 batch_size: u32,
1801 prepared: PreparedProjection,
1802 filter: FilterExpression,
1803 ) -> Result<BoxStream<'static, ReadBatchTask>> {
1804 FileReader::do_read_ranges(
1805 prepared.column_infos,
1806 self.scheduler.clone(),
1807 self.cache.clone(),
1808 self.decoder_plugins.clone(),
1809 ranges,
1810 batch_size,
1811 prepared.decoder_projection,
1812 filter,
1813 self.options.decoder_config.clone(),
1814 self.options.batch_size_bytes,
1815 )
1816 .await
1817 }
1818
1819 pub(crate) async fn read_prepared_tasks(
1820 &self,
1821 params: ReadBatchParams,
1822 batch_size: u32,
1823 prepared: PreparedProjection,
1824 read_len: u64,
1825 filter: FilterExpression,
1826 ) -> Result<Pin<Box<dyn Stream<Item = ReadBatchTask> + Send>>> {
1827 let verify_bound = |params: &ReadBatchParams, bound: u64, inclusive: bool| {
1828 if bound > read_len || (bound == read_len && inclusive) {
1829 Err(Error::invalid_input(format!(
1830 "cannot read {params:?} from columns with {read_len} rows"
1831 )))
1832 } else {
1833 Ok(())
1834 }
1835 };
1836 match ¶ms {
1837 ReadBatchParams::Indices(indices) => {
1838 for idx in indices {
1839 match idx {
1840 None => {
1841 return Err(Error::invalid_input("Null value in indices array"));
1842 }
1843 Some(idx) => {
1844 verify_bound(¶ms, idx as u64, true)?;
1845 }
1846 }
1847 }
1848 let indices = indices.iter().map(|idx| idx.unwrap() as u64).collect();
1849 self.take_rows(indices, batch_size, prepared).await
1850 }
1851 ReadBatchParams::Range(range) => {
1852 verify_bound(¶ms, range.end as u64, false)?;
1853 self.read_range(
1854 range.start as u64..range.end as u64,
1855 batch_size,
1856 prepared,
1857 filter,
1858 )
1859 .await
1860 }
1861 ReadBatchParams::Ranges(ranges) => {
1862 let mut ranges_u64 = Vec::with_capacity(ranges.len());
1863 for range in ranges.as_ref() {
1864 verify_bound(¶ms, range.end, false)?;
1865 ranges_u64.push(range.start..range.end);
1866 }
1867 self.read_ranges(ranges_u64, batch_size, prepared, filter)
1868 .await
1869 }
1870 ReadBatchParams::RangeFrom(range) => {
1871 verify_bound(¶ms, range.start as u64, true)?;
1872 self.read_range(range.start as u64..read_len, batch_size, prepared, filter)
1873 .await
1874 }
1875 ReadBatchParams::RangeTo(range) => {
1876 verify_bound(¶ms, range.end as u64, false)?;
1877 self.read_range(0..range.end as u64, batch_size, prepared, filter)
1878 .await
1879 }
1880 ReadBatchParams::RangeFull => {
1881 self.read_range(0..read_len, batch_size, prepared, filter)
1882 .await
1883 }
1884 }
1885 }
1886}
1887
1888impl ProjectedFileReader {
1889 pub(crate) fn base_projection(&self) -> &ReaderProjection {
1890 &self.core.base_projection
1891 }
1892
1893 pub(crate) async fn read_prepared_tasks(
1894 &self,
1895 params: ReadBatchParams,
1896 batch_size: u32,
1897 prepared: PreparedProjection,
1898 read_len: u64,
1899 filter: FilterExpression,
1900 ) -> Result<Pin<Box<dyn Stream<Item = ReadBatchTask> + Send>>> {
1901 self.core
1902 .read_prepared_tasks(params, batch_size, prepared, read_len, filter)
1903 .await
1904 }
1905
1906 pub fn with_scheduler(&self, scheduler: Arc<dyn EncodingsIo>) -> Self {
1908 Self {
1909 core: self.core.with_scheduler(scheduler),
1910 }
1911 }
1912
1913 pub fn num_rows(&self) -> u64 {
1915 self.core.num_rows()
1916 }
1917
1918 pub fn schema(&self) -> &Arc<Schema> {
1920 self.core.schema()
1921 }
1922
1923 pub fn file_statistics(&self) -> Option<FileStatistics> {
1925 self.core.metadata_provider.file_statistics()
1926 }
1927
1928 #[cfg(test)]
1929 pub(crate) fn metadata_index(&self) -> Option<&Arc<FileMetadataIndex>> {
1930 match &self.core.metadata_provider {
1931 FileMetadataProvider::Indexed(metadata_index) => Some(metadata_index),
1932 FileMetadataProvider::Full(_) => None,
1933 }
1934 }
1935
1936 pub async fn read_global_buffer(&self, index: u32) -> Result<Bytes> {
1938 self.core.read_global_buffer(index).await
1939 }
1940}
1941
1942impl FileReader {
1943 #[cfg(test)]
1944 fn scheduler(&self) -> Arc<dyn EncodingsIo> {
1945 self.core.scheduler.clone()
1946 }
1947
1948 pub async fn try_open(
1949 scheduler: FileScheduler,
1950 base_projection: Option<ReaderProjection>,
1951 decoder_plugins: Arc<DecoderPlugins>,
1952 cache: &LanceCache,
1953 options: FileReaderOptions,
1954 ) -> Result<Self> {
1955 match Self::try_open_for_dispatch(
1956 scheduler,
1957 base_projection,
1958 decoder_plugins,
1959 cache,
1960 options,
1961 )
1962 .await?
1963 {
1964 versions::OpenedFileReader::V1 {
1965 major_version,
1966 minor_version,
1967 } => Err(Error::version_conflict(
1968 "Attempt to use the Lance current-format reader to read a v1 file".to_string(),
1969 major_version,
1970 minor_version,
1971 )),
1972 versions::OpenedFileReader::Current(reader) => Ok(reader),
1973 }
1974 }
1975
1976 pub(crate) async fn try_open_for_dispatch(
1977 scheduler: FileScheduler,
1978 base_projection: Option<ReaderProjection>,
1979 decoder_plugins: Arc<DecoderPlugins>,
1980 cache: &LanceCache,
1981 options: FileReaderOptions,
1982 ) -> Result<versions::OpenedFileReader> {
1983 let metadata = match Self::read_raw_metadata_for_dispatch(&scheduler).await? {
1984 RawFileMetadataOpen::Legacy {
1985 major_version,
1986 minor_version,
1987 } => {
1988 return Ok(versions::OpenedFileReader::V1 {
1989 major_version,
1990 minor_version,
1991 });
1992 }
1993 RawFileMetadataOpen::Current { version, metadata } => {
1994 Arc::new(versions::finish_metadata(version, metadata)?)
1995 }
1996 };
1997 let path = scheduler.reader().path().clone();
1998 let io = Arc::new(
1999 LanceEncodingsIo::new(scheduler).with_read_chunk_size(options.read_chunk_size),
2000 );
2001 Self::try_open_with_file_metadata(
2002 io,
2003 path,
2004 base_projection,
2005 decoder_plugins,
2006 metadata,
2007 cache,
2008 options,
2009 )
2010 .await
2011 .map(versions::OpenedFileReader::Current)
2012 }
2013
2014 pub async fn try_open_with_file_metadata(
2015 scheduler: Arc<dyn EncodingsIo>,
2016 path: Path,
2017 base_projection: Option<ReaderProjection>,
2018 decoder_plugins: Arc<DecoderPlugins>,
2019 metadata: Arc<CachedFileMetadata>,
2020 cache: &LanceCache,
2021 options: FileReaderOptions,
2022 ) -> Result<Self> {
2023 if metadata.version == ConcreteFileVersion::V1 {
2024 return Err(Error::version_conflict(
2025 "Attempt to use the Lance current-format reader with v1 metadata".to_string(),
2026 metadata.major_version,
2027 metadata.minor_version,
2028 ));
2029 }
2030 let read_projection = versions::read_projection(metadata.version)?;
2031 let has_explicit_projection = base_projection.is_some();
2032 let base_projection = base_projection.unwrap_or_else(|| {
2033 versions::reader_projection_from_whole_schema(&metadata.file_schema, metadata.version)
2034 });
2035 if has_explicit_projection {
2036 Self::validate_projection(&base_projection, &metadata)?;
2037 }
2038 let cache = Arc::new(cache.with_key_prefix(path.as_ref()));
2039 let core = DecodeEngine::try_new(
2040 scheduler,
2041 base_projection,
2042 decoder_plugins,
2043 FileMetadataProvider::Full(metadata.clone()),
2044 read_projection,
2045 cache,
2046 options,
2047 )?;
2048 Ok(Self { core, metadata })
2049 }
2050
2051 pub async fn read_all_metadata(scheduler: &FileScheduler) -> Result<CachedFileMetadata> {
2052 match Self::read_raw_metadata_for_dispatch(scheduler).await? {
2053 RawFileMetadataOpen::Legacy {
2054 major_version,
2055 minor_version,
2056 } => Err(Error::version_conflict(
2057 "Attempt to use the Lance current-format reader to read v1 metadata".to_string(),
2058 major_version,
2059 minor_version,
2060 )),
2061 RawFileMetadataOpen::Current { version, metadata } => {
2062 versions::finish_metadata(version, metadata)
2063 }
2064 }
2065 }
2066
2067 pub async fn read_metadata_index(scheduler: &FileScheduler) -> Result<FileMetadataIndex> {
2068 let index = Self::read_raw_metadata_index(scheduler).await?;
2069 versions::finish_metadata_index(index)
2070 }
2071
2072 pub async fn read_metadata_index_with_schema(
2073 scheduler: &FileScheduler,
2074 file_schema: Arc<Schema>,
2075 num_rows: u64,
2076 ) -> Result<FileMetadataIndex> {
2077 let index =
2078 Self::read_raw_metadata_index_with_schema(scheduler, file_schema, num_rows).await?;
2079 versions::finish_metadata_index(index)
2080 }
2081
2082 pub fn version(&self) -> ConcreteFileVersion {
2083 self.metadata.version
2084 }
2085
2086 async fn prepare(&self, projection: ReaderProjection) -> Result<(PreparedProjection, u64)> {
2087 self.core
2088 .read_projection
2089 .prepare(
2090 &self.core.metadata_provider,
2091 &projection,
2092 &self.core.scheduler,
2093 &self.core.cache,
2094 )
2095 .await
2096 }
2097
2098 pub async fn read_tasks(
2099 &self,
2100 params: ReadBatchParams,
2101 batch_size: u32,
2102 projection: Option<ReaderProjection>,
2103 filter: FilterExpression,
2104 ) -> Result<Pin<Box<dyn Stream<Item = ReadBatchTask> + Send>>> {
2105 let projection = projection.unwrap_or_else(|| self.base_projection().clone());
2106 let (prepared, read_len) = self.prepare(projection).await?;
2107 self.read_prepared_tasks(params, batch_size, prepared, read_len, filter)
2108 .await
2109 }
2110
2111 pub async fn read_stream_projected(
2112 &self,
2113 params: ReadBatchParams,
2114 batch_size: u32,
2115 batch_readahead: u32,
2116 projection: ReaderProjection,
2117 filter: FilterExpression,
2118 ) -> Result<Pin<Box<dyn RecordBatchStream>>> {
2119 let schema = projection.schema.clone();
2120 let tasks = self
2121 .read_tasks(params, batch_size, Some(projection), filter)
2122 .await?;
2123 Ok(tasks_to_record_batch_stream(schema, tasks, batch_readahead))
2124 }
2125
2126 pub fn read_stream_projected_blocking(
2127 &self,
2128 params: ReadBatchParams,
2129 batch_size: u32,
2130 projection: Option<ReaderProjection>,
2131 filter: FilterExpression,
2132 ) -> Result<Box<dyn RecordBatchReader + Send + 'static>> {
2133 let projection = projection.unwrap_or_else(|| self.base_projection().clone());
2134 Self::validate_projection(&projection, self.metadata())?;
2135 let prepared = self.full_projection(projection);
2136 let read_len = self.core.read_projection.read_length(&prepared)?;
2137 self.read_prepared_blocking(params, batch_size, prepared, read_len, filter)
2138 }
2139
2140 pub async fn read_stream(
2141 &self,
2142 params: ReadBatchParams,
2143 batch_size: u32,
2144 batch_readahead: u32,
2145 filter: FilterExpression,
2146 ) -> Result<Pin<Box<dyn RecordBatchStream>>> {
2147 self.read_stream_projected(
2148 params,
2149 batch_size,
2150 batch_readahead,
2151 self.base_projection().clone(),
2152 filter,
2153 )
2154 .await
2155 }
2156}
2157
2158impl ProjectedFileReader {
2159 pub async fn try_open(
2160 scheduler: FileScheduler,
2161 base_projection: Option<ReaderProjection>,
2162 decoder_plugins: Arc<DecoderPlugins>,
2163 cache: &LanceCache,
2164 options: FileReaderOptions,
2165 ) -> Result<Self> {
2166 let base_projection = base_projection.ok_or_else(|| {
2167 Error::invalid_input("ProjectedReader requires an explicit base projection")
2168 })?;
2169 let metadata_index = Arc::new(FileReader::read_metadata_index(&scheduler).await?);
2170 let path = scheduler.reader().path().clone();
2171 let io = Arc::new(
2172 LanceEncodingsIo::new(scheduler).with_read_chunk_size(options.read_chunk_size),
2173 );
2174 Self::try_open_with_metadata_index(
2175 io,
2176 path,
2177 Some(base_projection),
2178 decoder_plugins,
2179 metadata_index,
2180 cache,
2181 options,
2182 )
2183 .await
2184 }
2185
2186 pub async fn try_open_with_metadata_index(
2187 scheduler: Arc<dyn EncodingsIo>,
2188 path: Path,
2189 base_projection: Option<ReaderProjection>,
2190 decoder_plugins: Arc<DecoderPlugins>,
2191 metadata_index: Arc<FileMetadataIndex>,
2192 cache: &LanceCache,
2193 options: FileReaderOptions,
2194 ) -> Result<Self> {
2195 if metadata_index.version == ConcreteFileVersion::V1 {
2196 return Err(Error::version_conflict(
2197 "Attempt to use the Lance projected current-format reader with v1 metadata"
2198 .to_string(),
2199 0,
2200 2,
2201 ));
2202 }
2203 let base_projection = base_projection.ok_or_else(|| {
2204 Error::invalid_input("ProjectedReader requires an explicit base projection")
2205 })?;
2206 let read_projection = versions::read_projection(metadata_index.version)?;
2207 read_projection.validate_indexed(&base_projection, &metadata_index)?;
2208 let cache = Arc::new(cache.with_key_prefix(path.as_ref()));
2209 let core = DecodeEngine::try_new(
2210 scheduler,
2211 base_projection,
2212 decoder_plugins,
2213 FileMetadataProvider::Indexed(metadata_index),
2214 read_projection,
2215 cache,
2216 options,
2217 )?;
2218 Ok(Self { core })
2219 }
2220
2221 pub async fn try_open_with_file_metadata(
2222 scheduler: Arc<dyn EncodingsIo>,
2223 path: Path,
2224 base_projection: Option<ReaderProjection>,
2225 decoder_plugins: Arc<DecoderPlugins>,
2226 metadata: Arc<CachedFileMetadata>,
2227 cache: &LanceCache,
2228 options: FileReaderOptions,
2229 ) -> Result<Self> {
2230 if metadata.version == ConcreteFileVersion::V1 {
2231 return Err(Error::version_conflict(
2232 "Attempt to use the Lance projected current-format reader with v1 metadata"
2233 .to_string(),
2234 metadata.major_version,
2235 metadata.minor_version,
2236 ));
2237 }
2238 let read_projection = versions::read_projection(metadata.version)?;
2239 let has_explicit_projection = base_projection.is_some();
2240 let base_projection = base_projection.unwrap_or_else(|| {
2241 versions::reader_projection_from_whole_schema(&metadata.file_schema, metadata.version)
2242 });
2243 if has_explicit_projection {
2244 FileReader::validate_projection(&base_projection, &metadata)?;
2245 }
2246 let cache = Arc::new(cache.with_key_prefix(path.as_ref()));
2247 let core = DecodeEngine::try_new(
2248 scheduler,
2249 base_projection,
2250 decoder_plugins,
2251 FileMetadataProvider::Full(metadata),
2252 read_projection,
2253 cache,
2254 options,
2255 )?;
2256 Ok(Self { core })
2257 }
2258
2259 pub fn version(&self) -> ConcreteFileVersion {
2260 self.core.metadata_provider.version()
2261 }
2262
2263 pub async fn read_tasks(
2264 &self,
2265 params: ReadBatchParams,
2266 batch_size: u32,
2267 projection: Option<ReaderProjection>,
2268 filter: FilterExpression,
2269 ) -> Result<Pin<Box<dyn Stream<Item = ReadBatchTask> + Send>>> {
2270 let projection = projection.unwrap_or_else(|| self.base_projection().clone());
2271 let (prepared, read_len) = self
2272 .core
2273 .read_projection
2274 .prepare(
2275 &self.core.metadata_provider,
2276 &projection,
2277 &self.core.scheduler,
2278 &self.core.cache,
2279 )
2280 .await?;
2281 self.read_prepared_tasks(params, batch_size, prepared, read_len, filter)
2282 .await
2283 }
2284}
2285
2286pub fn describe_encoding(page: &pbfile::column_metadata::Page) -> String {
2288 if let Some(encoding) = &page.encoding {
2289 if let Some(style) = &encoding.location {
2290 match style {
2291 pbfile::encoding::Location::Indirect(indirect) => {
2292 format!(
2293 "IndirectEncoding(pos={},size={})",
2294 indirect.buffer_location, indirect.buffer_length
2295 )
2296 }
2297 pbfile::encoding::Location::Direct(direct) => {
2298 let encoding_any =
2299 prost_types::Any::decode(Bytes::from(direct.encoding.clone()))
2300 .expect("failed to deserialize encoding as protobuf");
2301 if encoding_any.type_url == "/lance.encodings.ArrayEncoding" {
2302 let encoding = encoding_any.to_msg::<pbenc::ArrayEncoding>();
2303 match encoding {
2304 Ok(encoding) => {
2305 format!("{:#?}", encoding)
2306 }
2307 Err(err) => {
2308 format!("Unsupported(decode_err={})", err)
2309 }
2310 }
2311 } else if encoding_any.type_url == "/lance.encodings21.PageLayout" {
2312 let encoding = encoding_any.to_msg::<pbenc21::PageLayout>();
2313 match encoding {
2314 Ok(encoding) => {
2315 format!("{:#?}", encoding)
2316 }
2317 Err(err) => {
2318 format!("Unsupported(decode_err={})", err)
2319 }
2320 }
2321 } else {
2322 format!("Unrecognized(type_url={})", encoding_any.type_url)
2323 }
2324 }
2325 pbfile::encoding::Location::None(_) => "NoEncodingDescription".to_string(),
2326 }
2327 } else {
2328 "MISSING STYLE".to_string()
2329 }
2330 } else {
2331 "MISSING".to_string()
2332 }
2333}
2334
2335pub trait EncodedBatchReaderExt {
2336 fn try_from_mini_lance(bytes: Bytes, schema: &Schema) -> Result<Self>
2337 where
2338 Self: Sized;
2339 fn try_from_self_described_lance(bytes: Bytes) -> Result<Self>
2340 where
2341 Self: Sized;
2342}
2343
2344impl EncodedBatchReaderExt for EncodedBatch {
2345 fn try_from_mini_lance(bytes: Bytes, schema: &Schema) -> Result<Self>
2346 where
2347 Self: Sized,
2348 {
2349 let footer = FileReader::decode_footer(&bytes)?;
2350 let file_version = FileReader::current_file_version(&footer)?;
2351 let projection = versions::reader_projection_from_whole_schema(schema, file_version);
2352
2353 let column_metadata_start = footer.column_meta_start as usize;
2356 let column_metadata_end = footer.global_buff_offsets_start as usize;
2357 let column_metadata_bytes = bytes.slice(column_metadata_start..column_metadata_end);
2358 let column_metadatas =
2359 FileReader::read_all_column_metadata(column_metadata_bytes, &footer)?;
2360
2361 let page_table = versions::decode_column_metadata(file_version, &column_metadatas)?;
2362
2363 Ok(Self {
2364 data: bytes,
2365 num_rows: page_table
2366 .first()
2367 .map(|col| col.page_infos.iter().map(|page| page.num_rows).sum::<u64>())
2368 .unwrap_or(0),
2369 page_table,
2370 top_level_columns: projection.column_indices,
2371 schema: Arc::new(schema.clone()),
2372 })
2373 }
2374
2375 fn try_from_self_described_lance(bytes: Bytes) -> Result<Self>
2376 where
2377 Self: Sized,
2378 {
2379 let footer = FileReader::decode_footer(&bytes)?;
2380 let file_version = FileReader::current_file_version(&footer)?;
2381
2382 let file_len = bytes.len() as u64;
2383 let gbo_table = FileReader::do_decode_gbo_table(
2384 &bytes.slice(footer.global_buff_offsets_start as usize..),
2385 &footer,
2386 )?;
2387 FileReader::validate_gbo_table(&gbo_table, file_len, file_version)?;
2388 if gbo_table.is_empty() {
2389 return Err(Error::internal(
2390 "File did not contain any global buffers, schema expected".to_string(),
2391 ));
2392 }
2393 let schema_range = gbo_table[0].checked_range(0, file_len)?;
2394 let schema_start = schema_range.start as usize;
2395 let schema_end = schema_range.end as usize;
2396
2397 let schema_bytes = bytes.slice(schema_start..schema_end);
2398 let (_, schema) = FileReader::decode_schema(schema_bytes)?;
2399 let projection = versions::reader_projection_from_whole_schema(&schema, file_version);
2400
2401 let column_metadata_start = footer.column_meta_start as usize;
2404 let column_metadata_end = footer.global_buff_offsets_start as usize;
2405 let column_metadata_bytes = bytes.slice(column_metadata_start..column_metadata_end);
2406 let column_metadatas =
2407 FileReader::read_all_column_metadata(column_metadata_bytes, &footer)?;
2408
2409 let page_table = versions::decode_column_metadata(file_version, &column_metadatas)?;
2410
2411 Ok(Self {
2412 data: bytes,
2413 num_rows: page_table
2414 .first()
2415 .map(|col| col.page_infos.iter().map(|page| page.num_rows).sum::<u64>())
2416 .unwrap_or(0),
2417 page_table,
2418 top_level_columns: projection.column_indices,
2419 schema: Arc::new(schema),
2420 })
2421 }
2422}
2423
2424#[cfg(test)]
2425mod tests {
2426 use std::{
2427 collections::{BTreeMap, HashMap},
2428 pin::Pin,
2429 sync::Arc,
2430 };
2431
2432 use arrow_array::{
2433 DictionaryArray, Int8Array, Int32Array, ListArray, RecordBatch, RecordBatchIterator,
2434 StringArray, UInt32Array,
2435 types::{Float64Type, Int8Type, Int32Type},
2436 };
2437 use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
2438 use arrow_schema::{DataType, Field, Fields, Schema as ArrowSchema};
2439 use bytes::Bytes;
2440 use futures::{StreamExt, prelude::stream::TryStreamExt};
2441 use lance_arrow::{BLOB_META_KEY, RecordBatchExt};
2442 use lance_core::{ArrowResult, datatypes::Schema};
2443 use lance_datagen::{ArrayGeneratorExt, BatchCount, ByteCount, RowCount, array, gen_batch};
2444 use lance_encoding::{
2445 constants::{STRUCTURAL_ENCODING_META_KEY, STRUCTURAL_ENCODING_SPARSE},
2446 decoder::{
2447 DecodeBatchScheduler, DecoderPlugins, EncodedBatchLayout, FilterExpression,
2448 PageEncoding, ReadBatchTask, decode_batch,
2449 },
2450 encoder::{EncodedBatch, EncodingOptions, encode_batch},
2451 format::pb21,
2452 };
2453 use lance_io::{stream::RecordBatchStream, utils::CachedFileSize};
2454 use log::debug;
2455 use rstest::rstest;
2456 use tokio::sync::mpsc;
2457
2458 use crate::reader::{
2459 EncodedBatchReaderExt, FileReader, FileReaderOptions, ProjectedFileReader, ReaderProjection,
2460 };
2461 use crate::testing::{FsFixture, WrittenFile, test_cache, write_lance_file};
2462 use crate::version::{ConcreteFileVersion, LanceFileVersion};
2463 use crate::versions;
2464 use crate::writer::{FileWriterOptions, PAGE_BUFFER_ALIGNMENT};
2465 use lance_encoding::decoder::DecoderConfig;
2466
2467 fn footer_version(bytes: &[u8]) -> (u16, u16) {
2468 let version_start = bytes.len() - 8;
2469 (
2470 u16::from_le_bytes([bytes[version_start], bytes[version_start + 1]]),
2471 u16::from_le_bytes([bytes[version_start + 2], bytes[version_start + 3]]),
2472 )
2473 }
2474
2475 #[tokio::test]
2476 async fn sparse_file_writer_reader_scan_range_and_take_roundtrip() {
2477 let fs = FsFixture::default();
2478 let sparse_metadata = HashMap::from([(
2479 STRUCTURAL_ENCODING_META_KEY.to_string(),
2480 STRUCTURAL_ENCODING_SPARSE.to_string(),
2481 )]);
2482 let value_field =
2483 Field::new("values", DataType::Int32, true).with_metadata(sparse_metadata.clone());
2484 let item_field = Arc::new(Field::new("item", DataType::Int32, true));
2485 let list_field = Field::new("items", DataType::List(item_field.clone()), true)
2486 .with_metadata(sparse_metadata);
2487 let arrow_schema = Arc::new(ArrowSchema::new(vec![value_field, list_field]));
2488 let list = ListArray::try_new(
2489 item_field,
2490 OffsetBuffer::new(ScalarBuffer::from(vec![0_i32, 2, 2, 2, 3, 3, 5])),
2491 Arc::new(Int32Array::from(vec![
2492 Some(1),
2493 None,
2494 Some(3),
2495 Some(4),
2496 Some(5),
2497 ])),
2498 Some(NullBuffer::from(vec![true, false, true, true, true, true])),
2499 )
2500 .unwrap();
2501 let batch = RecordBatch::try_new(
2502 arrow_schema.clone(),
2503 vec![
2504 Arc::new(Int32Array::from(vec![
2505 Some(10),
2506 None,
2507 Some(30),
2508 Some(40),
2509 None,
2510 Some(60),
2511 ])),
2512 Arc::new(list),
2513 ],
2514 )
2515 .unwrap();
2516 let input = RecordBatchIterator::new(vec![Ok(batch.clone())], arrow_schema);
2517 write_lance_file(
2518 input,
2519 &fs,
2520 ConcreteFileVersion::V2_3,
2521 FileWriterOptions::default(),
2522 )
2523 .await;
2524
2525 let file_scheduler = fs
2526 .scheduler
2527 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
2528 .await
2529 .unwrap();
2530 let file_reader = FileReader::try_open(
2531 file_scheduler,
2532 None,
2533 Arc::<DecoderPlugins>::default(),
2534 &test_cache(),
2535 FileReaderOptions::default(),
2536 )
2537 .await
2538 .unwrap();
2539 assert_eq!(file_reader.metadata().column_infos.len(), 2);
2540 assert!(
2541 file_reader
2542 .metadata()
2543 .column_infos
2544 .iter()
2545 .flat_map(|column| column.page_infos.iter())
2546 .all(|page| {
2547 matches!(
2548 &page.encoding,
2549 PageEncoding::Structural(layout)
2550 if matches!(
2551 layout.layout,
2552 Some(pb21::page_layout::Layout::SparseLayout(_))
2553 )
2554 )
2555 })
2556 );
2557
2558 let scan = file_reader
2559 .read_stream(
2560 lance_io::ReadBatchParams::RangeFull,
2561 1024,
2562 1,
2563 FilterExpression::no_filter(),
2564 )
2565 .await
2566 .unwrap()
2567 .try_collect::<Vec<_>>()
2568 .await
2569 .unwrap();
2570 assert_eq!(scan, vec![batch.clone()]);
2571
2572 let range = file_reader
2573 .read_stream(
2574 lance_io::ReadBatchParams::Range(1..5),
2575 1024,
2576 1,
2577 FilterExpression::no_filter(),
2578 )
2579 .await
2580 .unwrap()
2581 .try_collect::<Vec<_>>()
2582 .await
2583 .unwrap();
2584 assert_eq!(range, vec![batch.slice(1, 4)]);
2585
2586 let indices = UInt32Array::from(vec![0, 3, 5]);
2587 let take = file_reader
2588 .read_stream(
2589 lance_io::ReadBatchParams::Indices(indices.clone()),
2590 1024,
2591 1,
2592 FilterExpression::no_filter(),
2593 )
2594 .await
2595 .unwrap()
2596 .try_collect::<Vec<_>>()
2597 .await
2598 .unwrap();
2599 assert_eq!(take, vec![batch.take(&indices).unwrap()]);
2600 }
2601
2602 #[tokio::test]
2603 async fn full_int8_dictionary_v2_2_roundtrip() {
2604 let fs = FsFixture::default();
2605 let values = Arc::new(StringArray::from(
2606 (0..=i8::MAX)
2607 .map(|value| format!("value-{value}"))
2608 .collect::<Vec<_>>(),
2609 ));
2610 let keys = Int8Array::from((0..=i8::MAX).collect::<Vec<_>>());
2611 let dictionary = Arc::new(DictionaryArray::<Int8Type>::new(keys, values));
2612 let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
2613 "dictionary",
2614 DataType::Dictionary(Box::new(DataType::Int8), Box::new(DataType::Utf8)),
2615 true,
2616 )]));
2617 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![dictionary]).unwrap();
2618
2619 write_lance_file(
2620 RecordBatchIterator::new([Ok(batch.clone())], arrow_schema),
2621 &fs,
2622 ConcreteFileVersion::V2_2,
2623 FileWriterOptions::default(),
2624 )
2625 .await;
2626
2627 let file_scheduler = fs
2628 .scheduler
2629 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
2630 .await
2631 .unwrap();
2632 let file_reader = FileReader::try_open(
2633 file_scheduler,
2634 None,
2635 Arc::<DecoderPlugins>::default(),
2636 &test_cache(),
2637 FileReaderOptions::default(),
2638 )
2639 .await
2640 .unwrap();
2641 let actual = file_reader
2642 .read_stream(
2643 lance_io::ReadBatchParams::RangeFull,
2644 1024,
2645 1,
2646 FilterExpression::no_filter(),
2647 )
2648 .await
2649 .unwrap()
2650 .try_collect::<Vec<_>>()
2651 .await
2652 .unwrap();
2653
2654 assert_eq!(actual, vec![batch]);
2655 }
2656
2657 async fn create_some_file(fs: &FsFixture, version: ConcreteFileVersion) -> WrittenFile {
2658 let location_type = DataType::Struct(Fields::from(vec![
2659 Field::new("x", DataType::Float64, true),
2660 Field::new("y", DataType::Float64, true),
2661 ]));
2662 let categories_type = DataType::List(Arc::new(Field::new("item", DataType::Utf8, true)));
2663
2664 let mut reader = gen_batch()
2665 .col("score", array::rand::<Float64Type>())
2666 .col("location", array::rand_type(&location_type))
2667 .col("categories", array::rand_type(&categories_type))
2668 .col("binary", array::rand_type(&DataType::Binary));
2669 if version == ConcreteFileVersion::V2_0 {
2670 reader = reader.col("large_bin", array::rand_type(&DataType::LargeBinary));
2671 }
2672 let reader = reader.into_reader_rows(RowCount::from(1000), BatchCount::from(100));
2673
2674 write_lance_file(reader, fs, version, FileWriterOptions::default()).await
2675 }
2676
2677 async fn create_wide_direct_file(fs: &FsFixture, num_columns: usize) -> WrittenFile {
2678 let mut reader = gen_batch();
2679 for column_idx in 0..num_columns {
2680 reader = reader.col(format!("c{column_idx}"), array::step::<Int32Type>());
2681 }
2682 let reader = reader.into_reader_rows(RowCount::from(1000), BatchCount::from(100));
2683
2684 write_lance_file(
2685 reader,
2686 fs,
2687 ConcreteFileVersion::V2_1,
2688 FileWriterOptions::default(),
2689 )
2690 .await
2691 }
2692
2693 async fn create_wide_fixed_size_list_file(fs: &FsFixture, num_columns: usize) -> WrittenFile {
2694 let data_type =
2695 DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Float32, true)), 4);
2696 let mut reader = gen_batch();
2697 for column_idx in 0..num_columns {
2698 reader = reader.col(
2699 format!("c{column_idx}"),
2700 array::rand_type(&data_type).with_random_nulls(0.1),
2701 );
2702 }
2703 let reader = reader.into_reader_rows(RowCount::from(64), BatchCount::from(4));
2704
2705 write_lance_file(
2706 reader,
2707 fs,
2708 ConcreteFileVersion::V2_1,
2709 FileWriterOptions::default(),
2710 )
2711 .await
2712 }
2713
2714 async fn create_wide_structural_file(fs: &FsFixture, num_groups: usize) -> WrittenFile {
2715 let struct_type = DataType::Struct(Fields::from(vec![
2716 Field::new("x", DataType::Int32, true),
2717 Field::new("y", DataType::Int32, true),
2718 ]));
2719 let list_type = DataType::List(Arc::new(Field::new("item", DataType::Int32, true)));
2720 let mut reader = gen_batch();
2721 for group_idx in 0..num_groups {
2722 reader = reader
2723 .col(
2724 format!("s{group_idx}"),
2725 array::rand_type(&struct_type).with_random_nulls(0.5),
2726 )
2727 .col(
2728 format!("l{group_idx}"),
2729 array::rand_type(&list_type).with_random_nulls(0.5),
2730 );
2731 }
2732 let reader = reader.into_reader_rows(RowCount::from(64), BatchCount::from(4));
2733
2734 write_lance_file(
2735 reader,
2736 fs,
2737 ConcreteFileVersion::V2_1,
2738 FileWriterOptions::default(),
2739 )
2740 .await
2741 }
2742
2743 type Transformer = Box<dyn Fn(&RecordBatch) -> RecordBatch>;
2744
2745 async fn verify_expected(
2746 expected: &[RecordBatch],
2747 mut actual: Pin<Box<dyn RecordBatchStream>>,
2748 read_size: u32,
2749 transform: Option<Transformer>,
2750 ) {
2751 let mut remaining = expected.iter().map(|batch| batch.num_rows()).sum::<usize>() as u32;
2752 let mut expected_iter = expected.iter().map(|batch| {
2753 if let Some(transform) = &transform {
2754 transform(batch)
2755 } else {
2756 batch.clone()
2757 }
2758 });
2759 let mut next_expected = expected_iter.next().unwrap().clone();
2760 while let Some(actual) = actual.next().await {
2761 let mut actual = actual.unwrap();
2762 let mut rows_to_verify = actual.num_rows() as u32;
2763 let expected_length = remaining.min(read_size);
2764 assert_eq!(expected_length, rows_to_verify);
2765
2766 while rows_to_verify > 0 {
2767 let next_slice_len = (next_expected.num_rows() as u32).min(rows_to_verify);
2768 assert_eq!(
2769 next_expected.slice(0, next_slice_len as usize),
2770 actual.slice(0, next_slice_len as usize)
2771 );
2772 remaining -= next_slice_len;
2773 rows_to_verify -= next_slice_len;
2774 if remaining > 0 {
2775 if next_slice_len == next_expected.num_rows() as u32 {
2776 next_expected = expected_iter.next().unwrap().clone();
2777 } else {
2778 next_expected = next_expected.slice(
2779 next_slice_len as usize,
2780 next_expected.num_rows() - next_slice_len as usize,
2781 );
2782 }
2783 }
2784 if rows_to_verify > 0 {
2785 actual = actual.slice(
2786 next_slice_len as usize,
2787 actual.num_rows() - next_slice_len as usize,
2788 );
2789 }
2790 }
2791 }
2792 assert_eq!(remaining, 0);
2793 }
2794
2795 async fn collect_read_tasks(
2796 tasks: Pin<Box<dyn futures::Stream<Item = ReadBatchTask> + Send>>,
2797 readahead: usize,
2798 ) -> Vec<RecordBatch> {
2799 tasks
2800 .map(|task| task.task)
2801 .buffered(readahead)
2802 .try_collect::<Vec<_>>()
2803 .await
2804 .unwrap()
2805 }
2806
2807 async fn read_file_with_mutated_bytes(
2811 version: LanceFileVersion,
2812 batch: RecordBatch,
2813 pattern: &[u8],
2814 patch_offset: usize,
2815 patch: &[u8],
2816 ) -> lance_core::Result<Vec<RecordBatch>> {
2817 let fs = FsFixture::default();
2818 let schema = batch.schema();
2819 write_lance_file(
2820 RecordBatchIterator::new(vec![Ok(batch)], schema),
2821 &fs,
2822 version.resolve(),
2823 FileWriterOptions::default(),
2824 )
2825 .await;
2826
2827 let mut bytes = fs
2828 .object_store
2829 .read_one_all(&fs.tmp_path)
2830 .await
2831 .unwrap()
2832 .to_vec();
2833 let matches = bytes
2834 .windows(pattern.len())
2835 .enumerate()
2836 .filter_map(|(position, window)| (window == pattern).then_some(position))
2837 .collect::<Vec<_>>();
2838 assert_eq!(
2839 matches.len(),
2840 1,
2841 "expected the byte pattern to appear exactly once in the file"
2842 );
2843 let patch_start = matches[0] + patch_offset;
2844 bytes[patch_start..patch_start + patch.len()].copy_from_slice(patch);
2845 fs.object_store.put(&fs.tmp_path, &bytes).await.unwrap();
2846
2847 let file_scheduler = fs
2848 .scheduler
2849 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
2850 .await
2851 .unwrap();
2852 let file_reader = FileReader::try_open(
2853 file_scheduler,
2854 None,
2855 Arc::<DecoderPlugins>::default(),
2856 &test_cache(),
2857 FileReaderOptions::default(),
2858 )
2859 .await
2860 .unwrap();
2861 file_reader
2862 .read_stream(
2863 lance_io::ReadBatchParams::RangeFull,
2864 1024,
2865 16,
2866 FilterExpression::no_filter(),
2867 )
2868 .await?
2869 .try_collect::<Vec<_>>()
2870 .await
2871 }
2872
2873 #[tokio::test]
2874 async fn test_reader_rejects_excess_miniblock_row_counts() {
2875 let batch =
2876 arrow_array::record_batch!(("id", UInt64, (0..2048_u64).collect::<Vec<_>>())).unwrap();
2877 let fs = FsFixture::default();
2878 write_lance_file(
2879 RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema()),
2880 &fs,
2881 ConcreteFileVersion::V2_1,
2882 FileWriterOptions::default(),
2883 )
2884 .await;
2885
2886 let mut bytes = fs
2887 .object_store
2888 .read_one_all(&fs.tmp_path)
2889 .await
2890 .unwrap()
2891 .to_vec();
2892 bytes[0] ^= 0xf7;
2895 fs.object_store.put(&fs.tmp_path, &bytes).await.unwrap();
2896
2897 let file_scheduler = fs
2898 .scheduler
2899 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
2900 .await
2901 .unwrap();
2902 let file_reader = FileReader::try_open(
2903 file_scheduler,
2904 None,
2905 Arc::<DecoderPlugins>::default(),
2906 &test_cache(),
2907 FileReaderOptions::default(),
2908 )
2909 .await
2910 .unwrap();
2911 let result = file_reader
2912 .read_stream(
2913 lance_io::ReadBatchParams::RangeFull,
2914 1024,
2915 16,
2916 FilterExpression::no_filter(),
2917 )
2918 .await;
2919 let error = match result {
2920 Ok(stream) => stream
2921 .try_collect::<Vec<_>>()
2922 .await
2923 .expect_err("excess mini-block row counts must fail the read"),
2924 Err(error) => error,
2925 };
2926 assert!(
2927 matches!(error, lance_core::Error::CorruptFile { .. }),
2928 "expected CorruptFile, got: {error}"
2929 );
2930 assert!(
2931 error.to_string().contains("exceeding items_in_page"),
2932 "unexpected message: {error}"
2933 );
2934 }
2935
2936 #[rstest]
2946 #[tokio::test]
2947 async fn test_default_reader_rejects_out_of_bounds_variable_width_offsets(
2948 #[values(LanceFileVersion::V2_1, LanceFileVersion::V2_2, LanceFileVersion::V2_3)]
2949 version: LanceFileVersion,
2950 ) {
2951 use arrow_array::{Array, DictionaryArray, Int32Array, StringArray};
2952
2953 let values = StringArray::from(vec!["alpha", "beta", "gamma"]);
2954 let indices = Int32Array::from((0..300).map(|i| i % 3).collect::<Vec<i32>>());
2955 let dictionary = DictionaryArray::new(indices, Arc::new(values));
2956 let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
2957 "category",
2958 dictionary.data_type().clone(),
2959 false,
2960 )]));
2961 let batch = RecordBatch::try_new(arrow_schema, vec![Arc::new(dictionary)]).unwrap();
2962
2963 let offsets_tail_pattern = [5_i32, 9, 14]
2970 .iter()
2971 .flat_map(|value| value.to_le_bytes())
2972 .collect::<Vec<u8>>();
2973 let error = read_file_with_mutated_bytes(
2974 version,
2975 batch,
2976 &offsets_tail_pattern,
2977 8,
2978 &100_000_i32.to_le_bytes(),
2979 )
2980 .await
2981 .expect_err("out-of-bounds offsets must fail the read");
2982 assert!(
2983 matches!(error, lance_core::Error::CorruptFile { .. }),
2984 "expected CorruptFile, got: {error}"
2985 );
2986 assert!(
2987 error.to_string().contains("out of bounds"),
2988 "unexpected message: {error}"
2989 );
2990 }
2991
2992 #[rstest]
2997 #[tokio::test]
2998 async fn test_default_reader_rejects_non_monotonic_storage_dictionary_offsets(
2999 #[values(LanceFileVersion::V2_1, LanceFileVersion::V2_2, LanceFileVersion::V2_3)]
3000 version: LanceFileVersion,
3001 ) {
3002 use arrow_array::StringArray;
3003
3004 let metadata = HashMap::from([
3005 (
3006 "lance-encoding:dict-size-ratio".to_string(),
3007 "0.99".to_string(),
3008 ),
3009 (
3010 "lance-encoding:dict-values-compression".to_string(),
3011 "none".to_string(),
3012 ),
3013 ]);
3014 let arrow_schema = Arc::new(ArrowSchema::new(vec![
3015 Field::new("category", DataType::Utf8, false).with_metadata(metadata),
3016 ]));
3017 let values = (0..300)
3018 .map(|index| match index % 3 {
3019 0 => "alpha",
3020 1 => "beta",
3021 _ => "gamma",
3022 })
3023 .collect::<Vec<_>>();
3024 let batch =
3025 RecordBatch::try_new(arrow_schema, vec![Arc::new(StringArray::from(values))]).unwrap();
3026
3027 let offsets_tail_pattern = [5_i32, 9, 14]
3028 .iter()
3029 .flat_map(|value| value.to_le_bytes())
3030 .collect::<Vec<u8>>();
3031 let error = read_file_with_mutated_bytes(
3032 version,
3033 batch,
3034 &offsets_tail_pattern,
3035 4,
3036 &2_i32.to_le_bytes(),
3037 )
3038 .await
3039 .expect_err("non-monotonic dictionary offsets must fail the read");
3040 assert!(
3041 matches!(error, lance_core::Error::CorruptFile { .. }),
3042 "expected CorruptFile, got: {error}"
3043 );
3044 assert!(
3045 error.to_string().contains("decreases"),
3046 "unexpected message: {error}"
3047 );
3048 }
3049
3050 #[rstest]
3055 #[tokio::test]
3056 async fn test_default_reader_rejects_out_of_bounds_miniblock_offsets(
3057 #[values(LanceFileVersion::V2_1, LanceFileVersion::V2_2, LanceFileVersion::V2_3)]
3058 version: LanceFileVersion,
3059 ) {
3060 use arrow_array::StringArray;
3061
3062 let arrow_schema = Arc::new(ArrowSchema::new(vec![Field::new(
3063 "strings",
3064 DataType::Utf8,
3065 false,
3066 )]));
3067 let batch = RecordBatch::try_new(
3068 arrow_schema,
3069 vec![Arc::new(StringArray::from(vec!["alpha", "beta", "gamma"]))],
3070 )
3071 .unwrap();
3072
3073 let chunk_offsets_pattern = [16_i32, 21, 25, 30]
3078 .iter()
3079 .flat_map(|value| value.to_le_bytes())
3080 .collect::<Vec<u8>>();
3081 let error = read_file_with_mutated_bytes(
3082 version,
3083 batch,
3084 &chunk_offsets_pattern,
3085 12,
3086 &100_000_i32.to_le_bytes(),
3087 )
3088 .await
3089 .expect_err("an out-of-bounds chunk offset must fail the read");
3090 assert!(
3091 matches!(error, lance_core::Error::CorruptFile { .. }),
3092 "expected CorruptFile, got: {error}"
3093 );
3094 assert!(
3095 error.to_string().contains("out of bounds"),
3096 "unexpected message: {error}"
3097 );
3098 }
3099
3100 #[tokio::test]
3101 async fn test_round_trip() {
3102 let fs = FsFixture::default();
3103
3104 let WrittenFile { data, .. } = create_some_file(&fs, ConcreteFileVersion::V2_0).await;
3105
3106 let file_size = fs.object_store.size(&fs.tmp_path).await.unwrap() as usize;
3107 let footer = fs
3108 .object_store
3109 .open(&fs.tmp_path)
3110 .await
3111 .unwrap()
3112 .get_range(file_size - 8..file_size)
3113 .await
3114 .unwrap();
3115 assert_eq!(footer_version(&footer), (0, 3));
3116 assert_eq!(
3117 crate::determine_file_version(&fs.object_store, &fs.tmp_path, Some(file_size))
3118 .await
3119 .unwrap(),
3120 ConcreteFileVersion::V2_0
3121 );
3122
3123 for read_size in [32, 1024, 1024 * 1024] {
3124 let file_scheduler = fs
3125 .scheduler
3126 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3127 .await
3128 .unwrap();
3129 let file_reader = FileReader::try_open(
3130 file_scheduler,
3131 None,
3132 Arc::<DecoderPlugins>::default(),
3133 &test_cache(),
3134 FileReaderOptions::default(),
3135 )
3136 .await
3137 .unwrap();
3138
3139 assert_eq!(
3140 (
3141 file_reader.metadata().major_version,
3142 file_reader.metadata().minor_version
3143 ),
3144 (0, 3)
3145 );
3146 let schema = file_reader.schema();
3147 assert_eq!(schema.metadata.get("foo").unwrap(), "bar");
3148
3149 let batch_stream = file_reader
3150 .read_stream(
3151 lance_io::ReadBatchParams::RangeFull,
3152 read_size,
3153 16,
3154 FilterExpression::no_filter(),
3155 )
3156 .await
3157 .unwrap();
3158
3159 verify_expected(&data, batch_stream, read_size, None).await;
3160 }
3161 }
3162
3163 #[rstest]
3164 #[test_log::test(tokio::test)]
3165 async fn test_encoded_batch_round_trip(
3166 #[values(ConcreteFileVersion::V2_0)] version: ConcreteFileVersion,
3168 ) {
3169 let data = gen_batch()
3170 .col("x", array::rand::<Int32Type>())
3171 .col("y", array::rand_utf8(ByteCount::from(16), false))
3172 .into_batch_rows(RowCount::from(10000))
3173 .unwrap();
3174
3175 let lance_schema = Arc::new(Schema::try_from(data.schema().as_ref()).unwrap());
3176
3177 let encoding_options = EncodingOptions {
3178 cache_bytes_per_column: 4096,
3179 max_page_bytes: 32 * 1024 * 1024,
3180 keep_original_array: true,
3181 buffer_alignment: 64,
3182 };
3183
3184 let encoding_strategy = crate::versions::v2_0::encoding_strategy();
3185
3186 let encoded_batch = encode_batch(
3187 &data,
3188 lance_schema.clone(),
3189 encoding_strategy.as_ref(),
3190 &encoding_options,
3191 )
3192 .await
3193 .unwrap();
3194
3195 let bytes = versions::encode_self_described_batch(version, &encoded_batch).unwrap();
3197 assert_eq!(footer_version(&bytes), (2, 0));
3198
3199 let decoded_batch = EncodedBatch::try_from_self_described_lance(bytes).unwrap();
3200
3201 let decoded = decode_batch(
3202 &decoded_batch,
3203 &FilterExpression::no_filter(),
3204 Arc::<DecoderPlugins>::default(),
3205 false,
3206 EncodedBatchLayout::Array,
3207 None,
3208 )
3209 .await
3210 .unwrap();
3211
3212 assert_eq!(data, decoded);
3213
3214 let bytes = versions::encode_mini_batch(version, &encoded_batch).unwrap();
3216 assert_eq!(footer_version(&bytes), (2, 0));
3217 let decoded_batch =
3218 EncodedBatch::try_from_mini_lance(bytes, lance_schema.as_ref()).unwrap();
3219 let decoded = decode_batch(
3220 &decoded_batch,
3221 &FilterExpression::no_filter(),
3222 Arc::<DecoderPlugins>::default(),
3223 false,
3224 EncodedBatchLayout::Array,
3225 None,
3226 )
3227 .await
3228 .unwrap();
3229
3230 assert_eq!(data, decoded);
3231 }
3232
3233 #[rstest]
3234 #[test_log::test(tokio::test)]
3235 async fn test_projection(
3236 #[values(
3237 ConcreteFileVersion::V2_0,
3238 ConcreteFileVersion::V2_1,
3239 ConcreteFileVersion::V2_2
3240 )]
3241 version: ConcreteFileVersion,
3242 ) {
3243 let fs = FsFixture::default();
3244
3245 let written_file = create_some_file(&fs, version).await;
3246 let file_scheduler = fs
3247 .scheduler
3248 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3249 .await
3250 .unwrap();
3251
3252 let field_id_mapping = written_file
3253 .field_id_mapping
3254 .iter()
3255 .copied()
3256 .collect::<BTreeMap<_, _>>();
3257
3258 let empty_projection = ReaderProjection {
3259 column_indices: Vec::default(),
3260 schema: Arc::new(Schema::default()),
3261 };
3262
3263 for columns in [
3264 vec!["score"],
3265 vec!["location"],
3266 vec!["categories"],
3267 vec!["score.x"],
3268 vec!["score", "categories"],
3269 vec!["score", "location"],
3270 vec!["location", "categories"],
3271 vec!["score.y", "location", "categories"],
3272 ] {
3273 debug!("Testing round trip with projection {:?}", columns);
3274 for use_field_ids in [true, false] {
3275 let file_reader = FileReader::try_open(
3277 file_scheduler.clone(),
3278 None,
3279 Arc::<DecoderPlugins>::default(),
3280 &test_cache(),
3281 FileReaderOptions::default(),
3282 )
3283 .await
3284 .unwrap();
3285
3286 let projected_schema = written_file.schema.project(&columns).unwrap();
3287 let projection = if use_field_ids {
3288 versions::reader_projection_from_field_ids(
3289 file_reader.metadata().version(),
3290 &projected_schema,
3291 &field_id_mapping,
3292 )
3293 .unwrap()
3294 } else {
3295 versions::reader_projection_from_column_names(
3296 file_reader.metadata().version(),
3297 &written_file.schema,
3298 &columns,
3299 )
3300 .unwrap()
3301 };
3302
3303 let batch_stream = file_reader
3304 .read_stream_projected(
3305 lance_io::ReadBatchParams::RangeFull,
3306 1024,
3307 16,
3308 projection.clone(),
3309 FilterExpression::no_filter(),
3310 )
3311 .await
3312 .unwrap();
3313
3314 let projection_arrow = ArrowSchema::from(projection.schema.as_ref());
3315 verify_expected(
3316 &written_file.data,
3317 batch_stream,
3318 1024,
3319 Some(Box::new(move |batch: &RecordBatch| {
3320 batch.project_by_schema(&projection_arrow).unwrap()
3321 })),
3322 )
3323 .await;
3324
3325 let file_reader = FileReader::try_open(
3327 file_scheduler.clone(),
3328 Some(projection.clone()),
3329 Arc::<DecoderPlugins>::default(),
3330 &test_cache(),
3331 FileReaderOptions::default(),
3332 )
3333 .await
3334 .unwrap();
3335
3336 let batch_stream = file_reader
3337 .read_stream(
3338 lance_io::ReadBatchParams::RangeFull,
3339 1024,
3340 16,
3341 FilterExpression::no_filter(),
3342 )
3343 .await
3344 .unwrap();
3345
3346 let projection_arrow = ArrowSchema::from(projection.schema.as_ref());
3347 verify_expected(
3348 &written_file.data,
3349 batch_stream,
3350 1024,
3351 Some(Box::new(move |batch: &RecordBatch| {
3352 batch.project_by_schema(&projection_arrow).unwrap()
3353 })),
3354 )
3355 .await;
3356
3357 assert!(
3358 file_reader
3359 .read_stream_projected(
3360 lance_io::ReadBatchParams::RangeFull,
3361 1024,
3362 16,
3363 empty_projection.clone(),
3364 FilterExpression::no_filter(),
3365 )
3366 .await
3367 .is_err()
3368 );
3369 }
3370 }
3371
3372 assert!(
3373 FileReader::try_open(
3374 file_scheduler.clone(),
3375 Some(empty_projection),
3376 Arc::<DecoderPlugins>::default(),
3377 &test_cache(),
3378 FileReaderOptions::default(),
3379 )
3380 .await
3381 .is_err()
3382 );
3383
3384 let arrow_schema = ArrowSchema::new(vec![
3385 Field::new("x", DataType::Int32, true),
3386 Field::new("y", DataType::Int32, true),
3387 ]);
3388 let schema = Schema::try_from(&arrow_schema).unwrap();
3389
3390 let projection_with_dupes = ReaderProjection {
3391 column_indices: vec![0, 0],
3392 schema: Arc::new(schema),
3393 };
3394
3395 assert!(
3396 FileReader::try_open(
3397 file_scheduler.clone(),
3398 Some(projection_with_dupes),
3399 Arc::<DecoderPlugins>::default(),
3400 &test_cache(),
3401 FileReaderOptions::default(),
3402 )
3403 .await
3404 .is_err()
3405 );
3406 }
3407
3408 #[tokio::test]
3409 async fn test_lazy_reader_direct_projection_matches_eager_reader() {
3410 let fs = FsFixture::default();
3411 let written_file = create_wide_direct_file(&fs, 16).await;
3412
3413 let file_scheduler = fs
3414 .scheduler
3415 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3416 .await
3417 .unwrap();
3418 let projection = versions::reader_projection_from_column_names(
3419 ConcreteFileVersion::V2_1,
3420 &written_file.schema,
3421 &["c10"],
3422 )
3423 .unwrap();
3424
3425 let eager_reader = FileReader::try_open(
3426 file_scheduler.clone(),
3427 None,
3428 Arc::<DecoderPlugins>::default(),
3429 &test_cache(),
3430 FileReaderOptions::default(),
3431 )
3432 .await
3433 .unwrap();
3434 let expected = eager_reader
3435 .read_stream_projected(
3436 lance_io::ReadBatchParams::RangeFull,
3437 127,
3438 16,
3439 projection.clone(),
3440 FilterExpression::no_filter(),
3441 )
3442 .await
3443 .unwrap()
3444 .try_collect::<Vec<_>>()
3445 .await
3446 .unwrap();
3447
3448 let cache = test_cache();
3449 let lazy_reader = ProjectedFileReader::try_open(
3450 file_scheduler,
3451 Some(projection.clone()),
3452 Arc::<DecoderPlugins>::default(),
3453 &cache,
3454 FileReaderOptions::default(),
3455 )
3456 .await
3457 .unwrap();
3458 let tasks = lazy_reader
3459 .read_tasks(
3460 lance_io::ReadBatchParams::RangeFull,
3461 127,
3462 None,
3463 FilterExpression::no_filter(),
3464 )
3465 .await
3466 .unwrap();
3467 let actual = collect_read_tasks(tasks, 16).await;
3468
3469 assert_eq!(expected, actual);
3470 }
3471
3472 #[tokio::test]
3473 async fn test_lazy_reader_loads_only_requested_column_metadata() {
3474 let fs = FsFixture::default();
3475 let written_file = create_wide_direct_file(&fs, 512).await;
3476
3477 let projection = versions::reader_projection_from_column_names(
3478 ConcreteFileVersion::V2_1,
3479 &written_file.schema,
3480 &["c0"],
3481 )
3482 .unwrap();
3483 let file_scheduler = fs
3484 .scheduler
3485 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3486 .await
3487 .unwrap();
3488 let lazy_reader = ProjectedFileReader::try_open(
3489 file_scheduler,
3490 Some(projection.clone()),
3491 Arc::<DecoderPlugins>::default(),
3492 &test_cache(),
3493 FileReaderOptions::default(),
3494 )
3495 .await
3496 .unwrap();
3497 let selected_column = projection.column_indices[0] as usize;
3498 let requested_metadata_bytes = lazy_reader
3499 .metadata_index()
3500 .unwrap()
3501 .column_metadata_offsets[selected_column]
3502 .1;
3503 let total_metadata_bytes = lazy_reader
3504 .metadata_index()
3505 .unwrap()
3506 .column_metadata_offsets
3507 .iter()
3508 .map(|(_, length)| *length)
3509 .sum::<u64>();
3510 assert!(
3511 total_metadata_bytes > 8 * fs.object_store.block_size() as u64,
3512 "test file metadata is too small to prove lazy loading: {total_metadata_bytes} bytes"
3513 );
3514
3515 fs.object_store.io_stats_incremental();
3516 let tasks = lazy_reader
3517 .read_tasks(
3518 lance_io::ReadBatchParams::Range(0..0),
3519 1024,
3520 Some(projection.clone()),
3521 FilterExpression::no_filter(),
3522 )
3523 .await
3524 .unwrap();
3525 let batches = collect_read_tasks(tasks, 1).await;
3526 assert!(batches.is_empty());
3527
3528 let stats = fs.object_store.io_stats_incremental();
3529 assert!(
3530 stats.read_bytes < total_metadata_bytes / 2,
3531 "lazy read fetched too much metadata: read {} bytes, requested column metadata is {} bytes, total column metadata is {} bytes",
3532 stats.read_bytes,
3533 requested_metadata_bytes,
3534 total_metadata_bytes
3535 );
3536
3537 fs.object_store.io_stats_incremental();
3538 let tasks = lazy_reader
3539 .read_tasks(
3540 lance_io::ReadBatchParams::Range(0..0),
3541 1024,
3542 Some(projection),
3543 FilterExpression::no_filter(),
3544 )
3545 .await
3546 .unwrap();
3547 let batches = collect_read_tasks(tasks, 1).await;
3548 assert!(batches.is_empty());
3549
3550 let stats = fs.object_store.io_stats_incremental();
3551 assert_eq!(
3552 stats.read_iops, 0,
3553 "cached column metadata should avoid repeat metadata I/O"
3554 );
3555 assert_eq!(
3556 stats.read_bytes, 0,
3557 "cached column metadata should avoid repeat metadata reads"
3558 );
3559 }
3560
3561 async fn assert_lazy_projection_matches_eager_and_reads_metadata_subset(
3562 fs: &FsFixture,
3563 projection: ReaderProjection,
3564 shape: &str,
3565 ) -> Vec<RecordBatch> {
3566 let file_scheduler = fs
3567 .scheduler
3568 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3569 .await
3570 .unwrap();
3571 let eager_reader = FileReader::try_open(
3572 file_scheduler.clone(),
3573 None,
3574 Arc::<DecoderPlugins>::default(),
3575 &test_cache(),
3576 FileReaderOptions::default(),
3577 )
3578 .await
3579 .unwrap();
3580 let expected = eager_reader
3581 .read_stream_projected(
3582 lance_io::ReadBatchParams::RangeFull,
3583 127,
3584 16,
3585 projection.clone(),
3586 FilterExpression::no_filter(),
3587 )
3588 .await
3589 .unwrap()
3590 .try_collect::<Vec<_>>()
3591 .await
3592 .unwrap();
3593
3594 let cache = test_cache();
3595 let lazy_reader = ProjectedFileReader::try_open(
3596 file_scheduler,
3597 Some(projection.clone()),
3598 Arc::<DecoderPlugins>::default(),
3599 &cache,
3600 FileReaderOptions::default(),
3601 )
3602 .await
3603 .unwrap();
3604 let metadata_index = lazy_reader.metadata_index().unwrap();
3605 let requested_metadata_bytes = projection
3606 .column_indices
3607 .iter()
3608 .map(|column_index| metadata_index.column_metadata_offsets[*column_index as usize].1)
3609 .sum::<u64>();
3610 let total_metadata_bytes = metadata_index
3611 .column_metadata_offsets
3612 .iter()
3613 .map(|(_, length)| *length)
3614 .sum::<u64>();
3615 assert!(total_metadata_bytes > requested_metadata_bytes * 8);
3616
3617 fs.object_store.io_stats_incremental();
3618 let tasks = lazy_reader
3619 .read_tasks(
3620 lance_io::ReadBatchParams::Range(0..0),
3621 127,
3622 None,
3623 FilterExpression::no_filter(),
3624 )
3625 .await
3626 .unwrap();
3627 assert!(collect_read_tasks(tasks, 1).await.is_empty());
3628 let metadata_stats = fs.object_store.io_stats_incremental();
3629 assert!(
3630 metadata_stats.read_bytes < total_metadata_bytes / 2,
3631 "lazy {shape} read fetched too much metadata: read {} bytes, requested column metadata is {} bytes, total column metadata is {} bytes",
3632 metadata_stats.read_bytes,
3633 requested_metadata_bytes,
3634 total_metadata_bytes
3635 );
3636
3637 let tasks = lazy_reader
3638 .read_tasks(
3639 lance_io::ReadBatchParams::RangeFull,
3640 127,
3641 None,
3642 FilterExpression::no_filter(),
3643 )
3644 .await
3645 .unwrap();
3646 let actual = collect_read_tasks(tasks, 16).await;
3647 assert_eq!(expected, actual);
3648 actual
3649 }
3650
3651 #[tokio::test]
3652 async fn test_lazy_reader_fixed_size_list_projection_matches_eager_reader() {
3653 let fs = FsFixture::default();
3654 let written_file = create_wide_fixed_size_list_file(&fs, 512).await;
3655 let projection = versions::reader_projection_from_column_names(
3656 ConcreteFileVersion::V2_1,
3657 &written_file.schema,
3658 &["c17", "c509"],
3659 )
3660 .unwrap();
3661 assert!(projection.prefers_indexed_metadata(512));
3662 assert_lazy_projection_matches_eager_and_reads_metadata_subset(
3663 &fs,
3664 projection,
3665 "fixed-size-list",
3666 )
3667 .await;
3668 }
3669
3670 #[tokio::test]
3671 async fn test_v2_0_rejects_indexed_metadata_reader() {
3672 let fs = FsFixture::default();
3673 let written_file = create_some_file(&fs, ConcreteFileVersion::V2_0).await;
3674 let projection = versions::reader_projection_from_column_names(
3675 ConcreteFileVersion::V2_0,
3676 &written_file.schema,
3677 &["score"],
3678 )
3679 .unwrap();
3680 assert!(projection.prefers_indexed_metadata(100));
3681
3682 let file_scheduler = fs
3683 .scheduler
3684 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3685 .await
3686 .unwrap();
3687 let err = ProjectedFileReader::try_open(
3688 file_scheduler,
3689 Some(projection),
3690 Arc::<DecoderPlugins>::default(),
3691 &test_cache(),
3692 FileReaderOptions::default(),
3693 )
3694 .await
3695 .unwrap_err();
3696 assert!(
3697 matches!(err, lance_core::Error::NotSupported { .. }),
3698 "expected V2.0 indexed metadata open to fail, got {err:?}"
3699 );
3700 }
3701
3702 #[tokio::test]
3703 async fn test_lazy_reader_nested_projection_compacts_physical_columns() {
3704 let fs = FsFixture::default();
3705 let written_file = create_wide_structural_file(&fs, 128).await;
3706 let projection = versions::reader_projection_from_column_names(
3707 ConcreteFileVersion::V2_1,
3708 &written_file.schema,
3709 &["s97.y", "l4", "s3"],
3710 )
3711 .unwrap();
3712
3713 assert_eq!(
3714 projection
3715 .schema
3716 .fields
3717 .iter()
3718 .map(|field| field.name.as_str())
3719 .collect::<Vec<_>>(),
3720 vec!["s97", "l4", "s3"]
3721 );
3722 assert_eq!(projection.schema.fields[0].children.len(), 1);
3723 assert_eq!(projection.schema.fields[0].children[0].name, "y");
3724 assert_eq!(projection.schema.fields[2].children.len(), 2);
3725 assert_eq!(projection.column_indices.len(), 4);
3726 assert!(
3727 projection
3728 .column_indices
3729 .windows(2)
3730 .any(|indices| indices[0] > indices[1]),
3731 "the projection must reorder physical columns to exercise compact remapping"
3732 );
3733 assert!(projection.prefers_indexed_metadata(128 * 4));
3734 let actual = assert_lazy_projection_matches_eager_and_reads_metadata_subset(
3735 &fs, projection, "nested",
3736 )
3737 .await;
3738 assert!(
3739 actual
3740 .iter()
3741 .flat_map(|batch| batch.columns())
3742 .any(|column| column.null_count() > 0),
3743 "the structural projection must exercise nullable arrays"
3744 );
3745 }
3746
3747 #[rstest]
3748 #[case::before_metadata_region(90, 5)]
3749 #[case::after_metadata_region(190, 20)]
3750 fn test_decode_cmo_table_rejects_out_of_range_offsets(
3751 #[case] position: u64,
3752 #[case] length: u64,
3753 ) {
3754 let mut cmo_table = [0; 16];
3755 cmo_table[0..8].copy_from_slice(&position.to_le_bytes());
3756 cmo_table[8..16].copy_from_slice(&length.to_le_bytes());
3757 let footer = super::Footer {
3758 column_meta_start: 100,
3759 column_meta_offsets_start: 200,
3760 global_buff_offsets_start: 200,
3761 num_global_buffers: 0,
3762 num_columns: 1,
3763 major_version: 2,
3764 minor_version: 1,
3765 };
3766
3767 let err = FileReader::decode_cmo_table(Bytes::copy_from_slice(&cmo_table), &footer)
3768 .expect_err("out-of-range CMO entries must be rejected");
3769 assert!(
3770 matches!(err, lance_core::Error::InvalidInput { .. }),
3771 "expected InvalidInput, got {err:?}"
3772 );
3773 }
3774
3775 #[rstest]
3776 #[case::blob(BLOB_META_KEY)]
3777 #[case::packed_struct("lance-encoding:packed")]
3778 #[tokio::test]
3779 async fn test_lazy_reader_rejects_opaque_projection(#[case] metadata_key: &str) {
3780 let fs = FsFixture::default();
3781 let written_file = create_some_file(&fs, ConcreteFileVersion::V2_1).await;
3782
3783 let ordinary_projection = versions::reader_projection_from_column_names(
3784 ConcreteFileVersion::V2_1,
3785 &written_file.schema,
3786 &["location.x"],
3787 )
3788 .unwrap();
3789 assert_eq!(ordinary_projection.schema.fields[0].children.len(), 1);
3790 assert!(ordinary_projection.prefers_indexed_metadata(100));
3791
3792 let file_scheduler = fs
3793 .scheduler
3794 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3795 .await
3796 .unwrap();
3797 let err = ProjectedFileReader::try_open(
3798 file_scheduler.clone(),
3799 None,
3800 Arc::<DecoderPlugins>::default(),
3801 &test_cache(),
3802 FileReaderOptions::default(),
3803 )
3804 .await
3805 .unwrap_err();
3806 assert!(
3807 matches!(err, lance_core::Error::InvalidInput { .. }),
3808 "expected InvalidInput, got {err:?}"
3809 );
3810
3811 let mut projection = ordinary_projection;
3812 Arc::make_mut(&mut projection.schema).fields[0]
3813 .metadata
3814 .insert(metadata_key.to_string(), "true".to_string());
3815 assert!(!projection.prefers_indexed_metadata(100));
3816
3817 let err = ProjectedFileReader::try_open(
3818 file_scheduler,
3819 Some(projection),
3820 Arc::<DecoderPlugins>::default(),
3821 &test_cache(),
3822 FileReaderOptions::default(),
3823 )
3824 .await
3825 .unwrap_err();
3826 assert!(
3827 matches!(err, lance_core::Error::NotSupported { .. }),
3828 "expected NotSupported for {metadata_key}, got {err:?}"
3829 );
3830 }
3831
3832 #[tokio::test]
3839 async fn test_lazy_reader_validates_unequal_length_projection() {
3840 use arrow_array::Int32Array;
3841 use lance_io::ReadBatchParams;
3842
3843 let arrow_schema = Arc::new(ArrowSchema::new(vec![
3844 Field::new("a", DataType::Int32, true),
3845 Field::new("c", DataType::Int32, true),
3846 ]));
3847 let lance_schema = Schema::try_from(arrow_schema.as_ref()).unwrap();
3848
3849 let fs = FsFixture::default();
3850 let mut writer = versions::v2_1::create_writer(
3851 fs.object_store.create(&fs.tmp_path).await.unwrap(),
3852 lance_schema.clone(),
3853 FileWriterOptions::default(),
3854 )
3855 .unwrap();
3856 writer
3858 .write_column(0, Arc::new(Int32Array::from(vec![1, 2, 3, 4, 5])))
3859 .await
3860 .unwrap();
3861 writer
3862 .write_column(1, Arc::new(Int32Array::from(vec![100])))
3863 .await
3864 .unwrap();
3865 writer.finish().await.unwrap();
3866
3867 let file_scheduler = fs
3868 .scheduler
3869 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3870 .await
3871 .unwrap();
3872 let cache = test_cache();
3873 let open_indexed = |names: &[&str]| {
3874 let projection = versions::reader_projection_from_column_names(
3875 ConcreteFileVersion::V2_1,
3876 &lance_schema,
3877 names,
3878 )
3879 .unwrap();
3880 ProjectedFileReader::try_open(
3881 file_scheduler.clone(),
3882 Some(projection),
3883 Arc::<DecoderPlugins>::default(),
3884 &cache,
3885 FileReaderOptions::default(),
3886 )
3887 };
3888
3889 let lazy = open_indexed(&["a", "c"]).await.unwrap();
3892 let err = match lazy
3895 .read_tasks(
3896 ReadBatchParams::RangeFull,
3897 1024,
3898 None,
3899 FilterExpression::no_filter(),
3900 )
3901 .await
3902 {
3903 Ok(_) => panic!("expected the mismatched-length projection to be rejected"),
3904 Err(e) => e.to_string(),
3905 };
3906 assert!(
3907 err.contains("a=5") && err.contains("c=1"),
3908 "error should name each column's length, got: {err}"
3909 );
3910
3911 let lazy = open_indexed(&["c"]).await.unwrap();
3914 let tasks = lazy
3915 .read_tasks(
3916 ReadBatchParams::RangeFull,
3917 1024,
3918 None,
3919 FilterExpression::no_filter(),
3920 )
3921 .await
3922 .unwrap();
3923 let batches = collect_read_tasks(tasks, 16).await;
3924 let values: Vec<Option<i32>> = batches
3925 .iter()
3926 .flat_map(|b| {
3927 b.column(0)
3928 .as_any()
3929 .downcast_ref::<Int32Array>()
3930 .unwrap()
3931 .iter()
3932 .collect::<Vec<_>>()
3933 })
3934 .collect();
3935 assert_eq!(values, vec![Some(100)]);
3936 }
3937
3938 #[test_log::test(tokio::test)]
3939 async fn test_compressing_buffer() {
3940 let fs = FsFixture::default();
3941
3942 let written_file = create_some_file(&fs, ConcreteFileVersion::V2_0).await;
3943 let file_scheduler = fs
3944 .scheduler
3945 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
3946 .await
3947 .unwrap();
3948
3949 let file_reader = FileReader::try_open(
3951 file_scheduler.clone(),
3952 None,
3953 Arc::<DecoderPlugins>::default(),
3954 &test_cache(),
3955 FileReaderOptions::default(),
3956 )
3957 .await
3958 .unwrap();
3959
3960 let mut projection = written_file.schema.project(&["score"]).unwrap();
3961 for field in projection.fields.iter_mut() {
3962 field
3963 .metadata
3964 .insert("lance:compression".to_string(), "zstd".to_string());
3965 }
3966 let projection = ReaderProjection {
3967 column_indices: projection.fields.iter().map(|f| f.id as u32).collect(),
3968 schema: Arc::new(projection),
3969 };
3970
3971 let batch_stream = file_reader
3972 .read_stream_projected(
3973 lance_io::ReadBatchParams::RangeFull,
3974 1024,
3975 16,
3976 projection.clone(),
3977 FilterExpression::no_filter(),
3978 )
3979 .await
3980 .unwrap();
3981
3982 let projection_arrow = Arc::new(ArrowSchema::from(projection.schema.as_ref()));
3983 verify_expected(
3984 &written_file.data,
3985 batch_stream,
3986 1024,
3987 Some(Box::new(move |batch: &RecordBatch| {
3988 batch.project_by_schema(&projection_arrow).unwrap()
3989 })),
3990 )
3991 .await;
3992 }
3993
3994 #[tokio::test]
3995 async fn test_read_all() {
3996 let fs = FsFixture::default();
3997 let WrittenFile { data, .. } = create_some_file(&fs, ConcreteFileVersion::V2_0).await;
3998 let total_rows = data.iter().map(|batch| batch.num_rows()).sum::<usize>();
3999
4000 let file_scheduler = fs
4001 .scheduler
4002 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4003 .await
4004 .unwrap();
4005 let file_reader = FileReader::try_open(
4006 file_scheduler.clone(),
4007 None,
4008 Arc::<DecoderPlugins>::default(),
4009 &test_cache(),
4010 FileReaderOptions::default(),
4011 )
4012 .await
4013 .unwrap();
4014
4015 let batches = file_reader
4016 .read_stream(
4017 lance_io::ReadBatchParams::RangeFull,
4018 total_rows as u32,
4019 16,
4020 FilterExpression::no_filter(),
4021 )
4022 .await
4023 .unwrap()
4024 .try_collect::<Vec<_>>()
4025 .await
4026 .unwrap();
4027 assert_eq!(batches.len(), 1);
4028 assert_eq!(batches[0].num_rows(), total_rows);
4029 }
4030
4031 #[rstest]
4032 #[tokio::test]
4033 async fn test_blocking_take(
4034 #[values(
4035 ConcreteFileVersion::V2_0,
4036 ConcreteFileVersion::V2_1,
4037 ConcreteFileVersion::V2_2
4038 )]
4039 version: ConcreteFileVersion,
4040 ) {
4041 let fs = FsFixture::default();
4042 let WrittenFile { data, schema, .. } = create_some_file(&fs, version).await;
4043 let total_rows = data.iter().map(|batch| batch.num_rows()).sum::<usize>();
4044
4045 let file_scheduler = fs
4046 .scheduler
4047 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4048 .await
4049 .unwrap();
4050 let file_reader = FileReader::try_open(
4051 file_scheduler.clone(),
4052 Some(
4053 versions::reader_projection_from_column_names(version, &schema, &["score"])
4054 .unwrap(),
4055 ),
4056 Arc::<DecoderPlugins>::default(),
4057 &test_cache(),
4058 FileReaderOptions::default(),
4059 )
4060 .await
4061 .unwrap();
4062
4063 let batches = tokio::task::spawn_blocking(move || {
4064 file_reader
4065 .read_stream_projected_blocking(
4066 lance_io::ReadBatchParams::Indices(UInt32Array::from(vec![0, 1, 2, 3, 4])),
4067 total_rows as u32,
4068 None,
4069 FilterExpression::no_filter(),
4070 )
4071 .unwrap()
4072 .collect::<ArrowResult<Vec<_>>>()
4073 .unwrap()
4074 })
4075 .await
4076 .unwrap();
4077
4078 assert_eq!(batches.len(), 1);
4079 assert_eq!(batches[0].num_rows(), 5);
4080 assert_eq!(batches[0].num_columns(), 1);
4081 }
4082
4083 #[tokio::test(flavor = "multi_thread")]
4084 async fn test_drop_in_progress() {
4085 let fs = FsFixture::default();
4086 let WrittenFile { data, .. } = create_some_file(&fs, ConcreteFileVersion::V2_0).await;
4087 let total_rows = data.iter().map(|batch| batch.num_rows()).sum::<usize>();
4088
4089 let file_scheduler = fs
4090 .scheduler
4091 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4092 .await
4093 .unwrap();
4094 let file_reader = FileReader::try_open(
4095 file_scheduler.clone(),
4096 None,
4097 Arc::<DecoderPlugins>::default(),
4098 &test_cache(),
4099 FileReaderOptions::default(),
4100 )
4101 .await
4102 .unwrap();
4103
4104 let mut batches = file_reader
4105 .read_stream(
4106 lance_io::ReadBatchParams::RangeFull,
4107 (total_rows / 10) as u32,
4108 16,
4109 FilterExpression::no_filter(),
4110 )
4111 .await
4112 .unwrap();
4113
4114 drop(file_reader);
4115
4116 let batch = batches.next().await.unwrap().unwrap();
4117 assert!(batch.num_rows() > 0);
4118
4119 drop(batches);
4121 }
4122
4123 #[tokio::test]
4124 async fn drop_while_scheduling() {
4125 let fs = FsFixture::default();
4135 let written_file = create_some_file(&fs, ConcreteFileVersion::V2_0).await;
4136 let total_rows = written_file
4137 .data
4138 .iter()
4139 .map(|batch| batch.num_rows())
4140 .sum::<usize>();
4141
4142 let file_scheduler = fs
4143 .scheduler
4144 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4145 .await
4146 .unwrap();
4147 let file_reader = FileReader::try_open(
4148 file_scheduler.clone(),
4149 None,
4150 Arc::<DecoderPlugins>::default(),
4151 &test_cache(),
4152 FileReaderOptions::default(),
4153 )
4154 .await
4155 .unwrap();
4156
4157 let projection = versions::reader_projection_from_whole_schema(
4158 &written_file.schema,
4159 ConcreteFileVersion::V2_0,
4160 );
4161 let column_infos = file_reader.metadata().column_infos.clone();
4162 let mut decode_scheduler = DecodeBatchScheduler::try_new(
4163 &projection.schema,
4164 &projection.column_indices,
4165 &column_infos,
4166 &vec![],
4167 total_rows as u64,
4168 Arc::<DecoderPlugins>::default(),
4169 file_reader.scheduler(),
4170 test_cache(),
4171 &FilterExpression::no_filter(),
4172 &DecoderConfig::default(),
4173 )
4174 .await
4175 .unwrap();
4176
4177 let range = 0..total_rows as u64;
4178
4179 let (tx, rx) = mpsc::unbounded_channel();
4180
4181 drop(rx);
4183
4184 decode_scheduler.schedule_range(
4186 range,
4187 &FilterExpression::no_filter(),
4188 tx,
4189 file_reader.scheduler(),
4190 )
4191 }
4192
4193 #[tokio::test]
4194 async fn test_read_empty_range() {
4195 let fs = FsFixture::default();
4196 create_some_file(&fs, ConcreteFileVersion::V2_0).await;
4197
4198 let file_scheduler = fs
4199 .scheduler
4200 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4201 .await
4202 .unwrap();
4203 let file_reader = FileReader::try_open(
4204 file_scheduler.clone(),
4205 None,
4206 Arc::<DecoderPlugins>::default(),
4207 &test_cache(),
4208 FileReaderOptions::default(),
4209 )
4210 .await
4211 .unwrap();
4212
4213 let batches = file_reader
4215 .read_stream(
4216 lance_io::ReadBatchParams::Range(0..0),
4217 1024,
4218 16,
4219 FilterExpression::no_filter(),
4220 )
4221 .await
4222 .unwrap()
4223 .try_collect::<Vec<_>>()
4224 .await
4225 .unwrap();
4226
4227 assert_eq!(batches.len(), 0);
4228
4229 let batches = file_reader
4231 .read_stream(
4232 lance_io::ReadBatchParams::Ranges(Arc::new([0..1, 2..2])),
4233 1024,
4234 16,
4235 FilterExpression::no_filter(),
4236 )
4237 .await
4238 .unwrap()
4239 .try_collect::<Vec<_>>()
4240 .await
4241 .unwrap();
4242 assert_eq!(batches.len(), 1);
4243 }
4244
4245 async fn write_file_with_global_buffer(fs: &FsFixture, buffer: Bytes) {
4246 let lance_schema =
4247 lance_core::datatypes::Schema::try_from(&ArrowSchema::new(vec![Field::new(
4248 "foo",
4249 DataType::Int32,
4250 true,
4251 )]))
4252 .unwrap();
4253
4254 let mut file_writer = versions::v2_1::create_writer(
4255 fs.object_store.create(&fs.tmp_path).await.unwrap(),
4256 lance_schema,
4257 FileWriterOptions::default(),
4258 )
4259 .unwrap();
4260
4261 let buf_index = file_writer.add_global_buffer(buffer).await.unwrap();
4262 assert_eq!(buf_index, 1);
4263
4264 file_writer.finish().await.unwrap();
4265 }
4266
4267 #[derive(Clone, Copy, Debug)]
4268 enum MetadataReadPath {
4269 Full,
4270 Indexed,
4271 }
4272
4273 #[derive(Clone, Copy, Debug)]
4274 enum InvalidGboDescriptor {
4275 Unaligned,
4276 PastEof,
4277 Overflowing,
4278 }
4279
4280 #[rstest]
4281 #[case::full_unaligned(MetadataReadPath::Full, InvalidGboDescriptor::Unaligned, "not aligned")]
4282 #[case::full_past_eof(MetadataReadPath::Full, InvalidGboDescriptor::PastEof, "outside file")]
4283 #[case::full_overflowing(
4284 MetadataReadPath::Full,
4285 InvalidGboDescriptor::Overflowing,
4286 "overflows"
4287 )]
4288 #[case::indexed_unaligned(
4289 MetadataReadPath::Indexed,
4290 InvalidGboDescriptor::Unaligned,
4291 "not aligned"
4292 )]
4293 #[case::indexed_past_eof(
4294 MetadataReadPath::Indexed,
4295 InvalidGboDescriptor::PastEof,
4296 "outside file"
4297 )]
4298 #[case::indexed_overflowing(
4299 MetadataReadPath::Indexed,
4300 InvalidGboDescriptor::Overflowing,
4301 "overflows"
4302 )]
4303 #[tokio::test]
4304 async fn test_metadata_rejects_invalid_gbo_descriptor(
4305 #[case] read_path: MetadataReadPath,
4306 #[case] invalid_descriptor: InvalidGboDescriptor,
4307 #[case] expected_message: &str,
4308 ) {
4309 let fs = FsFixture::default();
4310 write_file_with_global_buffer(&fs, Bytes::from_static(b"hello")).await;
4311
4312 let mut file_bytes = fs
4313 .object_store
4314 .read_one_all(&fs.tmp_path)
4315 .await
4316 .unwrap()
4317 .to_vec();
4318 let file_len = file_bytes.len() as u64;
4319 let footer = FileReader::decode_footer(&Bytes::copy_from_slice(&file_bytes)).unwrap();
4320 let gbo_table_start = usize::try_from(footer.global_buff_offsets_start).unwrap();
4321 let alignment = PAGE_BUFFER_ALIGNMENT as u64;
4322 let (position, size) = match invalid_descriptor {
4323 InvalidGboDescriptor::Unaligned => (1, 0),
4324 InvalidGboDescriptor::PastEof => (((file_len + alignment) / alignment) * alignment, 0),
4325 InvalidGboDescriptor::Overflowing => (u64::MAX - (u64::MAX % alignment), alignment),
4326 };
4327 file_bytes[gbo_table_start..gbo_table_start + 8].copy_from_slice(&position.to_le_bytes());
4328 file_bytes[gbo_table_start + 8..gbo_table_start + 16].copy_from_slice(&size.to_le_bytes());
4329 fs.object_store
4330 .put(&fs.tmp_path, &file_bytes)
4331 .await
4332 .unwrap();
4333
4334 let scheduler = fs
4335 .scheduler
4336 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4337 .await
4338 .unwrap();
4339 let error = match read_path {
4340 MetadataReadPath::Full => FileReader::read_all_metadata(&scheduler).await.map(|_| ()),
4341 MetadataReadPath::Indexed => FileReader::read_metadata_index(&scheduler)
4342 .await
4343 .map(|_| ()),
4344 }
4345 .expect_err("invalid GBO descriptor must fail before metadata I/O");
4346
4347 assert!(
4348 matches!(error, lance_core::Error::InvalidInput { .. }),
4349 "expected InvalidInput, got {error:?}"
4350 );
4351 assert!(
4352 error.to_string().contains(expected_message),
4353 "unexpected error: {error}"
4354 );
4355 }
4356
4357 #[rstest]
4361 #[case::within_tail_window(true)]
4362 #[case::outside_tail_window(false)]
4363 #[tokio::test]
4364 async fn test_read_global_buffer(#[case] within_window: bool) {
4365 let fs = FsFixture::default();
4366
4367 let block_size = fs.object_store.block_size();
4368 let buffer = if within_window {
4369 Bytes::from_static(b"hello")
4370 } else {
4371 Bytes::from(vec![7u8; 2 * block_size])
4372 };
4373 let expected_read_iops = if within_window { 0 } else { 1 };
4374
4375 write_file_with_global_buffer(&fs, buffer.clone()).await;
4376
4377 let file_scheduler = fs
4378 .scheduler
4379 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4380 .await
4381 .unwrap();
4382 let file_reader = FileReader::try_open(
4383 file_scheduler,
4384 None,
4385 Arc::<DecoderPlugins>::default(),
4386 &test_cache(),
4387 FileReaderOptions::default(),
4388 )
4389 .await
4390 .unwrap();
4391
4392 let retained = &file_reader.metadata().retained_global_buffers;
4395 assert!(!retained.contains_key(&0), "schema must not be retained");
4396 assert_eq!(retained.contains_key(&1), within_window);
4397
4398 fs.object_store.io_stats_incremental();
4400
4401 let buf = file_reader.read_global_buffer(1).await.unwrap();
4402 assert_eq!(buf, buffer);
4403
4404 let stats = fs.object_store.io_stats_incremental();
4405 assert_eq!(stats.read_iops, expected_read_iops);
4406 }
4407
4408 #[tokio::test]
4411 async fn test_read_global_buffer_no_user_buffers() {
4412 let fs = FsFixture::default();
4413 create_some_file(&fs, ConcreteFileVersion::V2_1).await;
4414
4415 let file_scheduler = fs
4416 .scheduler
4417 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4418 .await
4419 .unwrap();
4420 let file_reader = FileReader::try_open(
4421 file_scheduler,
4422 None,
4423 Arc::<DecoderPlugins>::default(),
4424 &test_cache(),
4425 FileReaderOptions::default(),
4426 )
4427 .await
4428 .unwrap();
4429
4430 let metadata = file_reader.metadata();
4431 assert_eq!(metadata.file_buffers.len(), 1, "expected only the schema");
4432 assert!(
4433 metadata.retained_global_buffers.is_empty(),
4434 "a file with no user global buffers must retain nothing"
4435 );
4436 }
4437
4438 #[rstest]
4439 #[tokio::test]
4440 async fn test_deep_size_of_includes_column_metadata(
4441 #[values(
4442 ConcreteFileVersion::V2_0,
4443 ConcreteFileVersion::V2_1,
4444 ConcreteFileVersion::V2_2,
4445 ConcreteFileVersion::V2_3
4446 )]
4447 version: ConcreteFileVersion,
4448 ) {
4449 use lance_core::deepsize::DeepSizeOf;
4454
4455 let fs = FsFixture::default();
4456 let _written = create_some_file(&fs, version).await;
4457 let cache = test_cache();
4458 let file_scheduler = fs
4459 .scheduler
4460 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4461 .await
4462 .unwrap();
4463 let file_reader = FileReader::try_open(
4464 file_scheduler,
4465 None,
4466 Arc::<DecoderPlugins>::default(),
4467 &cache,
4468 FileReaderOptions::default(),
4469 )
4470 .await
4471 .unwrap();
4472
4473 let metadata = file_reader.metadata();
4474 let deep_size = metadata.deep_size_of();
4475
4476 assert!(
4482 deep_size > 1024,
4483 "deep_size_of ({deep_size}) is suspiciously small — \
4484 column_metadatas and column_infos may not be accounted for"
4485 );
4486
4487 assert!(
4489 !metadata.column_metadatas.is_empty(),
4490 "Expected non-empty column_metadatas"
4491 );
4492
4493 let num_columns = metadata.column_metadatas.len();
4496 assert!(
4497 deep_size > num_columns * 50,
4498 "deep_size_of ({deep_size}) should scale with column count ({num_columns})"
4499 );
4500 }
4501
4502 #[tokio::test]
4503 async fn test_read_global_buffer_out_of_range() {
4504 let fs = FsFixture::default();
4505
4506 write_file_with_global_buffer(&fs, Bytes::from_static(b"hello")).await;
4507
4508 let file_scheduler = fs
4509 .scheduler
4510 .open_file(&fs.tmp_path, &CachedFileSize::unknown())
4511 .await
4512 .unwrap();
4513 let file_reader = FileReader::try_open(
4514 file_scheduler,
4515 None,
4516 Arc::<DecoderPlugins>::default(),
4517 &test_cache(),
4518 FileReaderOptions::default(),
4519 )
4520 .await
4521 .unwrap();
4522
4523 let err = file_reader.read_global_buffer(2).await.unwrap_err();
4526 assert!(
4527 matches!(err, lance_core::Error::InvalidInput { .. }),
4528 "expected InvalidInput, got: {err:?}"
4529 );
4530 let msg = err.to_string();
4531 assert!(msg.contains('2'), "error should mention the index: {msg}");
4532 }
4533
4534 #[rstest]
4539 fn test_validate_struct_child_lengths(#[values(false, true)] is_structural: bool) {
4540 let run = |dt: DataType, indices: &[u32], lengths: Vec<u64>| -> lance_core::Result<u64> {
4541 let arrow = ArrowSchema::new(vec![Field::new("s", dt, true)]);
4542 let schema = Schema::try_from(&arrow).unwrap();
4543 if is_structural {
4544 versions::v2_1::test_projection_length(&schema, indices, &lengths)
4545 } else {
4546 versions::v2_0::test_projection_length(&schema, indices, &lengths)
4547 }
4548 };
4549
4550 let struct_ty = || {
4551 DataType::Struct(Fields::from(vec![
4552 Field::new("a", DataType::Int32, true),
4553 Field::new("b", DataType::Int32, true),
4554 ]))
4555 };
4556
4557 let (indices, equal, unequal): (&[u32], Vec<u64>, Vec<u64>) = if is_structural {
4560 (&[0, 1], vec![5, 5], vec![5, 3])
4561 } else {
4562 (&[0, 1, 2], vec![5, 5, 5], vec![5, 5, 3])
4563 };
4564
4565 assert_eq!(run(struct_ty(), indices, equal).unwrap(), 5);
4566
4567 let err = run(struct_ty(), indices, unequal).unwrap_err();
4568 let msg = err.to_string();
4569 assert!(
4570 msg.contains("differing lengths") && msg.contains('b'),
4571 "expected a child-length error naming 'b', got: {msg}"
4572 );
4573 }
4574
4575 #[test]
4576 fn test_validate_v2_0_unloaded_blob_projection_is_opaque() {
4577 let metadata = HashMap::from([(BLOB_META_KEY.to_string(), "true".to_string())]);
4578 let arrow = ArrowSchema::new(vec![
4579 Field::new("blob", DataType::LargeBinary, true).with_metadata(metadata),
4580 ]);
4581 let mut schema = Schema::try_from(&arrow).unwrap();
4582 schema.fields[0].unloaded_mut();
4583 let projection = ReaderProjection {
4584 schema: Arc::new(schema),
4585 column_indices: vec![0],
4586 };
4587 let rows = versions::v2_0::test_projection_length(
4588 &projection.schema,
4589 &projection.column_indices,
4590 &[3],
4591 )
4592 .unwrap();
4593
4594 assert_eq!(rows, 3);
4595 }
4596
4597 #[test]
4598 fn test_validate_length_list_and_empty_struct() {
4599 let validate = |dt: DataType,
4600 is_structural: bool,
4601 indices: &[u32],
4602 lengths: Vec<u64>|
4603 -> lance_core::Result<u64> {
4604 let arrow = ArrowSchema::new(vec![Field::new("f", dt, true)]);
4605 let schema = Schema::try_from(&arrow).unwrap();
4606 if is_structural {
4607 versions::v2_1::test_projection_length(&schema, indices, &lengths)
4608 } else {
4609 versions::v2_0::test_projection_length(&schema, indices, &lengths)
4610 }
4611 };
4612
4613 let list_ty = DataType::List(Arc::new(Field::new("item", DataType::Int32, true)));
4618 assert_eq!(
4619 validate(list_ty.clone(), false, &[0, 1], vec![5, 17]).unwrap(),
4620 5
4621 );
4622 assert_eq!(validate(list_ty, true, &[0], vec![5]).unwrap(), 5);
4623
4624 let list_of_struct = DataType::List(Arc::new(Field::new(
4631 "item",
4632 DataType::Struct(Fields::from(vec![
4633 Field::new("a", DataType::Int32, true),
4634 Field::new("b", DataType::Int32, true),
4635 ])),
4636 true,
4637 )));
4638 assert_eq!(
4639 validate(
4640 list_of_struct.clone(),
4641 false,
4642 &[0, 1, 2, 3],
4643 vec![6, 6, 29, 29]
4644 )
4645 .unwrap(),
4646 6
4647 );
4648 assert_eq!(
4650 validate(list_of_struct, true, &[0, 1], vec![29, 29]).unwrap(),
4651 29
4652 );
4653
4654 let empty_struct = DataType::Struct(Fields::empty());
4656 assert_eq!(
4657 validate(empty_struct.clone(), false, &[0], vec![9]).unwrap(),
4658 9
4659 );
4660 assert_eq!(validate(empty_struct, true, &[0], vec![9]).unwrap(), 9);
4661 }
4662}