Skip to main content

lance_file/
reader.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use 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
56/// Default chunk size for reading large pages (8MiB)
57/// Pages larger than this will be split into multiple chunks during read
58pub const DEFAULT_READ_CHUNK_SIZE: u64 = 8 * 1024 * 1024;
59
60// For now, we don't use global buffers for anything other than schema.  If we
61// use these later we should make them lazily loaded and then cached once loaded.
62//
63// We store their position / length for debugging purposes
64#[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/// Statistics summarize some of the file metadata for quick summary info
104#[derive(Debug)]
105pub struct FileStatistics {
106    /// Statistics about each of the columns in the file
107    pub columns: Vec<ColumnStatistics>,
108}
109
110/// Summary information describing a column
111#[derive(Debug)]
112pub struct ColumnStatistics {
113    /// The number of pages in the column
114    pub num_pages: usize,
115    /// The total number of data & metadata bytes in the column
116    ///
117    /// This is the compressed on-disk size
118    pub size_bytes: u64,
119}
120
121// TODO: Caching
122#[derive(Debug)]
123pub struct CachedFileMetadata {
124    /// The schema of the file
125    pub file_schema: Arc<Schema>,
126    /// The column metadatas
127    pub column_metadatas: Vec<pbfile::ColumnMetadata>,
128    pub column_infos: Vec<Arc<ColumnInfo>>,
129    /// The number of rows in the file
130    pub num_rows: u64,
131    pub file_buffers: Vec<BufferDescriptor>,
132    /// The number of bytes contained in the data page section of the file
133    pub num_data_bytes: u64,
134    /// The number of bytes contained in the column metadata (not including buffers
135    /// referenced by the metadata)
136    pub num_column_metadata_bytes: u64,
137    /// The number of bytes contained in global buffers
138    pub num_global_buffer_bytes: u64,
139    /// The number of bytes contained in the CMO and GBO tables
140    pub num_footer_bytes: u64,
141    /// The major version number stored in the file footer.
142    pub major_version: u16,
143    /// The minor version number stored in the file footer.
144    pub minor_version: u16,
145    pub version: ConcreteFileVersion,
146    /// The actual total file size in bytes, as reported by the object store.
147    pub file_size_bytes: u64,
148    /// User global buffers (index >= 1) whose bytes were already captured by the
149    /// tail read that `read_all_metadata` performs at open, keyed by buffer index.
150    ///
151    /// All global buffers are laid out contiguously starting at the schema, so on
152    /// small/medium files they land inside the captured tail window. Retaining
153    /// those bytes lets `read_global_buffer` serve them with zero additional I/O.
154    /// The bytes are copied out of the tail (rather than sliced) so the much
155    /// larger tail allocation can be dropped — we only hold what we will serve.
156    ///
157    /// The schema (buffer 0) is excluded: it is already decoded at open and is
158    /// not fetched through `read_global_buffer`. Buffers that fall outside the
159    /// window (large files) are absent here and fall back to a dedicated read.
160    pub retained_global_buffers: BTreeMap<u32, Bytes>,
161}
162
163impl CachedFileMetadata {
164    /// Total file size in bytes.
165    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        // column_metadatas is Vec<pbfile::ColumnMetadata> (protobuf generated,
189        // does not implement DeepSizeOf). We use prost::Message::encoded_len()
190        // as a proxy for in-memory size. The decoded representation is typically
191        // several times larger than the wire format due to heap-allocated
192        // repeated/string/bytes fields, so we apply a 4x multiplier.
193        let column_metadatas_size = column_metadata_deep_size(self.column_metadatas.as_slice());
194
195        // column_infos is Vec<Arc<ColumnInfo>>. Each ColumnInfo contains
196        // page_infos (with protobuf PageEncoding), buffer offsets, and a
197        // column-level ColumnEncoding protobuf.
198        let column_infos_size = self.column_infos.deep_size_of_children(context);
199
200        // Global buffer bytes retained for zero-IO reads (copied out of the tail).
201        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/// Lightweight file metadata used to locate per-column metadata on demand.
212///
213/// This contains the file-level schema, row count, global buffer descriptors,
214/// and column metadata offset table. Unlike [`CachedFileMetadata`], it does not
215/// hold decoded metadata for every column.
216#[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    /// Returns the total size of the file in bytes.
230    pub fn file_size(&self) -> u64 {
231        self.file_size_bytes
232    }
233
234    /// Returns the number of physical columns in the file.
235    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/// Selecting columns from a lance file requires specifying both the
285/// index of the column and the data type of the column
286///
287/// Partly, this is because it is not strictly required that columns
288/// be read into the same type.  For example, a string column may be
289/// read as a string, large_string or string_view type.
290///
291/// A read will only succeed if the decoder for a column is capable
292/// of decoding into the requested type.
293///
294/// Note that this should generally be limited to different in-memory
295/// representations of the same semantic type.  An encoding could
296/// theoretically support "casting" (e.g. int to string, etc.) but
297/// there is little advantage in doing so here.
298///
299/// Note: in order to specify a projection the user will need some way
300/// to figure out the column indices.  In the table format we do this
301/// using field IDs and keeping track of the field id->column index mapping.
302///
303/// If users are not using the table format then they will need to figure
304/// out some way to do this themselves.
305#[derive(Debug, Clone)]
306pub struct ReaderProjection {
307    /// The data types (schema) of the selected columns.  The names
308    /// of the schema are arbitrary and ignored.
309    pub schema: Arc<Schema>,
310    /// The indices of the columns to load.
311    ///
312    /// The content of this vector depends on the file version.
313    ///
314    /// In Lance File Version 2.0 we need ids for structural fields as
315    /// well as leaf fields:
316    ///
317    ///   - Primitive: the index of the column in the schema
318    ///   - List: the index of the list column in the schema
319    ///     followed by the column indices of the children
320    ///   - FixedSizeList (of primitive): the index of the column in the schema
321    ///     (this case is not nested)
322    ///   - FixedSizeList (of non-primitive): not yet implemented
323    ///   - Dictionary: same as primitive
324    ///   - Struct: the index of the struct column in the schema
325    ///     followed by the column indices of the children
326    ///
327    ///   In other words, this should be a DFS listing of the desired schema.
328    ///
329    /// In Lance File Version 2.1 we only need ids for leaf fields.  Any structural
330    /// fields are completely transparent.
331    ///
332    /// For example, if the goal is to load:
333    ///
334    ///   x: int32
335    ///   y: `struct<z: int32, w: string>`
336    ///   z: `list<int32>`
337    ///
338    /// and the schema originally used to store the data was:
339    ///
340    ///   a: `struct<x: int32>`
341    ///   b: int64
342    ///   y: `struct<z: int32, c: int64, w: string>`
343    ///   z: `list<int32>`
344    ///
345    /// Then the column_indices should be:
346    ///
347    /// - 2.0: [1, 3, 4, 6, 7, 8]
348    /// - 2.1: [0, 2, 4, 5]
349    pub column_indices: Vec<u32>,
350}
351
352impl ReaderProjection {
353    /// Returns whether this projection is selective enough to benefit from
354    /// loading column metadata through the file's metadata index.
355    ///
356    /// The caller must already have selected a file format that supports indexed
357    /// metadata. This method only evaluates the projection shape and selectivity.
358    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/// File Reader Options that can control reading behaviors, such as whether to enable caching on repetition indices
365#[derive(Clone, Debug)]
366pub struct FileReaderOptions {
367    pub decoder_config: DecoderConfig,
368    /// Size of chunks when reading large pages. Pages larger than this
369    /// will be read in multiple chunks to control memory usage.
370    /// Default: 8MB (DEFAULT_READ_CHUNK_SIZE)
371    pub read_chunk_size: u64,
372    /// If set, the reader will produce batches whose total size in bytes
373    /// is approximately this value. The row-based `batch_size` remains an
374    /// independent upper bound, and the limit reached first determines the batch size.
375    ///
376    /// This can be set at the dataset level (via `ReadParams::file_reader_options`)
377    /// to provide a default for all scans, or at the scanner level (via
378    /// `Scanner::batch_size_bytes`) to override per scan.
379    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/// Executable projection behavior selected by an exact file-version module.
405///
406/// The shared reader invokes this behavior but never interprets a version or
407/// accepted-grammar profile.
408#[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/// A projection-scoped reader for a current-format Lance file.
439///
440/// This reader fixes a base projection at construction time. All later reads
441/// must stay within that projection, which lets the reader load only the column
442/// metadata needed by the base projection when opening from a [`FileMetadataIndex`].
443/// It intentionally does not expose APIs that require synchronous access to full
444/// file metadata.
445#[derive(Debug, Clone)]
446pub struct ProjectedFileReader {
447    core: DecodeEngine,
448}
449
450/// A current-format Lance file reader backed by fully decoded metadata.
451#[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    // We don't use this today because we always load metadata for every column
500    // and don't yet support "metadata projection"
501    #[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
512// Count the V2.1 physical columns required to reconstruct a projected field.
513// This is the same DFS shape consumed by `ColumnInfoIter`: ordinary structural
514// nodes are transparent and leaves contribute columns. Indexed metadata loading
515// can therefore compact any ordinary structural projection into 0..N while
516// preserving this order.
517//
518// Blob and packed-struct fields remain unsupported by indexed projection. Their
519// opaque decode semantics are handled by the existing full-metadata reader.
520fn 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
534// The reader combines a projection's columns into rectangular batches, so they
535// must all have the same length.  Returns that common length, or a descriptive
536// error (naming each column's length) when they differ. Ordinary files always
537// pass; only files written with `FileWriter::write_column` whose columns ended up
538// unequal can fail, and those must be read separately.
539pub(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    /// Returns a clone of this reader whose I/O is additionally recorded into
612    /// `stats`, on top of the scheduler's global accounting.
613    ///
614    /// All cached metadata is shared with `self`, so no file is re-opened and
615    /// only a few `Arc` clones are performed.  If the underlying I/O service
616    /// does not support per-scope statistics (e.g. an in-memory scheduler), the
617    /// returned reader is an ordinary, uninstrumented clone.
618    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    /// The number of rows stored in a single physical column.
633    ///
634    /// For ordinary (rectangular) files every column has the same length, equal
635    /// to [`num_rows`](Self::num_rows). Files written with
636    /// [`FileWriter::write_column`](crate::writer::FileWriter::write_column)
637    /// may have columns of differing lengths; this returns the length of one
638    /// such column, derived by summing its pages' row counts. Errors if
639    /// `column_index` is out of bounds.
640    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    // Checks to make sure the footer is written correctly and returns the
741    // position of the file descriptor (which comes from the footer)
742    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    // TODO: Once we have coalesced I/O we should only read the column metadatas that we need
795    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        // cmo == column_metadata_offsets
801        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        // This could, in theory, trigger another IOP but the GBO table should never be large
929        // enough for that to happen
930        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        // The tail read above already pulled in any global buffer that lives within
1002        // the captured window. Copy those user buffers (index >= 1; the schema at 0
1003        // is decoded above and never fetched via read_global_buffer) out of the tail
1004        // so read_global_buffer can serve them without I/O. We copy rather than slice
1005        // so the much larger tail allocation can be released once decoding is done.
1006        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    /// Reads the lightweight metadata index from a file.
1094    ///
1095    /// This reads the file schema from the schema global buffer. Use
1096    /// [`Self::read_metadata_index_with_schema`] when the caller already has
1097    /// the schema and row count from a higher-level metadata source.
1098    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    /// Reads the metadata index without fetching the schema global buffer.
1105    ///
1106    /// Use this when the caller already has the file schema and physical row
1107    /// count from an enclosing metadata layer, such as a dataset manifest.
1108    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    // The actual decoder needs all the column infos that make up a type.  In other words, if
1147    // the first type in the schema is Struct<i32, i32> then the decoder will need 3 column infos.
1148    //
1149    // This is a file reader concern because the file reader needs to support late projection of columns
1150    // and so it will need to figure this out anyways.
1151    //
1152    // It's a bit of a tricky process though because the number of column infos may depend on the
1153    // encoding.  Considering the above example, if we wrote it with a packed encoding, then there would
1154    // only be a single column in the file (and not 3).
1155    //
1156    // At the moment this method words because our rules are simple and we just repeat them here.  See
1157    // Self::default_projection for a similar problem.  In the future this is something the encodings
1158    // registry will need to figure out.
1159    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 &params {
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(&params, 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(&params, 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(&params, 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(&params, 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(&params, 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 &params {
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(&params, 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(&params, 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(&params, 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(&params, 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(&params, 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    /// Returns a clone of this reader using a different scheduler.
1907    pub fn with_scheduler(&self, scheduler: Arc<dyn EncodingsIo>) -> Self {
1908        Self {
1909            core: self.core.with_scheduler(scheduler),
1910        }
1911    }
1912
1913    /// Returns the number of rows in the file.
1914    pub fn num_rows(&self) -> u64 {
1915        self.core.num_rows()
1916    }
1917
1918    /// Returns the file schema visible to this reader.
1919    pub fn schema(&self) -> &Arc<Schema> {
1920        self.core.schema()
1921    }
1922
1923    /// Returns file statistics when this reader has full file metadata.
1924    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    /// Reads a global buffer by index.
1937    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
2286/// Inspects a page and returns a String describing the page's encoding
2287pub 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        // Next, read the metadata for the columns
2354        // This is both the column metadata and the CMO table
2355        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        // Next, read the metadata for the columns
2402        // This is both the column metadata and the CMO table
2403        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    /// Writes `batch` to a fresh file, overwrites `patch` bytes at `patch_offset`
2808    /// into the single occurrence of `pattern`, and reads the file back with the
2809    /// default reader configuration.
2810    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        // V2.1 places this column's first mini-block metadata word at byte zero.
2893        // The issue's mutation makes its non-final item count exceed the page total.
2894        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    /// A corrupt file whose variable-width offsets point outside the value bytes
2937    /// must fail with a typed error under the default reader configuration
2938    /// (`validate_on_decode` disabled) instead of materializing values outside
2939    /// the data buffer.
2940    ///
2941    /// Uses a dictionary-encoded string column because its values page stores
2942    /// the offsets verbatim, so flipping the tail offset in the file reaches the
2943    /// Arrow conversion boundary without being rejected by an intermediate
2944    /// decompressor.
2945    #[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        // The dictionary values page stores the value offsets as plain
2964        // little-endian i32s ending with [5, 9, 14] (2.1 also stores the leading
2965        // zero, 2.2+ omits it).  If a future encoding change stops storing these
2966        // offsets verbatim this lookup fails loudly and the test needs a new
2967        // byte pattern.  The patch rewrites the tail offset so it points far
2968        // beyond the value bytes.
2969        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    /// Storage dictionaries expand their values through `DataBlockBuilder`
2993    /// before the final Arrow layout validation.  Corrupt dictionary offsets
2994    /// must therefore fail at the append boundary instead of reaching a slice
2995    /// operation with a decreasing range.
2996    #[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    /// Same contract as the test above, but for a plain (non-dictionary) string
3051    /// column: the mini-block chunk stores chunk-relative value offsets that are
3052    /// used to slice the chunk, so a corrupt tail offset must surface as a typed
3053    /// error from the chunk decompressor instead of a panic in the decode task.
3054    #[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        // For ["alpha", "beta", "gamma"] the chunk stores LE i32 offsets
3074        // [16, 21, 25, 30] (chunk-relative: a 16-byte offsets region precedes
3075        // the value bytes).  The patch rewrites the tail offset to point far
3076        // past the chunk.
3077        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        // TODO: Add V2_1 (currently fails)
3167        #[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        // Test self described
3196        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        // Test mini
3215        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                // We can specify the projection as part of the read operation via read_stream_projected
3276                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                // We can also specify the projection as a base projection when we open the file
3326                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    // The projection-length validation lives in `DecodeEngine`, shared by the
3833    // eager and the lazy (indexed) metadata providers. The indexed provider loads
3834    // only the projected columns and renumbers them 0..N, so this checks that the
3835    // renumbered `column_infos`/`column_indices` still line up for the length
3836    // check: a mismatched-length projection is rejected through
3837    // `ProjectedFileReader`, and a single short column resolves to its own length.
3838    #[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        // "a" has 5 rows, "c" has 1 -- an unequal-length file.
3857        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        // A mismatched-length projection [a, c] (5 vs 1) is rejected at read time,
3890        // through the indexed provider's renumbered column infos.
3891        let lazy = open_indexed(&["a", "c"]).await.unwrap();
3892        // (`read_tasks` yields a stream, which is not `Debug`, so match rather
3893        // than `unwrap_err`.)
3894        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        // A single short column resolves to its own length (1), not the file's
3912        // longest column.
3913        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        // We can specify the projection as part of the read operation via read_stream_projected
3950        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 in-progress scan
4120        drop(batches);
4121    }
4122
4123    #[tokio::test]
4124    async fn drop_while_scheduling() {
4125        // This is a bit of a white-box test, pokes at the internals.  We want to
4126        // test the case where the read stream is dropped before the scheduling
4127        // thread finishes.  We can't do that in a black-box fashion because the
4128        // scheduling thread runs in the background and there is no easy way to
4129        // pause / gate it.
4130
4131        // It's a regression for a bug where the scheduling thread would panic
4132        // if the stream was dropped before it finished.
4133
4134        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        // Simulate the stream / decoder being dropped
4182        drop(rx);
4183
4184        // Scheduling should not panic
4185        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        // All ranges empty, no data
4214        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        // Some ranges empty
4230        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    /// A global buffer that fits inside the tail region captured at open is served
4358    /// from memory with no additional I/O.  A buffer larger than that window cannot
4359    /// fit and falls back to a dedicated read.  Both must round-trip correctly.
4360    #[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        // The user buffer should be retained only when it fits the tail window, and
4393        // the schema (buffer 0) is never retained.
4394        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        // Reset the IO counters so we only measure the read_global_buffer call.
4399        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    /// A file whose only global buffer is the schema (i.e. a plain data file, the
4409    /// common case) must retain nothing — there is no user buffer to serve.
4410    #[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        // Regression test: CachedFileMetadata::deep_size_of must account for
4450        // column_metadatas and column_infos, otherwise the moka cache weigher
4451        // dramatically underestimates entry sizes and never evicts, causing
4452        // unbounded memory growth on random-access workloads.
4453        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        // The file has multiple columns (score, location, categories, binary,
4477        // maybe large_bin). The deep_size_of must be substantially more than
4478        // just the schema — it should include column_metadatas + column_infos.
4479        // A naive implementation that ignores these fields reports < 1 KB;
4480        // a correct one should report at least several KB for this test file.
4481        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        // Verify column_metadatas is non-empty (sanity check).
4488        assert!(
4489            !metadata.column_metadatas.is_empty(),
4490            "Expected non-empty column_metadatas"
4491        );
4492
4493        // Verify the size scales with the number of columns: a file with more
4494        // columns should have a larger deep_size_of.
4495        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        // The file has two global buffers (schema at 0, "hello" at 1); index 2 is
4524        // out of range and must surface a descriptive error rather than panicking.
4525        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    // Exercises the projection length-validation walk in isolation, feeding
4535    // synthetic per-column lengths so we can reach cases no current writer can
4536    // actually produce -- in particular a struct whose children diverge in
4537    // length, which the decoders would otherwise panic on or misread.
4538    #[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        // In 2.1 a struct contributes no column of its own (just its two leaves);
4558        // in 2.0 it also has its own column first.
4559        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        // A list's items have a different cardinality than its rows; that gap
4614        // must not be flagged as a mismatch. In 2.0 the list is an offsets column
4615        // (rows) plus an items column (item count); in 2.1 it is a single column
4616        // whose page rows are the top-level row count.
4617        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        // A list of structs: the struct's children sit below the list boundary, so
4625        // their (item-count) lengths must NOT be compared against the list's row
4626        // count. In 2.0 each field has a column [list, struct, a, b] = [6, 6, 29,
4627        // 29]; the list resolves to its own row count (6) and the struct's longer
4628        // children are not flagged. (Regression: this previously errored because
4629        // the nested struct's children were checked against the struct's count.)
4630        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        // In 2.1 only the two leaves carry columns; still no false mismatch.
4649        assert_eq!(
4650            validate(list_of_struct, true, &[0, 1], vec![29, 29]).unwrap(),
4651            29
4652        );
4653
4654        // An empty struct contributes a single column and validates to its length.
4655        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}