Skip to main content

iceberg/arrow/reader/
pipeline.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! The main `ArrowReader` pipeline: reading a stream of `FileScanTask`s,
19//! opening Parquet files and resolving schemas, then wiring projection,
20//! predicates, row-group / row selection, and delete handling into a stream
21//! of transformed Arrow `RecordBatch`es.
22
23use std::collections::HashMap;
24use std::sync::Arc;
25use std::sync::atomic::AtomicU64;
26
27use arrow_schema::{DataType, Field};
28use futures::{StreamExt, TryStreamExt};
29use parquet::arrow::arrow_reader::{ArrowReaderMetadata, ArrowReaderOptions};
30use parquet::arrow::{
31    PARQUET_FIELD_ID_META_KEY, ParquetRecordBatchStreamBuilder, ProjectionMask, RowNumber,
32};
33use parquet::encryption::decrypt::FileDecryptionProperties;
34
35use super::row_lineage::synthesize_row_id_column;
36use super::{
37    ArrowFileReader, ArrowReader, ParquetReadOptions, add_fallback_field_ids_to_arrow_schema,
38    apply_name_mapping_to_arrow_schema, find_leaf_by_field_id,
39};
40use crate::arrow::build_partition_constant;
41use crate::arrow::caching_delete_file_loader::CachingDeleteFileLoader;
42use crate::arrow::int96::coerce_int96_timestamps;
43use crate::arrow::record_batch_transformer::RecordBatchTransformerBuilder;
44use crate::arrow::scan_metrics::{CountingFileRead, ScanMetrics, ScanResult};
45use crate::encryption::StandardKeyMetadata;
46use crate::error::{Result, invalid_data};
47use crate::expr::BoundPredicate;
48use crate::expr::visitors::bloom_filter_evaluator::{
49    BloomFilterEvaluator, ColumnBloomFilter, collect_bloom_filter_field_ids,
50};
51use crate::io::{FileIO, FileMetadata, FileRead};
52use crate::metadata_columns::{
53    RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER, RESERVED_COL_NAME_POS,
54    RESERVED_COL_NAME_ROW_ID, RESERVED_FIELD_ID_FILE,
55    RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER, RESERVED_FIELD_ID_PARTITION,
56    RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID, RESERVED_FIELD_ID_SPEC_ID, is_metadata_field,
57};
58use crate::scan::{ArrowRecordBatchStream, FileScanTask, FileScanTaskStream};
59use crate::spec::{Datum, PartitionSpec, Struct};
60use crate::{Error, ErrorKind};
61
62impl ArrowReader {
63    /// Take a stream of FileScanTasks and reads all the files.
64    /// Returns a [`ScanResult`] containing the record batch stream and scan metrics.
65    pub fn read(self, tasks: FileScanTaskStream) -> Result<ScanResult> {
66        let concurrency_limit_data_files = self.concurrency_limit_data_files;
67        let scan_metrics = ScanMetrics::new();
68
69        let task_reader = FileScanTaskReader {
70            batch_size: self.batch_size,
71            file_io: self.file_io,
72            delete_file_loader: self
73                .delete_file_loader
74                .with_scan_metrics(scan_metrics.clone()),
75            row_group_filtering_enabled: self.row_group_filtering_enabled,
76            row_selection_enabled: self.row_selection_enabled,
77            bloom_filter_enabled: self.bloom_filter_enabled,
78            parquet_read_options: self.parquet_read_options,
79            scan_metrics: scan_metrics.clone(),
80        };
81
82        // Fast-path for single concurrency to avoid overhead of try_flatten_unordered
83        let stream: ArrowRecordBatchStream = if concurrency_limit_data_files == 1 {
84            Box::pin(
85                tasks
86                    .and_then(move |task| task_reader.clone().process(task))
87                    .map_err(|err| {
88                        Error::new(ErrorKind::Unexpected, "file scan task generate failed")
89                            .with_source(err)
90                    })
91                    .try_flatten(),
92            )
93        } else {
94            Box::pin(
95                tasks
96                    .map_ok(move |task| task_reader.clone().process(task))
97                    .map_err(|err| {
98                        Error::new(ErrorKind::Unexpected, "file scan task generate failed")
99                            .with_source(err)
100                    })
101                    .try_buffer_unordered(concurrency_limit_data_files)
102                    .try_flatten_unordered(concurrency_limit_data_files),
103            )
104        };
105
106        Ok(ScanResult::new(stream, scan_metrics))
107    }
108}
109
110// Metadata columns synthesized without reading any data column, so a projection of only
111// these can be pruned to zero data columns. Narrower than `is_metadata_field`, which also
112// matches `_deleted` -- excluded here because it has no synthesis handler.
113const PRUNABLE_METADATA_FIELDS: &[i32] = &[
114    RESERVED_FIELD_ID_FILE,
115    RESERVED_FIELD_ID_SPEC_ID,
116    RESERVED_FIELD_ID_PARTITION,
117    RESERVED_FIELD_ID_POS,
118    RESERVED_FIELD_ID_ROW_ID,
119    RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
120];
121
122/// Per-scan state for processing [`FileScanTask`]s. Created once per
123/// [`ArrowReader::read`] call and cloned per task.
124#[derive(Clone)]
125struct FileScanTaskReader {
126    batch_size: Option<usize>,
127    file_io: FileIO,
128    delete_file_loader: CachingDeleteFileLoader,
129    row_group_filtering_enabled: bool,
130    row_selection_enabled: bool,
131    bloom_filter_enabled: bool,
132    parquet_read_options: ParquetReadOptions,
133    scan_metrics: ScanMetrics,
134}
135
136impl FileScanTaskReader {
137    async fn process(self, task: FileScanTask) -> Result<ArrowRecordBatchStream> {
138        let should_load_page_index = (self.row_selection_enabled && task.predicate().is_some())
139            || !task.deletes().is_empty();
140        let mut parquet_read_options = self.parquet_read_options;
141        parquet_read_options.preload_page_index = should_load_page_index;
142
143        let delete_filter_rx = self
144            .delete_file_loader
145            .load_deletes(task.deletes(), task.schema_ref());
146
147        // Open the Parquet file once, loading its metadata
148        let (parquet_file_reader, arrow_metadata) = ArrowReader::open_parquet_file(
149            task.data_file_path(),
150            &self.file_io,
151            task.file_size_in_bytes(),
152            parquet_read_options,
153            self.scan_metrics.bytes_read_counter(),
154            task.key_metadata(),
155        )
156        .await?;
157
158        // Check if Parquet file has embedded field IDs
159        // Corresponds to Java's ParquetSchemaUtil.hasIds()
160        // Reference: parquet/src/main/java/org/apache/iceberg/parquet/ParquetSchemaUtil.java:118
161        let missing_field_ids = arrow_metadata
162            .schema()
163            .fields()
164            .iter()
165            .next()
166            .is_some_and(|f| f.metadata().get(PARQUET_FIELD_ID_META_KEY).is_none());
167
168        // Position-based fallback applies only when the file has no embedded field IDs
169        // AND no name mapping is available. With a name mapping, field IDs are assigned
170        // to the Arrow schema below, and projection/predicate planning must use them
171        // (see #2403).
172        let use_position_fallback = missing_field_ids && task.name_mapping().is_none();
173
174        let project_pos = task.project_field_ids().contains(&RESERVED_FIELD_ID_POS);
175        let project_row_id = task.project_field_ids().contains(&RESERVED_FIELD_ID_ROW_ID);
176
177        // The RowNumber virtual column materializes `_pos`. It is also the per-row
178        // positional fallback for `_row_id` (`first_row_id + pos`), so add it whenever
179        // `_row_id` is synthesized. A null `first_row_id` nulls the whole `_row_id`
180        // column, so nothing is synthesized and the column is not needed.
181        let need_row_number = project_pos || (project_row_id && task.first_row_id().is_some());
182
183        let field_ids = task.project_field_ids();
184        let metadata_only_projection = !field_ids.is_empty()
185            && field_ids
186                .iter()
187                .all(|id| PRUNABLE_METADATA_FIELDS.contains(id));
188
189        let install_row_number = need_row_number || metadata_only_projection;
190
191        let arrow_metadata = Self::configure_arrow_reader_metadata(
192            arrow_metadata,
193            &task,
194            missing_field_ids,
195            install_row_number,
196        )?;
197
198        // Build the stream reader, reusing the already-opened file reader
199        let mut record_batch_stream_builder =
200            ParquetRecordBatchStreamBuilder::new_with_metadata(parquet_file_reader, arrow_metadata);
201
202        // Whether the file physically carries the `_last_updated_sequence_number` column
203        // (some engines, e.g. Iceberg Java on rewrite, write it per-row), resolved by its
204        // embedded field id against the Parquet schema.
205        let project_last_updated_seq = task
206            .project_field_ids()
207            .contains(&RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER);
208
209        // Parquet leaf index of the physically-stored column, resolved by its embedded
210        // reserved field id. `find_leaf_by_field_id` tolerates id-less leaves (e.g. a
211        // Variant column's internal metadata/value leaves, which the spec requires to have
212        // no id), so an unprojected variant alongside a metadata column with correct ID does
213        // not hide it.
214        let phys_last_updated_seq_leaf = if project_last_updated_seq {
215            find_leaf_by_field_id(
216                record_batch_stream_builder.parquet_schema(),
217                RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
218            )
219        } else {
220            None
221        };
222
223        // Present by name but not by the embedded id (only meaningful when no by-id column
224        // was found). An unthreadable shape we reject rather than coalesce incorrectly.
225        let last_updated_seq_present_by_name_only = project_last_updated_seq
226            && phys_last_updated_seq_leaf.is_none()
227            && record_batch_stream_builder
228                .schema()
229                .fields()
230                .iter()
231                .any(|f| f.name() == RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER);
232
233        // Read the physical column only when first_row_id is set (with a data sequence
234        // number to fall back to). A null first_row_id drops the leaf and nulls the whole
235        // column below, discarding any per-row values the file carries -- matching Java
236        // (`ValueReaders.lastUpdated` nulls when the base row id is null).
237        let coalesce_last_updated_seq_leaf = phys_last_updated_seq_leaf
238            .filter(|_| task.first_row_id().is_some() && task.data_sequence_number().is_some());
239
240        let phys_row_id_leaf = if project_row_id {
241            find_leaf_by_field_id(
242                record_batch_stream_builder.parquet_schema(),
243                RESERVED_FIELD_ID_ROW_ID,
244            )
245        } else {
246            None
247        };
248
249        // A column named `_row_id` that carries no embedded field id, in a file that DOES
250        // use embedded ids -- an unthreadable physically-stored `_row_id`, rejected below.
251        // Gated on `!use_position_fallback`: under positional fallback every column lacks an
252        // embedded id and synthetic ids are assigned by position, so a user column that
253        // happens to be named `_row_id` is real data, not a reserved metadata column.
254        let row_id_present_by_name_only = project_row_id
255            && !use_position_fallback
256            && phys_row_id_leaf.is_none()
257            && record_batch_stream_builder
258                .schema()
259                .fields()
260                .iter()
261                .any(|f| f.name() == RESERVED_COL_NAME_ROW_ID);
262
263        // Read the physical column only when first_row_id is set. A null first_row_id
264        // nulls the whole column below (matching Java `ValueReaders.rowIds`).
265        let coalesce_row_id_leaf = phys_row_id_leaf.filter(|_| task.first_row_id().is_some());
266
267        // Filter out metadata fields for Parquet projection (they don't exist in files)
268        let project_field_ids_without_metadata: Vec<i32> = task
269            .project_field_ids()
270            .iter()
271            .filter(|&&id| !is_metadata_field(id))
272            .copied()
273            .collect();
274
275        // Create projection mask based on field IDs
276        // - If file has embedded IDs: field-ID-based projection
277        // - If name mapping applied: field-ID-based projection using the IDs the name
278        //   mapping assigned to the Arrow schema
279        // - Otherwise: position-based fallback projection
280        let mut projection_mask = ArrowReader::get_arrow_projection_mask(
281            &project_field_ids_without_metadata,
282            task.schema(),
283            record_batch_stream_builder.parquet_schema(),
284            record_batch_stream_builder.schema(),
285            use_position_fallback, // Whether to use position-based (true) or field-ID-based (false) projection
286        )?;
287
288        // A metadata-only projection leaves `project_field_ids_without_metadata` empty,
289        // which `get_arrow_projection_mask` maps to "read all columns" (so `COUNT(*)` still
290        // gets a row count). Downgrade that to "read no data columns": `install_row_number`
291        // put the RowNumber virtual column on every metadata-only projection as a row-count
292        // source independent of the data columns, so the count survives with zero data
293        // columns read. `COUNT(*)` (an empty projection) has no RowNumber and keeps reading
294        // all columns to preserve the row count.
295        //
296        // This runs BEFORE the union so any physical metadata leaf is added onto a `none`
297        // base, pruning the read to just that leaf (`union` with an `all` base stays `all`).
298        if project_field_ids_without_metadata.is_empty() && install_row_number {
299            projection_mask =
300                ProjectionMask::none(record_batch_stream_builder.parquet_schema().num_columns());
301        }
302
303        // Union in the physical leaves of any metadata columns we will coalesce. Their
304        // reserved field ids are not in the task schema, so they can't be requested through
305        // `get_arrow_projection_mask` (which resolves ids against the task schema); add
306        // their Parquet leaves directly.
307        for leaf in [coalesce_last_updated_seq_leaf, coalesce_row_id_leaf]
308            .into_iter()
309            .flatten()
310        {
311            let phys_mask =
312                ProjectionMask::leaves(record_batch_stream_builder.parquet_schema(), vec![leaf]);
313            projection_mask.union(&phys_mask);
314        }
315
316        record_batch_stream_builder =
317            record_batch_stream_builder.with_projection(projection_mask.clone());
318
319        // RecordBatchTransformer performs any transformations required on the RecordBatches
320        // that come back from the file, such as type promotion, default column insertion,
321        // column re-ordering, partition constants, and virtual field addition (like _file)
322        let mut record_batch_transformer_builder =
323            RecordBatchTransformerBuilder::new(task.schema_ref(), task.project_field_ids());
324
325        // Add the _file metadata column if it's in the projected fields
326        if task.project_field_ids().contains(&RESERVED_FIELD_ID_FILE) {
327            let file_datum = Datum::string(task.data_file_path().to_string());
328            record_batch_transformer_builder =
329                record_batch_transformer_builder.with_constant(RESERVED_FIELD_ID_FILE, file_datum);
330        }
331
332        if task
333            .project_field_ids()
334            .contains(&RESERVED_FIELD_ID_SPEC_ID)
335        {
336            let partition_spec = task
337                .partition_spec()
338                .ok_or_else(|| Error::new(ErrorKind::Unexpected, "Partition spec is missing"))?;
339
340            let spec_id_datum = Datum::int(partition_spec.spec_id());
341            record_batch_transformer_builder = record_batch_transformer_builder
342                .with_constant(RESERVED_FIELD_ID_SPEC_ID, spec_id_datum);
343        }
344
345        if project_last_updated_seq {
346            // Materialize the column, gated on the data file's `first_row_id`. Java gates
347            // it this way (`ValueReaders.lastUpdated` returns nulls when the base row id is
348            // null); the spec itself only says the column is assigned the manifest entry's
349            // sequence number on read.
350            record_batch_transformer_builder =
351                match (task.first_row_id(), task.data_sequence_number()) {
352                    (Some(_), Some(seq)) => {
353                        let datum = Datum::long(seq);
354                        if coalesce_last_updated_seq_leaf.is_some() {
355                            // The file physically carries the column: read the per-row value,
356                            // falling back to the data sequence number only where null.
357                            record_batch_transformer_builder.with_coalesced_last_updated_seq_column(
358                                RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
359                                datum,
360                            )
361                        } else if last_updated_seq_present_by_name_only {
362                            // Present by name but without the embedded field id (name mapping /
363                            // positional fallback). The transformer keys the source column by
364                            // field id, so we can't thread it; no real writer produces this, so
365                            // reject loudly rather than silently overwrite with the constant.
366                            // Arm-local by design: only this arm reads the physical column, so
367                            // only here can a name-only column defeat us. The `(None, _)` arm
368                            // nulls the column without reading it, so it needs no such guard.
369                            return Err(Error::new(
370                                ErrorKind::FeatureUnsupported,
371                                "Reading a physically-stored _last_updated_sequence_number column \
372                             without an embedded field id is not supported",
373                            ));
374                        } else {
375                            // Column absent: derive it from the data sequence number.
376                            record_batch_transformer_builder.with_constant(
377                                RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
378                                datum,
379                            )
380                        }
381                    }
382                    // Null first_row_id (v1/v2, or a pre-upgrade v3 snapshot): the column is null.
383                    (None, _) => record_batch_transformer_builder.with_null_metadata_column(
384                        RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
385                    )?,
386                    // first_row_id present but no data sequence number: after manifest
387                    // inheritance a committed entry always has one, so this is a malformed
388                    // manifest rather than a legitimate null.
389                    (Some(_), None) => {
390                        return Err(invalid_data!(
391                            "Data file {} has a first_row_id but no data sequence number",
392                            task.data_file_path()
393                        ));
394                    }
395                };
396        }
397
398        if project_row_id {
399            // A name-only physical `_row_id` can't be threaded (synthesis keys the leaf by
400            // its reserved field id). Reject only when `first_row_id` is set; otherwise the
401            // leaf is never read and `_row_id` is nulled out downstream (matching Java
402            // `ValueReaders.rowIds`), so a pre-v3 file with a user column named `_row_id`
403            // reads back as null rather than erroring.
404            if task.first_row_id().is_some() && row_id_present_by_name_only {
405                return Err(Error::new(
406                    ErrorKind::FeatureUnsupported,
407                    "Reading a physically-stored _row_id column without an embedded field id \
408                     is not supported",
409                ));
410            }
411
412            // `_row_id` is synthesized downstream over the record-batch stream (see
413            // `row_lineage::synthesize_row_id_column`); the transformer only passes the
414            // resulting column through, like `_pos`.
415            record_batch_transformer_builder =
416                record_batch_transformer_builder.with_virtual_field(RESERVED_FIELD_ID_ROW_ID);
417        }
418
419        if let (Some(partition_spec), Some(partition_data)) =
420            (task.partition_spec(), task.partition())
421        {
422            record_batch_transformer_builder = record_batch_transformer_builder
423                .with_partition(Arc::clone(partition_spec), partition_data.clone())?;
424        }
425
426        if project_pos {
427            record_batch_transformer_builder =
428                record_batch_transformer_builder.with_virtual_field(RESERVED_FIELD_ID_POS);
429        }
430
431        // Add the _partition metadata struct column if it's in the projected fields.
432        // Computed lazily here at read time from the unified partition type + task's spec + data.
433        if task
434            .project_field_ids()
435            .contains(&RESERVED_FIELD_ID_PARTITION)
436            && let Some(unified_type) = task.unified_partition_type()
437        {
438            let (spec, partition_data) = match (task.partition_spec(), task.partition()) {
439                (Some(spec), Some(data)) => (Arc::clone(spec), data.clone()),
440                // A missing spec/data is only acceptable when there are no partition
441                // fields to fill (unpartitioned table). If the unified type has fields
442                // but we lack a spec or data, the task is inconsistent and we cannot
443                // build the _partition column.
444                _ if unified_type.fields().is_empty() => {
445                    (Arc::new(PartitionSpec::unpartition_spec()), Struct::empty())
446                }
447                _ => {
448                    return Err(Error::new(
449                        ErrorKind::Unexpected,
450                        "cannot build _partition column: unified partition type has fields \
451                         but the scan task is missing its partition spec or data",
452                    ));
453                }
454            };
455            let constant = build_partition_constant(unified_type, &spec, &partition_data)?;
456            record_batch_transformer_builder =
457                record_batch_transformer_builder.with_partition_constant(constant);
458        }
459
460        let mut record_batch_transformer = record_batch_transformer_builder.build();
461
462        if let Some(batch_size) = self.batch_size {
463            record_batch_stream_builder = record_batch_stream_builder.with_batch_size(batch_size);
464        }
465
466        let delete_filter = delete_filter_rx.await.unwrap()?;
467        let delete_predicate = delete_filter.build_equality_delete_predicate(&task).await?;
468
469        // In addition to the optional predicate supplied in the `FileScanTask`,
470        // we also have an optional predicate resulting from equality delete files.
471        // If both are present, we logical-AND them together to form a single filter
472        // predicate that we can pass to the `RecordBatchStreamBuilder`.
473        let final_predicate = match (task.predicate(), delete_predicate) {
474            (None, None) => None,
475            (Some(predicate), None) => Some(predicate.clone()),
476            (None, Some(ref predicate)) => Some(predicate.clone()),
477            (Some(filter_predicate), Some(delete_predicate)) => {
478                Some(filter_predicate.clone().and(delete_predicate))
479            }
480        };
481
482        // There are three possible sources for potential lists of selected RowGroup indices,
483        // and two for `RowSelection`s.
484        // Selected RowGroup index lists can come from three sources:
485        //   * When task.start and task.length specify a byte range (file splitting);
486        //   * When there are equality delete files that are applicable;
487        //   * When there is a scan predicate and row_group_filtering_enabled = true.
488        // `RowSelection`s can be created in either or both of the following cases:
489        //   * When there are positional delete files that are applicable;
490        //   * When there is a scan predicate and row_selection_enabled = true
491        // Note that row group filtering from predicates only happens when
492        // there is a scan predicate AND row_group_filtering_enabled = true,
493        // but we perform row selection filtering if there are applicable
494        // equality delete files OR (there is a scan predicate AND row_selection_enabled),
495        // since the only implemented method of applying positional deletes is
496        // by using a `RowSelection`.
497        let mut selected_row_group_indices = None;
498        let mut row_selection = None;
499
500        // Filter row groups based on byte range from task.start and task.length.
501        // If both start and length are 0, read the entire file (backwards compatibility).
502        if task.start() != 0 || task.length() != 0 {
503            let byte_range_filtered_row_groups = ArrowReader::filter_row_groups_by_byte_range(
504                record_batch_stream_builder.metadata(),
505                task.start(),
506                task.length(),
507            )?;
508            selected_row_group_indices = Some(byte_range_filtered_row_groups);
509        }
510
511        if let Some(predicate) = final_predicate {
512            let (iceberg_field_ids, field_id_map) = ArrowReader::build_field_id_set_and_map(
513                record_batch_stream_builder.parquet_schema(),
514                record_batch_stream_builder.schema(),
515                &predicate,
516                use_position_fallback,
517            )?;
518
519            let row_filter = ArrowReader::get_row_filter(
520                &predicate,
521                record_batch_stream_builder.parquet_schema(),
522                &iceberg_field_ids,
523                &field_id_map,
524            )?;
525            record_batch_stream_builder = record_batch_stream_builder.with_row_filter(row_filter);
526
527            if self.row_group_filtering_enabled {
528                let predicate_filtered_row_groups = ArrowReader::get_selected_row_group_indices(
529                    &predicate,
530                    record_batch_stream_builder.metadata(),
531                    &field_id_map,
532                    task.schema(),
533                )?;
534
535                // Merge predicate-based filtering with byte range filtering (if present)
536                // by taking the intersection of both filters
537                selected_row_group_indices = match selected_row_group_indices {
538                    Some(byte_range_filtered) => {
539                        // Keep only row groups that are in both filters
540                        let intersection: Vec<usize> = byte_range_filtered
541                            .into_iter()
542                            .filter(|idx| predicate_filtered_row_groups.contains(idx))
543                            .collect();
544                        Some(intersection)
545                    }
546                    None => Some(predicate_filtered_row_groups),
547                };
548            }
549
550            if self.bloom_filter_enabled {
551                let all_rgs;
552                let candidate_rgs = match &selected_row_group_indices {
553                    Some(indices) => indices.as_slice(),
554                    None => {
555                        all_rgs = (0..record_batch_stream_builder.metadata().num_row_groups())
556                            .collect::<Vec<_>>();
557                        &all_rgs
558                    }
559                };
560
561                let bloom_filtered = Self::filter_row_groups_by_bloom_filter(
562                    &predicate,
563                    &mut record_batch_stream_builder,
564                    candidate_rgs,
565                    &field_id_map,
566                )
567                .await?;
568
569                if bloom_filtered.len() < candidate_rgs.len() {
570                    selected_row_group_indices = Some(bloom_filtered);
571                }
572            }
573
574            if self.row_selection_enabled {
575                row_selection = ArrowReader::get_row_selection_for_filter_predicate(
576                    &predicate,
577                    record_batch_stream_builder.metadata(),
578                    &selected_row_group_indices,
579                    &field_id_map,
580                    task.schema(),
581                )?;
582            }
583        }
584
585        let positional_delete_indexes = delete_filter.get_delete_vector(&task);
586
587        if let Some(positional_delete_indexes) = positional_delete_indexes {
588            let delete_row_selection = {
589                let positional_delete_indexes = positional_delete_indexes.lock().unwrap();
590
591                ArrowReader::build_deletes_row_selection(
592                    record_batch_stream_builder.metadata().row_groups(),
593                    &selected_row_group_indices,
594                    &positional_delete_indexes,
595                )
596            }?;
597
598            // merge the row selection from the delete files with the row selection
599            // from the filter predicate, if there is one from the filter predicate
600            row_selection = match row_selection {
601                None => Some(delete_row_selection),
602                Some(filter_row_selection) => {
603                    Some(filter_row_selection.intersection(&delete_row_selection))
604                }
605            };
606        }
607
608        if let Some(row_selection) = row_selection {
609            record_batch_stream_builder =
610                record_batch_stream_builder.with_row_selection(row_selection);
611        }
612
613        if let Some(selected_row_group_indices) = selected_row_group_indices {
614            record_batch_stream_builder =
615                record_batch_stream_builder.with_row_groups(selected_row_group_indices);
616        }
617
618        // Build the batch stream and send all the RecordBatches that it generates
619        // to the requester. When `_row_id` is projected, synthesize it over the raw parquet
620        // batches (using the reader-produced `_pos` position) before the transformer, which
621        // then passes it through as a virtual field.
622        let first_row_id = task.first_row_id();
623        let record_batch_stream = record_batch_stream_builder.build()?.map(move |batch| {
624            let mut batch = batch.map_err(|err| -> Error { err.into() })?;
625            if project_row_id {
626                batch = synthesize_row_id_column(batch, first_row_id)?;
627            }
628            // Process the record batch (type promotion, column reordering, virtual fields, etc.)
629            record_batch_transformer.process_record_batch(batch)
630        });
631
632        Ok(Box::pin(record_batch_stream) as ArrowRecordBatchStream)
633    }
634
635    /// Applies all task-specific schema and virtual-column options, rebuilding the
636    /// Arrow reader metadata at most once.
637    fn configure_arrow_reader_metadata(
638        arrow_metadata: ArrowReaderMetadata,
639        task: &FileScanTask,
640        missing_field_ids: bool,
641        install_row_number: bool,
642    ) -> Result<ArrowReaderMetadata> {
643        // Schema resolution follows the Iceberg Column Projection rule:
644        // "Columns in Iceberg data files are selected by field id."
645        // https://iceberg.apache.org/spec/#column-projection
646        //
647        // This mirrors Java's ReadConf strategy: use embedded IDs with pruneColumns();
648        // otherwise, use applyNameMapping() followed by pruneColumns() when configured, or
649        // addFallbackIds() followed by pruneColumnsFallback() for position-based fallback.
650        // The fast path returns early without materializing an owned schema.
651        let arrow_schema = if missing_field_ids {
652            let schema = if let Some(name_mapping) = task.name_mapping() {
653                apply_name_mapping_to_arrow_schema(
654                    Arc::clone(arrow_metadata.schema()),
655                    name_mapping,
656                )?
657            } else {
658                add_fallback_field_ids_to_arrow_schema(arrow_metadata.schema())
659            };
660            // Coerce INT96 timestamp columns before building the stream reader to avoid
661            // i64 overflow in arrow-rs. Apply this after assigning any missing field IDs
662            // so the final schema contains both changes.
663            coerce_int96_timestamps(&schema, task.schema()).unwrap_or(schema)
664        } else if let Some(coerced) =
665            coerce_int96_timestamps(arrow_metadata.schema(), task.schema())
666        {
667            coerced
668        } else if install_row_number {
669            Arc::clone(arrow_metadata.schema())
670        } else {
671            return Ok(arrow_metadata);
672        };
673
674        let mut options = ArrowReaderOptions::new().with_schema(Arc::clone(&arrow_schema));
675        if install_row_number {
676            let row_number_field = Arc::new(
677                Field::new(RESERVED_COL_NAME_POS, DataType::Int64, false)
678                    .with_metadata(HashMap::from([(
679                        PARQUET_FIELD_ID_META_KEY.to_string(),
680                        RESERVED_FIELD_ID_POS.to_string(),
681                    )]))
682                    .with_extension_type(RowNumber),
683            );
684            options = options.with_virtual_columns(vec![row_number_field])?;
685        }
686
687        ArrowReaderMetadata::try_new(Arc::clone(arrow_metadata.metadata()), options).map_err(|e| {
688            Error::new(
689                ErrorKind::Unexpected,
690                format!(
691                    "Failed to create ArrowReaderMetadata with the configured reader options \
692                     (missing_field_ids: {missing_field_ids}, \
693                      install_row_number: {install_row_number}, schema: {arrow_schema})"
694                ),
695            )
696            .with_source(e)
697        })
698    }
699
700    /// Reads bloom filters for relevant columns and evaluates the predicate
701    /// against them to filter out row groups that definitely don't match.
702    async fn filter_row_groups_by_bloom_filter(
703        predicate: &BoundPredicate,
704        builder: &mut ParquetRecordBatchStreamBuilder<ArrowFileReader>,
705        candidate_row_groups: &[usize],
706        field_id_map: &HashMap<i32, usize>,
707    ) -> Result<Vec<usize>> {
708        // Only collect field IDs from eq/in predicates — the only types
709        // bloom filters can help with. Skip columns not in the parquet schema.
710        let bloom_filter_field_ids: Vec<i32> = collect_bloom_filter_field_ids(predicate)?
711            .into_iter()
712            .filter(|id| field_id_map.contains_key(id))
713            .collect();
714
715        if bloom_filter_field_ids.is_empty() {
716            return Ok(candidate_row_groups.to_vec());
717        }
718
719        let mut result = Vec::with_capacity(candidate_row_groups.len());
720
721        for &rg_idx in candidate_row_groups {
722            let mut bloom_filters: HashMap<i32, ColumnBloomFilter> = HashMap::new();
723
724            for &field_id in &bloom_filter_field_ids {
725                let col_idx = field_id_map[&field_id];
726                let col_meta = builder.metadata().row_group(rg_idx).column(col_idx);
727
728                // Only attempt to load if this column chunk actually has a bloom filter
729                if col_meta.bloom_filter_offset().is_none() {
730                    continue;
731                }
732
733                let physical_type = col_meta.column_type();
734                let type_length = col_meta.column_descr().type_length();
735
736                match builder
737                    .get_row_group_column_bloom_filter(rg_idx, col_idx)
738                    .await
739                {
740                    Ok(Some(sbbf)) => {
741                        bloom_filters.insert(
742                            field_id,
743                            ColumnBloomFilter::new(sbbf, physical_type, type_length),
744                        );
745                    }
746                    Ok(None) => {}
747                    Err(e) => {
748                        // Left absent from the map, so the evaluator treats the column
749                        // as might-match and the row group survives.
750                        tracing::debug!(
751                            "Bloom filter for field {field_id} in row group {rg_idx} could not be read: {e}"
752                        );
753                    }
754                }
755            }
756
757            match BloomFilterEvaluator::eval(predicate, &bloom_filters) {
758                Ok(true) => result.push(rg_idx),
759                Ok(false) => { /* Row group pruned by bloom filter */ }
760                Err(e) => {
761                    tracing::debug!(
762                        "Bloom filter evaluation failed for row group {rg_idx}, including it: {e}"
763                    );
764                    result.push(rg_idx);
765                }
766            }
767        }
768
769        Ok(result)
770    }
771}
772
773impl ArrowReader {
774    /// Opens a Parquet file and loads its metadata, wrapping the reader with
775    /// [`CountingFileRead`] so all I/O is accumulated into `bytes_read`.
776    pub(crate) async fn open_parquet_file(
777        data_file_path: &str,
778        file_io: &FileIO,
779        file_size_in_bytes: u64,
780        parquet_read_options: ParquetReadOptions,
781        bytes_read: &Arc<AtomicU64>,
782        key_metadata: Option<&[u8]>,
783    ) -> Result<(ArrowFileReader, ArrowReaderMetadata)> {
784        let parquet_file = file_io.new_input(data_file_path)?;
785        let counting_reader =
786            CountingFileRead::new(parquet_file.reader().await?, Arc::clone(bytes_read));
787        Self::build_parquet_reader(
788            Box::new(counting_reader),
789            file_size_in_bytes,
790            parquet_read_options,
791            key_metadata,
792        )
793        .await
794    }
795
796    async fn build_parquet_reader(
797        parquet_reader: Box<dyn FileRead>,
798        file_size_in_bytes: u64,
799        parquet_read_options: ParquetReadOptions,
800        key_metadata: Option<&[u8]>,
801    ) -> Result<(ArrowFileReader, ArrowReaderMetadata)> {
802        let mut reader = ArrowFileReader::new(
803            FileMetadata {
804                size: file_size_in_bytes,
805            },
806            parquet_reader,
807        )
808        .with_parquet_read_options(parquet_read_options);
809
810        let arrow_reader_options = Self::build_arrow_reader_options(key_metadata)?;
811
812        let arrow_metadata = ArrowReaderMetadata::load_async(&mut reader, arrow_reader_options)
813            .await
814            .map_err(|e| {
815                Error::new(ErrorKind::Unexpected, "Failed to load Parquet metadata").with_source(e)
816            })?;
817
818        Ok((reader, arrow_metadata))
819    }
820
821    /// Builds `ArrowReaderOptions`, adding `FileDecryptionProperties` when
822    /// key metadata is present for Parquet Modular Encryption.
823    fn build_arrow_reader_options(key_metadata: Option<&[u8]>) -> Result<ArrowReaderOptions> {
824        match key_metadata {
825            Some(km) => {
826                let standard_key_metadata = StandardKeyMetadata::decode(km)?;
827                let mut builder = FileDecryptionProperties::builder(
828                    standard_key_metadata.encryption_key().as_bytes().to_vec(),
829                );
830                if let Some(aad) = standard_key_metadata.aad_prefix() {
831                    builder = builder.with_aad_prefix(aad.to_vec());
832                }
833                let decryption_properties = builder.build().map_err(|e| {
834                    Error::new(
835                        ErrorKind::Unexpected,
836                        "Failed to build Parquet file decryption properties",
837                    )
838                    .with_source(e)
839                })?;
840                Ok(
841                    ArrowReaderOptions::new()
842                        .with_file_decryption_properties(decryption_properties),
843                )
844            }
845            None => Ok(ArrowReaderOptions::default()),
846        }
847    }
848}
849
850#[cfg(test)]
851mod tests {
852    use std::collections::HashMap;
853    use std::fs::File;
854    use std::sync::Arc;
855
856    use arrow_array::cast::AsArray;
857    use arrow_array::{Array, ArrayRef, Int32Array, Int64Array, RecordBatch, StringArray};
858    use arrow_cast::cast;
859    use arrow_schema::{DataType, Field, Schema as ArrowSchema};
860    use futures::TryStreamExt;
861    use parquet::arrow::{ArrowWriter, PARQUET_FIELD_ID_META_KEY};
862    use parquet::basic::Compression;
863    use parquet::file::properties::WriterProperties;
864    use tempfile::TempDir;
865
866    use crate::Runtime;
867    use crate::arrow::ArrowReaderBuilder;
868    use crate::arrow::test_utils::write_encrypted_parquet;
869    use crate::io::FileIO;
870    use crate::metadata_columns::{
871        RESERVED_COL_NAME_POS, RESERVED_COL_NAME_ROW_ID, RESERVED_FIELD_ID_FILE,
872        RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID,
873    };
874    use crate::scan::{FileScanTask, FileScanTaskDeleteFile, FileScanTaskStream};
875    use crate::spec::{DataFileFormat, NestedField, PrimitiveType, Schema, SchemaRef, Type};
876
877    // INT96 encoding: [nanos_low_u32, nanos_high_u32, julian_day_u32]
878    // Julian day 2_440_588 = Unix epoch (1970-01-01)
879    const UNIX_EPOCH_JULIAN: i64 = 2_440_588;
880    const MICROS_PER_DAY: i64 = 86_400_000_000;
881    // Noon on 3333-01-01 (Julian day 2_953_529) — outside the i64 nanosecond range (~1677-2262).
882    const INT96_TEST_NANOS_WITHIN_DAY: u64 = 43_200_000_000_000;
883    const INT96_TEST_JULIAN_DAY: u32 = 2_953_529;
884
885    fn make_int96_test_value() -> (parquet::data_type::Int96, i64) {
886        let mut val = parquet::data_type::Int96::new();
887        val.set_data(
888            (INT96_TEST_NANOS_WITHIN_DAY & 0xFFFFFFFF) as u32,
889            (INT96_TEST_NANOS_WITHIN_DAY >> 32) as u32,
890            INT96_TEST_JULIAN_DAY,
891        );
892        let expected_micros = (INT96_TEST_JULIAN_DAY as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY
893            + (INT96_TEST_NANOS_WITHIN_DAY / 1_000) as i64;
894        (val, expected_micros)
895    }
896
897    async fn read_int96_batches(
898        file_path: &str,
899        schema: SchemaRef,
900        project_field_ids: Vec<i32>,
901    ) -> Vec<RecordBatch> {
902        let file_io = FileIO::new_with_fs();
903        let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
904
905        let file_size = std::fs::metadata(file_path).unwrap().len();
906        let task = FileScanTask::builder()
907            .with_file_size_in_bytes(file_size)
908            .with_start(0)
909            .with_length(file_size)
910            .with_data_file_path(file_path.to_string())
911            .with_data_file_format(DataFileFormat::Parquet)
912            .with_schema(schema)
913            .with_project_field_ids(project_field_ids)
914            .with_case_sensitive(false)
915            .build()
916            .unwrap();
917
918        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
919        reader
920            .read(tasks)
921            .unwrap()
922            .stream()
923            .try_collect()
924            .await
925            .unwrap()
926    }
927
928    // ArrowWriter cannot write INT96, so we use SerializedFileWriter directly.
929    fn write_int96_parquet_file(
930        table_location: &str,
931        filename: &str,
932        with_field_ids: bool,
933    ) -> (String, Vec<i64>) {
934        use parquet::basic::{Repetition, Type as PhysicalType};
935        use parquet::data_type::{Int32Type, Int96, Int96Type};
936        use parquet::file::writer::SerializedFileWriter;
937        use parquet::schema::types::Type as SchemaType;
938
939        let file_path = format!("{table_location}/{filename}");
940
941        let mut ts_builder = SchemaType::primitive_type_builder("ts", PhysicalType::INT96)
942            .with_repetition(Repetition::OPTIONAL);
943        let mut id_builder = SchemaType::primitive_type_builder("id", PhysicalType::INT32)
944            .with_repetition(Repetition::REQUIRED);
945
946        if with_field_ids {
947            ts_builder = ts_builder.with_id(Some(1));
948            id_builder = id_builder.with_id(Some(2));
949        }
950
951        let schema = SchemaType::group_type_builder("schema")
952            .with_fields(vec![
953                Arc::new(ts_builder.build().unwrap()),
954                Arc::new(id_builder.build().unwrap()),
955            ])
956            .build()
957            .unwrap();
958
959        // Dates outside the i64 nanosecond range (~1677-2262) overflow without coercion.
960        const NOON_NANOS: u64 = INT96_TEST_NANOS_WITHIN_DAY;
961        const JULIAN_3333: u32 = INT96_TEST_JULIAN_DAY;
962        const JULIAN_2100: u32 = 2_488_070;
963
964        let test_data: Vec<(u32, u32, u32, i64)> = vec![
965            // 3333-01-01 00:00:00
966            (
967                0,
968                0,
969                JULIAN_3333,
970                (JULIAN_3333 as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY,
971            ),
972            // 3333-01-01 12:00:00
973            (
974                (NOON_NANOS & 0xFFFFFFFF) as u32,
975                (NOON_NANOS >> 32) as u32,
976                JULIAN_3333,
977                (JULIAN_3333 as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY
978                    + (NOON_NANOS / 1_000) as i64,
979            ),
980            // 2100-01-01 00:00:00
981            (
982                0,
983                0,
984                JULIAN_2100,
985                (JULIAN_2100 as i64 - UNIX_EPOCH_JULIAN) * MICROS_PER_DAY,
986            ),
987        ];
988
989        let int96_values: Vec<Int96> = test_data
990            .iter()
991            .map(|(lo, hi, day, _)| {
992                let mut v = Int96::new();
993                v.set_data(*lo, *hi, *day);
994                v
995            })
996            .collect();
997
998        let id_values: Vec<i32> = (0..test_data.len() as i32).collect();
999        let expected_micros: Vec<i64> = test_data.iter().map(|(_, _, _, m)| *m).collect();
1000
1001        let file = File::create(&file_path).unwrap();
1002        let mut writer =
1003            SerializedFileWriter::new(file, Arc::new(schema), Default::default()).unwrap();
1004
1005        let mut row_group = writer.next_row_group().unwrap();
1006        {
1007            // def=1: ts is OPTIONAL and present. No repetition levels (top-level columns).
1008            let mut col = row_group.next_column().unwrap().unwrap();
1009            col.typed::<Int96Type>()
1010                .write_batch(&int96_values, Some(&vec![1; test_data.len()]), None)
1011                .unwrap();
1012            col.close().unwrap();
1013        }
1014        {
1015            let mut col = row_group.next_column().unwrap().unwrap();
1016            col.typed::<Int32Type>()
1017                .write_batch(&id_values, None, None)
1018                .unwrap();
1019            col.close().unwrap();
1020        }
1021        row_group.close().unwrap();
1022        writer.close().unwrap();
1023
1024        (file_path, expected_micros)
1025    }
1026
1027    async fn assert_int96_read_matches(
1028        file_path: &str,
1029        schema: SchemaRef,
1030        project_field_ids: Vec<i32>,
1031        expected_micros: &[i64],
1032    ) {
1033        use arrow_array::TimestampMicrosecondArray;
1034
1035        let batches = read_int96_batches(file_path, schema, project_field_ids).await;
1036
1037        assert_eq!(batches.len(), 1);
1038        let ts_array = batches[0]
1039            .column(0)
1040            .as_any()
1041            .downcast_ref::<TimestampMicrosecondArray>()
1042            .expect("Expected TimestampMicrosecondArray");
1043
1044        for (i, expected) in expected_micros.iter().enumerate() {
1045            assert_eq!(
1046                ts_array.value(i),
1047                *expected,
1048                "Row {i}: got {}, expected {expected}",
1049                ts_array.value(i)
1050            );
1051        }
1052    }
1053
1054    /// Writes a single-column Parquet file encrypted with `encryption_key`, then reads it
1055    /// back through `ArrowReader` and asserts the round-tripped values. The key length
1056    /// selects the AES-GCM variant in arrow-rs (16 -> AES-128, 32 -> AES-256).
1057    async fn assert_encrypted_parquet_roundtrip(encryption_key: &[u8]) {
1058        let aad_prefix = b"aad_prefix";
1059
1060        let schema = Arc::new(
1061            Schema::builder()
1062                .with_schema_id(1)
1063                .with_fields(vec![
1064                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1065                ])
1066                .build()
1067                .unwrap(),
1068        );
1069
1070        let arrow_schema = Arc::new(ArrowSchema::new(vec![
1071            Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1072                PARQUET_FIELD_ID_META_KEY.to_string(),
1073                "1".to_string(),
1074            )])),
1075        ]));
1076
1077        let tmp_dir = TempDir::new().unwrap();
1078        let table_location = tmp_dir.path().to_str().unwrap().to_string();
1079        let file_io = FileIO::new_with_fs();
1080
1081        let id_data = Arc::new(Int32Array::from(vec![10, 20, 30])) as ArrayRef;
1082        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_data]).unwrap();
1083
1084        let file_path = format!("{table_location}/encrypted.parquet");
1085        write_encrypted_parquet(&file_path, &batch, encryption_key, Some(aad_prefix));
1086
1087        let key_metadata = crate::encryption::StandardKeyMetadata::try_new(encryption_key)
1088            .unwrap()
1089            .with_aad_prefix(aad_prefix)
1090            .encode()
1091            .unwrap();
1092
1093        let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
1094
1095        let task = FileScanTask::builder()
1096            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1097            .with_start(0)
1098            .with_length(0)
1099            .with_data_file_path(file_path)
1100            .with_data_file_format(DataFileFormat::Parquet)
1101            .with_schema(schema)
1102            .with_project_field_ids(vec![1])
1103            .with_case_sensitive(false)
1104            .with_key_metadata(Some(key_metadata))
1105            .build()
1106            .unwrap();
1107
1108        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1109        let batches: Vec<RecordBatch> = reader
1110            .read(tasks)
1111            .unwrap()
1112            .stream()
1113            .try_collect()
1114            .await
1115            .unwrap();
1116
1117        assert_eq!(batches.len(), 1);
1118        let ids = batches[0]
1119            .column(0)
1120            .as_any()
1121            .downcast_ref::<Int32Array>()
1122            .unwrap();
1123        assert_eq!(ids.values(), &[10, 20, 30]);
1124    }
1125
1126    #[tokio::test]
1127    async fn test_read_encrypted_parquet_aes_128() {
1128        assert_encrypted_parquet_roundtrip(b"0123456789abcdef").await;
1129    }
1130
1131    #[tokio::test]
1132    async fn test_read_encrypted_parquet_aes_256() {
1133        assert_encrypted_parquet_roundtrip(b"0123456789abcdef0123456789abcdef").await;
1134    }
1135
1136    #[tokio::test]
1137    async fn test_read_encrypted_parquet_without_key_metadata_fails() {
1138        let encryption_key = b"0123456789abcdef";
1139
1140        let schema = Arc::new(
1141            Schema::builder()
1142                .with_schema_id(1)
1143                .with_fields(vec![
1144                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1145                ])
1146                .build()
1147                .unwrap(),
1148        );
1149
1150        let arrow_schema = Arc::new(ArrowSchema::new(vec![
1151            Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1152                PARQUET_FIELD_ID_META_KEY.to_string(),
1153                "1".to_string(),
1154            )])),
1155        ]));
1156
1157        let tmp_dir = TempDir::new().unwrap();
1158        let table_location = tmp_dir.path().to_str().unwrap().to_string();
1159        let file_io = FileIO::new_with_fs();
1160
1161        let id_data = Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef;
1162        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_data]).unwrap();
1163
1164        let file_path = format!("{table_location}/encrypted_no_key.parquet");
1165        write_encrypted_parquet(&file_path, &batch, encryption_key, None);
1166
1167        let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
1168
1169        let task = FileScanTask::builder()
1170            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1171            .with_start(0)
1172            .with_length(0)
1173            .with_data_file_path(file_path)
1174            .with_data_file_format(DataFileFormat::Parquet)
1175            .with_schema(schema)
1176            .with_project_field_ids(vec![1])
1177            .with_case_sensitive(false)
1178            .build()
1179            .unwrap();
1180
1181        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1182        let result: Result<Vec<RecordBatch>, _> =
1183            reader.read(tasks).unwrap().stream().try_collect().await;
1184
1185        let err = result.unwrap_err();
1186        assert_eq!(err.kind(), crate::ErrorKind::Unexpected);
1187        let err_str = format!("{err}");
1188        assert!(
1189            err_str.contains("encrypted footer"),
1190            "Expected error about encrypted footer, got: {err_str}"
1191        );
1192        assert!(
1193            err_str.contains("decryption properties were not provided"),
1194            "Expected error about missing decryption properties, got: {err_str}"
1195        );
1196    }
1197
1198    /// Writes a plain (unencrypted) single-column Int32 "id" parquet file with the
1199    /// given extra Arrow fields/columns appended, returning the file path.
1200    fn write_plain_parquet(
1201        dir: &str,
1202        name: &str,
1203        extra_fields: Vec<Field>,
1204        extra_columns: Vec<ArrayRef>,
1205    ) -> String {
1206        let mut fields =
1207            vec![
1208                Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
1209                    PARQUET_FIELD_ID_META_KEY.to_string(),
1210                    "1".to_string(),
1211                )])),
1212            ];
1213        fields.extend(extra_fields);
1214        let arrow_schema = Arc::new(ArrowSchema::new(fields));
1215
1216        let mut columns: Vec<ArrayRef> = vec![Arc::new(Int32Array::from(vec![1, 2, 3]))];
1217        columns.extend(extra_columns);
1218        let batch = RecordBatch::try_new(arrow_schema.clone(), columns).unwrap();
1219
1220        let file_path = format!("{dir}/{name}");
1221        let file = File::create(&file_path).unwrap();
1222        let props = WriterProperties::builder()
1223            .set_compression(Compression::SNAPPY)
1224            .build();
1225        let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
1226        writer.write(&batch).unwrap();
1227        writer.close().unwrap();
1228        file_path
1229    }
1230
1231    fn last_updated_seq_task(
1232        file_path: String,
1233        first_row_id: Option<i64>,
1234        data_sequence_number: Option<i64>,
1235    ) -> FileScanTask {
1236        use crate::metadata_columns::RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER;
1237
1238        let schema = Arc::new(
1239            Schema::builder()
1240                .with_schema_id(1)
1241                .with_fields(vec![
1242                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1243                ])
1244                .build()
1245                .unwrap(),
1246        );
1247
1248        FileScanTask::builder()
1249            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1250            .with_start(0)
1251            .with_length(0)
1252            .with_data_file_path(file_path)
1253            .with_data_file_format(DataFileFormat::Parquet)
1254            .with_schema(schema)
1255            .with_project_field_ids(vec![1, RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER])
1256            .with_first_row_id(first_row_id)
1257            .with_data_sequence_number(data_sequence_number)
1258            .with_case_sensitive(false)
1259            .build()
1260            .unwrap()
1261    }
1262
1263    /// Asserts the logical per-row values of the `_last_updated_sequence_number`
1264    /// column across all batches, independent of the physical (run-end) encoding.
1265    fn assert_last_updated_seq_column(batches: &[RecordBatch], expected: &[Option<i64>]) {
1266        use arrow_array::cast::AsArray;
1267        use arrow_cast::cast;
1268        use arrow_schema::DataType;
1269
1270        use crate::metadata_columns::RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER;
1271
1272        let mut actual = Vec::new();
1273        for batch in batches {
1274            let col = batch
1275                .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
1276                .expect("_last_updated_sequence_number column should be present");
1277            let logical = cast(col, &DataType::Int64).unwrap();
1278            let values = logical.as_primitive::<arrow_array::types::Int64Type>();
1279            for i in 0..values.len() {
1280                actual.push((!values.is_null(i)).then(|| values.value(i)));
1281            }
1282        }
1283        assert_eq!(actual, expected);
1284    }
1285
1286    #[tokio::test]
1287    async fn test_last_updated_sequence_number_null_when_no_first_row_id() {
1288        let tmp_dir = TempDir::new().unwrap();
1289        let dir = tmp_dir.path().to_str().unwrap();
1290        let file_path = write_plain_parquet(dir, "no_first_row_id.parquet", vec![], vec![]);
1291
1292        // A file with a null first_row_id (v1/v2, or a pre-upgrade v3 snapshot) produces
1293        // a null _last_updated_sequence_number column, even though it has a data
1294        // sequence number; the spec gates both lineage columns on first_row_id.
1295        let task = last_updated_seq_task(file_path, None, Some(9));
1296
1297        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1298        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1299        let batches: Vec<RecordBatch> = reader
1300            .read(tasks)
1301            .unwrap()
1302            .stream()
1303            .try_collect()
1304            .await
1305            .unwrap();
1306
1307        assert_last_updated_seq_column(&batches, &[None, None, None]);
1308    }
1309
1310    #[tokio::test]
1311    async fn test_last_updated_sequence_number_error_when_no_data_seq() {
1312        let tmp_dir = TempDir::new().unwrap();
1313        let dir = tmp_dir.path().to_str().unwrap();
1314        let file_path = write_plain_parquet(dir, "no_data_seq.parquet", vec![], vec![]);
1315
1316        // first_row_id present but data_sequence_number absent: after manifest
1317        // inheritance a committed entry always has one, so this is a malformed
1318        // manifest and must error rather than fabricate or null the column.
1319        let task = last_updated_seq_task(file_path, Some(42), None);
1320
1321        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1322        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1323        let result: Result<Vec<RecordBatch>, _> =
1324            reader.read(tasks).unwrap().stream().try_collect().await;
1325
1326        let err = result.unwrap_err();
1327        assert_eq!(err.kind(), crate::ErrorKind::DataInvalid);
1328        assert!(
1329            format!("{err}").contains("no data sequence number"),
1330            "unexpected error: {err}"
1331        );
1332    }
1333
1334    #[tokio::test]
1335    async fn test_last_updated_sequence_number_derived_from_data_seq() {
1336        let tmp_dir = TempDir::new().unwrap();
1337        let dir = tmp_dir.path().to_str().unwrap();
1338        let file_path = write_plain_parquet(dir, "with_first_row_id.parquet", vec![], vec![]);
1339
1340        // Non-null first_row_id + data sequence number -> the derived value (the data
1341        // sequence number) for every row. This is the only value-producing arm.
1342        let task = last_updated_seq_task(file_path, Some(42), Some(7));
1343
1344        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1345        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1346        let batches: Vec<RecordBatch> = reader
1347            .read(tasks)
1348            .unwrap()
1349            .stream()
1350            .try_collect()
1351            .await
1352            .unwrap();
1353
1354        assert_last_updated_seq_column(&batches, &[Some(7), Some(7), Some(7)]);
1355    }
1356
1357    #[tokio::test]
1358    async fn test_last_updated_sequence_number_mixed_files_share_schema() {
1359        use arrow_select::concat::concat_batches;
1360
1361        let tmp_dir = TempDir::new().unwrap();
1362        let dir = tmp_dir.path().to_str().unwrap();
1363
1364        // Three files in one scan exercising all three column paths, which must all
1365        // produce the SAME Arrow type (run-end-encoded) or concatenation fails:
1366        //   - constant: first_row_id set, no physical column -> derived constant
1367        //   - null gate: no first_row_id -> null column
1368        //   - coalesce: first_row_id set, physical column present -> per-row + fallback
1369        let constant = last_updated_seq_task(
1370            write_plain_parquet(dir, "constant.parquet", vec![], vec![]),
1371            Some(42),
1372            Some(7),
1373        );
1374        let nulled = last_updated_seq_task(
1375            write_plain_parquet(dir, "nulled.parquet", vec![], vec![]),
1376            None,
1377            Some(7),
1378        );
1379        let coalesced = last_updated_seq_task(
1380            write_plain_parquet(
1381                dir,
1382                "coalesced.parquet",
1383                vec![physical_last_updated_seq_field()],
1384                vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
1385            ),
1386            Some(50),
1387            Some(7),
1388        );
1389
1390        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1391        let tasks = Box::pin(futures::stream::iter(vec![
1392            Ok(constant),
1393            Ok(nulled),
1394            Ok(coalesced),
1395        ])) as FileScanTaskStream;
1396        let batches: Vec<RecordBatch> = reader
1397            .read(tasks)
1398            .unwrap()
1399            .stream()
1400            .try_collect()
1401            .await
1402            .unwrap();
1403
1404        assert_eq!(batches.len(), 3);
1405        // Identical schema across all three paths -> concat succeeds.
1406        let schema = batches[0].schema();
1407        concat_batches(&schema, &batches)
1408            .expect("constant, null and coalesce files must share one column type");
1409    }
1410
1411    /// A parquet field carrying the embedded `_last_updated_sequence_number` field id.
1412    fn physical_last_updated_seq_field() -> Field {
1413        use crate::metadata_columns::{
1414            RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
1415            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
1416        };
1417        Field::new(
1418            RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
1419            DataType::Int64,
1420            true,
1421        )
1422        .with_metadata(HashMap::from([(
1423            PARQUET_FIELD_ID_META_KEY.to_string(),
1424            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER.to_string(),
1425        )]))
1426    }
1427
1428    #[tokio::test]
1429    async fn test_last_updated_sequence_number_physical_column_coalesced() {
1430        let tmp_dir = TempDir::new().unwrap();
1431        let dir = tmp_dir.path().to_str().unwrap();
1432        // A file that physically carries the column, as Iceberg Java writes when
1433        // carrying rows forward across a rewrite: some rows have a stored value, some
1434        // are null (added/modified rows, inherited on read).
1435        let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1436        let file_path = write_plain_parquet(
1437            dir,
1438            "with_seq.parquet",
1439            vec![physical_last_updated_seq_field()],
1440            vec![seq_col],
1441        );
1442
1443        let task = last_updated_seq_task(file_path, Some(100), Some(9));
1444
1445        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1446        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1447        let batches: Vec<RecordBatch> = reader
1448            .read(tasks)
1449            .unwrap()
1450            .stream()
1451            .try_collect()
1452            .await
1453            .unwrap();
1454
1455        // Per-row value where non-null; the data sequence number (9) where null.
1456        assert_last_updated_seq_column(&batches, &[Some(5), Some(9), Some(8)]);
1457    }
1458
1459    #[tokio::test]
1460    async fn test_last_updated_sequence_number_coalesced_with_pos_column() {
1461        use crate::metadata_columns::{
1462            RESERVED_COL_NAME_POS, RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
1463            RESERVED_FIELD_ID_POS,
1464        };
1465
1466        let tmp_dir = TempDir::new().unwrap();
1467        let dir = tmp_dir.path().to_str().unwrap();
1468        let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1469        let file_path = write_plain_parquet(
1470            dir,
1471            "with_seq_and_pos.parquet",
1472            vec![physical_last_updated_seq_field()],
1473            vec![seq_col],
1474        );
1475
1476        // Co-project `_pos` (a virtual column appended to the Arrow output schema) with the
1477        // physical coalesce column. This guards that the physical column's index is
1478        // resolved in the Parquet schema, not the Arrow schema (whose indices shift once
1479        // virtual columns are appended).
1480        let schema = Arc::new(
1481            Schema::builder()
1482                .with_schema_id(1)
1483                .with_fields(vec![
1484                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1485                ])
1486                .build()
1487                .unwrap(),
1488        );
1489        let task = FileScanTask::builder()
1490            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1491            .with_start(0)
1492            .with_length(0)
1493            .with_data_file_path(file_path)
1494            .with_data_file_format(DataFileFormat::Parquet)
1495            .with_schema(schema)
1496            .with_project_field_ids(vec![
1497                1,
1498                RESERVED_FIELD_ID_POS,
1499                RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
1500            ])
1501            .with_first_row_id(Some(100))
1502            .with_data_sequence_number(Some(9))
1503            .with_case_sensitive(false)
1504            .build()
1505            .unwrap();
1506
1507        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1508        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1509        let batches: Vec<RecordBatch> = reader
1510            .read(tasks)
1511            .unwrap()
1512            .stream()
1513            .try_collect()
1514            .await
1515            .unwrap();
1516
1517        // The seq column still coalesces correctly...
1518        assert_last_updated_seq_column(&batches, &[Some(5), Some(9), Some(8)]);
1519        // ...and `_pos` is the row position, unaffected by the physical-column union.
1520        let pos_col = batches[0]
1521            .column_by_name(RESERVED_COL_NAME_POS)
1522            .expect("_pos column should be present")
1523            .as_primitive::<arrow_array::types::Int64Type>();
1524        assert_eq!(pos_col.values(), &[0, 1, 2]);
1525    }
1526
1527    #[tokio::test]
1528    async fn test_last_updated_sequence_number_physical_column_nulled_without_first_row_id() {
1529        let tmp_dir = TempDir::new().unwrap();
1530        let dir = tmp_dir.path().to_str().unwrap();
1531        let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1532        let file_path = write_plain_parquet(
1533            dir,
1534            "with_seq_no_first_row_id.parquet",
1535            vec![physical_last_updated_seq_field()],
1536            vec![seq_col],
1537        );
1538
1539        // Null first_row_id: the whole column is null even though the file physically
1540        // carries per-row values -- the gate wins, and the physical column is not read.
1541        let task = last_updated_seq_task(file_path, None, Some(9));
1542
1543        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1544        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1545        let batches: Vec<RecordBatch> = reader
1546            .read(tasks)
1547            .unwrap()
1548            .stream()
1549            .try_collect()
1550            .await
1551            .unwrap();
1552
1553        assert_last_updated_seq_column(&batches, &[None, None, None]);
1554    }
1555
1556    #[tokio::test]
1557    async fn test_last_updated_sequence_number_present_by_name_without_id_unsupported() {
1558        let tmp_dir = TempDir::new().unwrap();
1559        let dir = tmp_dir.path().to_str().unwrap();
1560        // Column present by name but WITHOUT the embedded field id (e.g. name mapping /
1561        // positional fallback). The transformer keys the source column by field id, so
1562        // this shape can't be threaded and is rejected loudly.
1563        let seq_field = Field::new(
1564            crate::metadata_columns::RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
1565            DataType::Int64,
1566            true,
1567        );
1568        let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1569        let file_path =
1570            write_plain_parquet(dir, "with_seq_by_name.parquet", vec![seq_field], vec![
1571                seq_col,
1572            ]);
1573
1574        let task = last_updated_seq_task(file_path, Some(100), Some(9));
1575
1576        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1577        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1578        let result: Result<Vec<RecordBatch>, _> =
1579            reader.read(tasks).unwrap().stream().try_collect().await;
1580
1581        let err = result.unwrap_err();
1582        assert_eq!(err.kind(), crate::ErrorKind::FeatureUnsupported);
1583        assert!(
1584            format!("{err}").contains("without an embedded field id"),
1585            "unexpected error: {err}"
1586        );
1587    }
1588
1589    #[tokio::test]
1590    async fn test_last_updated_sequence_number_physical_column_first_row_id_without_data_seq() {
1591        let tmp_dir = TempDir::new().unwrap();
1592        let dir = tmp_dir.path().to_str().unwrap();
1593        let seq_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1594        let file_path = write_plain_parquet(
1595            dir,
1596            "with_seq_no_data_seq.parquet",
1597            vec![physical_last_updated_seq_field()],
1598            vec![seq_col],
1599        );
1600
1601        // first_row_id set but no data sequence number: after manifest inheritance a
1602        // committed entry always has one, so this is a malformed manifest, rejected loudly
1603        // rather than nulled.
1604        let task = last_updated_seq_task(file_path, Some(100), None);
1605
1606        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1607        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1608        let result: Result<Vec<RecordBatch>, _> =
1609            reader.read(tasks).unwrap().stream().try_collect().await;
1610
1611        let err = result.unwrap_err();
1612        assert_eq!(err.kind(), crate::ErrorKind::DataInvalid);
1613        assert!(
1614            format!("{err}").contains("no data sequence number"),
1615            "unexpected error: {err}"
1616        );
1617    }
1618
1619    /// A scan task projecting `id` + `_row_id`, with the given `first_row_id`.
1620    fn row_id_task(file_path: String, first_row_id: Option<i64>) -> FileScanTask {
1621        row_id_task_with_options(file_path, first_row_id, 0, 0, vec![])
1622    }
1623
1624    fn row_id_task_with_options(
1625        file_path: String,
1626        first_row_id: Option<i64>,
1627        start: u64,
1628        length: u64,
1629        deletes: Vec<FileScanTaskDeleteFile>,
1630    ) -> FileScanTask {
1631        let schema = Arc::new(
1632            Schema::builder()
1633                .with_schema_id(1)
1634                .with_fields(vec![
1635                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1636                ])
1637                .build()
1638                .unwrap(),
1639        );
1640
1641        FileScanTask::builder()
1642            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1643            .with_start(start)
1644            .with_length(length)
1645            .with_data_file_path(file_path)
1646            .with_data_file_format(DataFileFormat::Parquet)
1647            .with_schema(schema)
1648            .with_project_field_ids(vec![1, RESERVED_FIELD_ID_ROW_ID])
1649            .with_first_row_id(first_row_id)
1650            .with_deletes(deletes)
1651            .with_case_sensitive(false)
1652            .build()
1653            .unwrap()
1654    }
1655
1656    /// Asserts the logical per-row values of the `_row_id` column across all batches,
1657    /// independent of the physical (run-end) encoding.
1658    fn assert_row_id_column(batches: &[RecordBatch], expected: &[Option<i64>]) {
1659        use arrow_array::cast::AsArray;
1660        use arrow_cast::cast;
1661        use arrow_schema::DataType;
1662
1663        let mut actual = Vec::new();
1664        for batch in batches {
1665            let col = batch
1666                .column_by_name(RESERVED_COL_NAME_ROW_ID)
1667                .expect("_row_id column should be present");
1668            let logical = cast(col, &DataType::Int64).unwrap();
1669            let values = logical.as_primitive::<arrow_array::types::Int64Type>();
1670            for i in 0..values.len() {
1671                actual.push((!values.is_null(i)).then(|| values.value(i)));
1672            }
1673        }
1674        assert_eq!(actual, expected);
1675    }
1676
1677    /// A parquet field carrying the embedded `_row_id` field id.
1678    fn physical_row_id_field() -> Field {
1679        Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true).with_metadata(HashMap::from([
1680            (
1681                PARQUET_FIELD_ID_META_KEY.to_string(),
1682                RESERVED_FIELD_ID_ROW_ID.to_string(),
1683            ),
1684        ]))
1685    }
1686
1687    #[tokio::test]
1688    async fn test_row_id_synthesized_from_first_row_id_and_pos() {
1689        let tmp_dir = TempDir::new().unwrap();
1690        let dir = tmp_dir.path().to_str().unwrap();
1691        let file_path = write_plain_parquet(dir, "row_id_synth.parquet", vec![], vec![]);
1692
1693        // No physical column: every row is first_row_id + pos.
1694        let task = row_id_task(file_path, Some(100));
1695
1696        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1697        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1698        let batches: Vec<RecordBatch> = reader
1699            .read(tasks)
1700            .unwrap()
1701            .stream()
1702            .try_collect()
1703            .await
1704            .unwrap();
1705
1706        assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
1707    }
1708
1709    #[tokio::test]
1710    async fn test_row_id_physical_column_coalesced() {
1711        let tmp_dir = TempDir::new().unwrap();
1712        let dir = tmp_dir.path().to_str().unwrap();
1713        // A file that physically carries `_row_id`, as written when carrying rows forward
1714        // across a rewrite: some rows have a stored value, some are null.
1715        let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1716        let file_path = write_plain_parquet(
1717            dir,
1718            "row_id_phys.parquet",
1719            vec![physical_row_id_field()],
1720            vec![id_col],
1721        );
1722
1723        let task = row_id_task(file_path, Some(100));
1724
1725        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1726        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1727        let batches: Vec<RecordBatch> = reader
1728            .read(tasks)
1729            .unwrap()
1730            .stream()
1731            .try_collect()
1732            .await
1733            .unwrap();
1734
1735        // Per-row value where non-null; first_row_id + pos (101) where null.
1736        assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
1737    }
1738
1739    #[tokio::test]
1740    async fn test_row_id_only_synthesis_reads_no_data_columns() {
1741        // The common v3 case: a new-row file with `first_row_id` set and NO physically
1742        // stored `_row_id`, projecting only `_row_id`. `_row_id` synthesis installs the
1743        // RowNumber virtual column (via `need_row_number`), so the row count comes from it
1744        // -- the scan must read no data columns, not fall back to reading everything.
1745        let tmp_dir = TempDir::new().unwrap();
1746        let dir = tmp_dir.path().to_str().unwrap();
1747
1748        let meta_only = metadata_projection_task_with_first_row_id(
1749            write_parquet_with_wide_column(dir, "row_id_only.parquet", vec![], vec![]),
1750            id_and_wide_schema(),
1751            vec![RESERVED_FIELD_ID_ROW_ID],
1752            Some(100),
1753        );
1754        let (batches, meta_only_bytes) = scan_task(meta_only).await;
1755
1756        assert_eq!(batches[0].num_columns(), 1);
1757        assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
1758
1759        // A scan that also projects the wide data column must read materially more.
1760        let with_data = metadata_projection_task_with_first_row_id(
1761            write_parquet_with_wide_column(dir, "row_id_only_ref.parquet", vec![], vec![]),
1762            id_and_wide_schema(),
1763            vec![2, RESERVED_FIELD_ID_ROW_ID],
1764            Some(100),
1765        );
1766        let (_, with_data_bytes) = scan_task(with_data).await;
1767
1768        assert!(
1769            meta_only_bytes < with_data_bytes,
1770            "_row_id-only synthesis should read fewer bytes than a scan of the wide column: \
1771             {meta_only_bytes} vs {with_data_bytes}"
1772        );
1773    }
1774
1775    #[tokio::test]
1776    async fn test_row_id_only_null_first_row_id_reads_no_data_columns() {
1777        // A null `first_row_id` (v1/v2, or a pre-upgrade v3 snapshot) nulls the whole
1778        // `_row_id` column, so nothing is synthesized -- but the column still has a length.
1779        // The RowNumber counter installed for the metadata-only projection supplies it, so
1780        // the scan reads no data columns instead of reading everything just for the count.
1781        let tmp_dir = TempDir::new().unwrap();
1782        let dir = tmp_dir.path().to_str().unwrap();
1783
1784        let meta_only = metadata_projection_task_with_first_row_id(
1785            write_parquet_with_wide_column(dir, "row_id_null.parquet", vec![], vec![]),
1786            id_and_wide_schema(),
1787            vec![RESERVED_FIELD_ID_ROW_ID],
1788            None,
1789        );
1790        let (batches, meta_only_bytes) = scan_task(meta_only).await;
1791
1792        assert_eq!(batches[0].num_columns(), 1);
1793        assert_row_id_column(&batches, &[None, None, None]);
1794
1795        let with_data = metadata_projection_task_with_first_row_id(
1796            write_parquet_with_wide_column(dir, "row_id_null_ref.parquet", vec![], vec![]),
1797            id_and_wide_schema(),
1798            vec![2, RESERVED_FIELD_ID_ROW_ID],
1799            None,
1800        );
1801        let (_, with_data_bytes) = scan_task(with_data).await;
1802
1803        assert!(
1804            meta_only_bytes < with_data_bytes,
1805            "_row_id-only scan with a null first_row_id should read fewer bytes than a scan \
1806             of the wide column: {meta_only_bytes} vs {with_data_bytes}"
1807        );
1808    }
1809
1810    #[tokio::test]
1811    async fn test_last_updated_seq_only_reads_no_data_columns() {
1812        use crate::metadata_columns::RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER;
1813
1814        // `SELECT _last_updated_sequence_number` derives a per-row constant (the data
1815        // sequence number) with no physical column to read. The metadata-only RowNumber
1816        // counter supplies the row count, so the scan prunes the data columns rather than
1817        // reading them all just to size the constant.
1818        let tmp_dir = TempDir::new().unwrap();
1819        let dir = tmp_dir.path().to_str().unwrap();
1820
1821        let lusn_only_task = |file_path: String, project_field_ids: Vec<i32>| {
1822            FileScanTask::builder()
1823                .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1824                .with_start(0)
1825                .with_length(0)
1826                .with_data_file_path(file_path)
1827                .with_data_file_format(DataFileFormat::Parquet)
1828                .with_schema(id_and_wide_schema())
1829                .with_project_field_ids(project_field_ids)
1830                .with_first_row_id(Some(42))
1831                .with_data_sequence_number(Some(7))
1832                .with_case_sensitive(false)
1833                .build()
1834                .unwrap()
1835        };
1836
1837        let meta_only = lusn_only_task(
1838            write_parquet_with_wide_column(dir, "lusn_only.parquet", vec![], vec![]),
1839            vec![RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER],
1840        );
1841        let (batches, meta_only_bytes) = scan_task(meta_only).await;
1842
1843        assert_eq!(batches[0].num_columns(), 1);
1844        assert_last_updated_seq_column(&batches, &[Some(7), Some(7), Some(7)]);
1845
1846        let with_data = lusn_only_task(
1847            write_parquet_with_wide_column(dir, "lusn_only_ref.parquet", vec![], vec![]),
1848            vec![2, RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER],
1849        );
1850        let (_, with_data_bytes) = scan_task(with_data).await;
1851
1852        assert!(
1853            meta_only_bytes < with_data_bytes,
1854            "_last_updated_sequence_number-only scan should read fewer bytes than a scan of \
1855             the wide column: {meta_only_bytes} vs {with_data_bytes}"
1856        );
1857    }
1858
1859    #[tokio::test]
1860    async fn test_partition_only_reads_no_data_columns() {
1861        use arrow_array::StructArray;
1862
1863        use crate::metadata_columns::{RESERVED_COL_NAME_PARTITION, RESERVED_FIELD_ID_PARTITION};
1864        use crate::spec::{Literal, PartitionSpec, Struct, Transform};
1865
1866        // `SELECT _partition` materializes a struct constant from the task's partition
1867        // metadata, with no physical column to read. It is the only struct-constant metadata
1868        // column; the RowNumber counter sizes it, so the scan prunes the data columns.
1869        let tmp_dir = TempDir::new().unwrap();
1870        let dir = tmp_dir.path().to_str().unwrap();
1871        let schema = id_and_wide_schema();
1872        let spec = Arc::new(
1873            PartitionSpec::builder(schema.clone())
1874                .with_spec_id(7)
1875                .add_partition_field("id", "id", Transform::Identity)
1876                .unwrap()
1877                .build()
1878                .unwrap(),
1879        );
1880        let unified_type = Arc::new(spec.partition_type(&schema).unwrap());
1881        let partition_data = Struct::from_iter(vec![Some(Literal::int(42))]);
1882
1883        let partition_task = |file_path: String, project_field_ids: Vec<i32>| {
1884            FileScanTask::builder()
1885                .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1886                .with_start(0)
1887                .with_length(0)
1888                .with_data_file_path(file_path)
1889                .with_data_file_format(DataFileFormat::Parquet)
1890                .with_schema(schema.clone())
1891                .with_project_field_ids(project_field_ids)
1892                .with_partition_spec(Some(spec.clone()))
1893                .with_partition(Some(partition_data.clone()))
1894                .with_unified_partition_type(Some(unified_type.clone()))
1895                .with_case_sensitive(false)
1896                .build()
1897                .unwrap()
1898        };
1899
1900        let meta_only = partition_task(
1901            write_parquet_with_wide_column(dir, "partition_only.parquet", vec![], vec![]),
1902            vec![RESERVED_FIELD_ID_PARTITION],
1903        );
1904        let (batches, meta_only_bytes) = scan_task(meta_only).await;
1905
1906        assert_eq!(batches[0].num_columns(), 1);
1907        let partition_col = batches[0]
1908            .column_by_name(RESERVED_COL_NAME_PARTITION)
1909            .expect("_partition column should be present")
1910            .as_any()
1911            .downcast_ref::<StructArray>()
1912            .unwrap();
1913        assert_eq!(partition_col.len(), 3);
1914        let inner = partition_col
1915            .column(0)
1916            .as_any()
1917            .downcast_ref::<Int32Array>()
1918            .unwrap();
1919        assert_eq!(inner.values(), &[42, 42, 42]);
1920
1921        let with_data = partition_task(
1922            write_parquet_with_wide_column(dir, "partition_only_ref.parquet", vec![], vec![]),
1923            vec![2, RESERVED_FIELD_ID_PARTITION],
1924        );
1925        let (_, with_data_bytes) = scan_task(with_data).await;
1926
1927        assert!(
1928            meta_only_bytes < with_data_bytes,
1929            "_partition-only scan should read fewer bytes than a scan of the wide column: \
1930             {meta_only_bytes} vs {with_data_bytes}"
1931        );
1932    }
1933
1934    #[tokio::test]
1935    async fn test_row_id_resolves_alongside_id_less_leaf() {
1936        // A file with an id-less leaf (mimicking a Variant column's internal metadata/value
1937        // leaves, which the spec requires to have no field id) plus a physical `_row_id`
1938        // that carries its embedded id. The reserved id must still resolve -- an
1939        // all-or-nothing field map would bail on the id-less leaf and wrongly reject the file.
1940        let tmp_dir = TempDir::new().unwrap();
1941        let dir = tmp_dir.path().to_str().unwrap();
1942        let idless_field = Field::new("variant_internal", DataType::Utf8, true);
1943        let idless_col = Arc::new(StringArray::from(vec!["a", "b", "c"])) as ArrayRef;
1944        let row_id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1945        let file_path = write_plain_parquet(
1946            dir,
1947            "row_id_with_idless_leaf.parquet",
1948            vec![idless_field, physical_row_id_field()],
1949            vec![idless_col, row_id_col],
1950        );
1951
1952        let task = row_id_task(file_path, Some(100));
1953
1954        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
1955        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
1956        let batches: Vec<RecordBatch> = reader
1957            .read(tasks)
1958            .unwrap()
1959            .stream()
1960            .try_collect()
1961            .await
1962            .unwrap();
1963
1964        // Physical value where non-null; first_row_id + pos (101) where null.
1965        assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
1966    }
1967
1968    #[tokio::test]
1969    async fn test_row_id_and_last_updated_seq_co_projected() {
1970        use crate::metadata_columns::RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER;
1971
1972        // Both lineage columns projected together over a file carrying both physical
1973        // leaves. Each must materialize independently -- neither leaf's mask clobbers the
1974        // other, and the two synthesized columns keep their own values.
1975        let tmp_dir = TempDir::new().unwrap();
1976        let dir = tmp_dir.path().to_str().unwrap();
1977        let row_id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
1978        let seq_col = Arc::new(Int64Array::from(vec![Some(50), None, Some(70)])) as ArrayRef;
1979        let file_path = write_plain_parquet(
1980            dir,
1981            "row_id_and_seq.parquet",
1982            vec![physical_row_id_field(), physical_last_updated_seq_field()],
1983            vec![row_id_col, seq_col],
1984        );
1985
1986        let schema = Arc::new(
1987            Schema::builder()
1988                .with_schema_id(1)
1989                .with_fields(vec![
1990                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
1991                ])
1992                .build()
1993                .unwrap(),
1994        );
1995        let task = FileScanTask::builder()
1996            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
1997            .with_start(0)
1998            .with_length(0)
1999            .with_data_file_path(file_path)
2000            .with_data_file_format(DataFileFormat::Parquet)
2001            .with_schema(schema)
2002            .with_project_field_ids(vec![
2003                1,
2004                RESERVED_FIELD_ID_ROW_ID,
2005                RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
2006            ])
2007            .with_first_row_id(Some(100))
2008            .with_data_sequence_number(Some(9))
2009            .with_case_sensitive(false)
2010            .build()
2011            .unwrap();
2012
2013        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2014        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2015        let batches: Vec<RecordBatch> = reader
2016            .read(tasks)
2017            .unwrap()
2018            .stream()
2019            .try_collect()
2020            .await
2021            .unwrap();
2022
2023        // _row_id: physical value where non-null, else first_row_id + pos (101).
2024        assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
2025        // _last_updated_sequence_number: physical value where non-null, else data seq (9).
2026        assert_last_updated_seq_column(&batches, &[Some(50), Some(9), Some(70)]);
2027    }
2028
2029    #[tokio::test]
2030    async fn test_row_id_null_when_no_first_row_id() {
2031        let tmp_dir = TempDir::new().unwrap();
2032        let dir = tmp_dir.path().to_str().unwrap();
2033        // Physically carries `_row_id`, but the file has a null first_row_id.
2034        let id_col = Arc::new(Int64Array::from(vec![Some(5), Some(6), Some(7)])) as ArrayRef;
2035        let file_path = write_plain_parquet(
2036            dir,
2037            "row_id_no_first.parquet",
2038            vec![physical_row_id_field()],
2039            vec![id_col],
2040        );
2041
2042        // Null first_row_id: the whole column is null; the physical values are not read.
2043        let task = row_id_task(file_path, None);
2044
2045        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2046        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2047        let batches: Vec<RecordBatch> = reader
2048            .read(tasks)
2049            .unwrap()
2050            .stream()
2051            .try_collect()
2052            .await
2053            .unwrap();
2054
2055        assert_row_id_column(&batches, &[None, None, None]);
2056    }
2057
2058    #[tokio::test]
2059    async fn test_row_id_with_pos_column() {
2060        use crate::metadata_columns::RESERVED_COL_NAME_POS;
2061
2062        let tmp_dir = TempDir::new().unwrap();
2063        let dir = tmp_dir.path().to_str().unwrap();
2064        let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2065        let file_path = write_plain_parquet(
2066            dir,
2067            "row_id_and_pos.parquet",
2068            vec![physical_row_id_field()],
2069            vec![id_col],
2070        );
2071
2072        // Co-project `_pos` and `_row_id`. `_row_id` synthesis consumes the position, and
2073        // `_pos` is also emitted -- the RowNumber column must be added once and the two
2074        // must not interfere.
2075        let schema = Arc::new(
2076            Schema::builder()
2077                .with_schema_id(1)
2078                .with_fields(vec![
2079                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2080                ])
2081                .build()
2082                .unwrap(),
2083        );
2084        let task = FileScanTask::builder()
2085            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
2086            .with_start(0)
2087            .with_length(0)
2088            .with_data_file_path(file_path)
2089            .with_data_file_format(DataFileFormat::Parquet)
2090            .with_schema(schema)
2091            .with_project_field_ids(vec![1, RESERVED_FIELD_ID_POS, RESERVED_FIELD_ID_ROW_ID])
2092            .with_first_row_id(Some(100))
2093            .with_case_sensitive(false)
2094            .build()
2095            .unwrap();
2096
2097        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2098        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2099        let batches: Vec<RecordBatch> = reader
2100            .read(tasks)
2101            .unwrap()
2102            .stream()
2103            .try_collect()
2104            .await
2105            .unwrap();
2106
2107        // `_row_id` coalesces correctly...
2108        assert_row_id_column(&batches, &[Some(5), Some(101), Some(8)]);
2109        // ...and `_pos` is the row position, not double-counted.
2110        let pos_col = batches[0]
2111            .column_by_name(RESERVED_COL_NAME_POS)
2112            .expect("_pos column should be present")
2113            .as_primitive::<arrow_array::types::Int64Type>();
2114        assert_eq!(pos_col.values(), &[0, 1, 2]);
2115    }
2116
2117    #[tokio::test]
2118    async fn test_row_id_mixed_files_share_schema() {
2119        use arrow_select::concat::concat_batches;
2120
2121        let tmp_dir = TempDir::new().unwrap();
2122        let dir = tmp_dir.path().to_str().unwrap();
2123
2124        // Three files in one scan exercising all three column paths, which must all
2125        // produce the SAME Arrow type (plain Int64) or concatenation fails:
2126        //   - synthesis: first_row_id set, no physical column -> first_row_id + pos
2127        //   - null gate: no first_row_id -> null column
2128        //   - coalesce: first_row_id set, physical column present -> per-row + fallback
2129        let synth = row_id_task(
2130            write_plain_parquet(dir, "row_id_synth2.parquet", vec![], vec![]),
2131            Some(42),
2132        );
2133        let nulled = row_id_task(
2134            write_plain_parquet(dir, "row_id_null2.parquet", vec![], vec![]),
2135            None,
2136        );
2137        let coalesced = row_id_task(
2138            write_plain_parquet(
2139                dir,
2140                "row_id_coalesced2.parquet",
2141                vec![physical_row_id_field()],
2142                vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
2143            ),
2144            Some(50),
2145        );
2146
2147        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2148        let tasks = Box::pin(futures::stream::iter(vec![
2149            Ok(synth),
2150            Ok(nulled),
2151            Ok(coalesced),
2152        ])) as FileScanTaskStream;
2153        let batches: Vec<RecordBatch> = reader
2154            .read(tasks)
2155            .unwrap()
2156            .stream()
2157            .try_collect()
2158            .await
2159            .unwrap();
2160
2161        assert_eq!(batches.len(), 3);
2162        let schema = batches[0].schema();
2163        concat_batches(&schema, &batches)
2164            .expect("synthesis, null and coalesce files must share one column type");
2165    }
2166
2167    #[tokio::test]
2168    async fn test_row_id_present_by_name_without_id_unsupported() {
2169        let tmp_dir = TempDir::new().unwrap();
2170        let dir = tmp_dir.path().to_str().unwrap();
2171        // Column present by name but WITHOUT the embedded field id. The transformer keys
2172        // the source column by field id, so this shape can't be threaded and is rejected.
2173        let id_field = Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true);
2174        let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2175        let file_path =
2176            write_plain_parquet(dir, "row_id_by_name.parquet", vec![id_field], vec![id_col]);
2177
2178        let task = row_id_task(file_path, Some(100));
2179
2180        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2181        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2182        let result: Result<Vec<RecordBatch>, _> =
2183            reader.read(tasks).unwrap().stream().try_collect().await;
2184
2185        let err = result.unwrap_err();
2186        assert_eq!(err.kind(), crate::ErrorKind::FeatureUnsupported);
2187        assert!(
2188            format!("{err}").contains("without an embedded field id"),
2189            "unexpected error: {err}"
2190        );
2191    }
2192
2193    #[tokio::test]
2194    async fn test_row_id_present_by_name_without_id_nulls_when_no_lineage() {
2195        let tmp_dir = TempDir::new().unwrap();
2196        let dir = tmp_dir.path().to_str().unwrap();
2197        // Same name-only shape as the reject test above, but the file carries no row lineage
2198        // (`first_row_id = None`), as a migrated pre-v3 file with a user column named
2199        // `_row_id` would. The physical leaf is never read, so `_row_id` is nulled out
2200        // rather than rejected (matching Java `ValueReaders.rowIds`).
2201        let id_field = Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true);
2202        let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2203        let file_path = write_plain_parquet(
2204            dir,
2205            "row_id_by_name_no_lineage.parquet",
2206            vec![id_field],
2207            vec![id_col],
2208        );
2209
2210        let task = row_id_task(file_path, None);
2211
2212        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2213        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2214        let batches: Vec<RecordBatch> = reader
2215            .read(tasks)
2216            .unwrap()
2217            .stream()
2218            .try_collect()
2219            .await
2220            .unwrap();
2221
2222        assert_row_id_column(&batches, &[None, None, None]);
2223    }
2224
2225    /// Builds a `row_id_task` (see above) that additionally carries a bound predicate,
2226    /// so a `RowSelection` is applied when the reader has row selection enabled.
2227    fn row_id_task_with_predicate(
2228        file_path: String,
2229        first_row_id: Option<i64>,
2230        extra_project_field_ids: Vec<i32>,
2231        predicate: crate::expr::Predicate,
2232    ) -> FileScanTask {
2233        use crate::expr::Bind;
2234
2235        let schema = Arc::new(
2236            Schema::builder()
2237                .with_schema_id(1)
2238                .with_fields(vec![
2239                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2240                ])
2241                .build()
2242                .unwrap(),
2243        );
2244        let bound = predicate.bind(Arc::clone(&schema), false).unwrap();
2245
2246        let mut project_field_ids = vec![1];
2247        project_field_ids.extend(extra_project_field_ids);
2248        project_field_ids.push(RESERVED_FIELD_ID_ROW_ID);
2249
2250        FileScanTask::builder()
2251            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
2252            .with_start(0)
2253            .with_length(0)
2254            .with_data_file_path(file_path)
2255            .with_data_file_format(DataFileFormat::Parquet)
2256            .with_schema(schema)
2257            .with_project_field_ids(project_field_ids)
2258            .with_predicate(Some(bound))
2259            .with_first_row_id(first_row_id)
2260            .with_case_sensitive(false)
2261            .build()
2262            .unwrap()
2263    }
2264
2265    #[tokio::test]
2266    async fn test_row_id_stable_under_row_selection() {
2267        use crate::expr::Reference;
2268        use crate::spec::Datum;
2269
2270        let tmp_dir = TempDir::new().unwrap();
2271        let dir = tmp_dir.path().to_str().unwrap();
2272        // id = [1, 2, 3]; drop the middle physical row via a predicate + row selection.
2273        let file_path = write_plain_parquet(dir, "row_id_selection.parquet", vec![], vec![]);
2274
2275        let task = row_id_task_with_predicate(
2276            file_path,
2277            Some(100),
2278            vec![],
2279            Reference::new("id").not_equal_to(Datum::int(2)),
2280        );
2281
2282        // Row selection must be enabled for the predicate to produce a RowSelection.
2283        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
2284            .with_row_selection_enabled(true)
2285            .build();
2286        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2287        let batches: Vec<RecordBatch> = reader
2288            .read(tasks)
2289            .unwrap()
2290            .stream()
2291            .try_collect()
2292            .await
2293            .unwrap();
2294
2295        // The survivors are physical rows 0 and 2, so their _row_id is first_row_id + the
2296        // PHYSICAL position: [100, 102]. A dense output index would wrongly give [100, 101].
2297        assert_row_id_column(&batches, &[Some(100), Some(102)]);
2298    }
2299
2300    #[tokio::test]
2301    async fn test_row_id_coalesce_stable_under_row_selection() {
2302        use crate::expr::Reference;
2303        use crate::spec::Datum;
2304
2305        let tmp_dir = TempDir::new().unwrap();
2306        let dir = tmp_dir.path().to_str().unwrap();
2307        // Physical _row_id = [Some(5), None, Some(8)] over id = [1, 2, 3]. Dropping the
2308        // middle row must keep the physical column and the RowNumber fallback row-aligned.
2309        let id_col = Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef;
2310        let file_path = write_plain_parquet(
2311            dir,
2312            "row_id_coalesce_selection.parquet",
2313            vec![physical_row_id_field()],
2314            vec![id_col],
2315        );
2316
2317        let task = row_id_task_with_predicate(
2318            file_path,
2319            Some(100),
2320            vec![],
2321            Reference::new("id").not_equal_to(Datum::int(2)),
2322        );
2323
2324        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
2325            .with_row_selection_enabled(true)
2326            .build();
2327        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2328        let batches: Vec<RecordBatch> = reader
2329            .read(tasks)
2330            .unwrap()
2331            .stream()
2332            .try_collect()
2333            .await
2334            .unwrap();
2335
2336        // Rows 0 and 2 survive: their stored values (5, 8) pass through. The dropped
2337        // row's null (which would have fallen back to 100 + 1) is gone -- proving the
2338        // physical column and the positional fallback are filtered by the same selection.
2339        assert_row_id_column(&batches, &[Some(5), Some(8)]);
2340    }
2341
2342    #[tokio::test]
2343    async fn test_row_id_global_across_row_groups() {
2344        let tmp_dir = TempDir::new().unwrap();
2345        let dir = tmp_dir.path().to_str().unwrap();
2346
2347        // 5 rows written with max_row_group_size = 2 -> 3 row groups. `_pos` must be the
2348        // GLOBAL file position, so `_row_id` continues across row-group boundaries rather
2349        // than restarting per group (which would silently duplicate ids).
2350        let file_path = format!("{dir}/row_id_multi_rg.parquet");
2351        let field = Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2352            PARQUET_FIELD_ID_META_KEY.to_string(),
2353            "1".to_string(),
2354        )]));
2355        let arrow_schema = Arc::new(ArrowSchema::new(vec![field]));
2356        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(Int32Array::from(
2357            vec![1, 2, 3, 4, 5],
2358        ))])
2359        .unwrap();
2360        let props = WriterProperties::builder()
2361            .set_compression(Compression::SNAPPY)
2362            .set_max_row_group_row_count(Some(2))
2363            .build();
2364        let file = File::create(&file_path).unwrap();
2365        let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2366        writer.write(&batch).unwrap();
2367        writer.close().unwrap();
2368
2369        let task = row_id_task(file_path, Some(0));
2370        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2371        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2372        let batches: Vec<RecordBatch> = reader
2373            .read(tasks)
2374            .unwrap()
2375            .stream()
2376            .try_collect()
2377            .await
2378            .unwrap();
2379
2380        assert_row_id_column(&batches, &[Some(0), Some(1), Some(2), Some(3), Some(4)]);
2381    }
2382
2383    #[tokio::test]
2384    async fn test_row_id_global_when_first_row_group_pruned() {
2385        use parquet::file::reader::{FileReader, SerializedFileReader};
2386
2387        let tmp_dir = TempDir::new().unwrap();
2388        let dir = tmp_dir.path().to_str().unwrap();
2389
2390        // 6 rows with max_row_group_size = 2 -> 3 row groups. A byte-range split that prunes
2391        // row group 0 reaches the reader via `with_row_groups()` -- a different path than the
2392        // `RowSelection` cases above. The survivors must keep their GLOBAL positions (starting
2393        // at 2), so a per-group RowNumber restart would surface as duplicate ids here.
2394        let file_path = format!("{dir}/row_id_prune_rg0.parquet");
2395        let field = Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2396            PARQUET_FIELD_ID_META_KEY.to_string(),
2397            "1".to_string(),
2398        )]));
2399        let arrow_schema = Arc::new(ArrowSchema::new(vec![field]));
2400        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![Arc::new(Int32Array::from(
2401            vec![1, 2, 3, 4, 5, 6],
2402        ))])
2403        .unwrap();
2404        let props = WriterProperties::builder()
2405            .set_compression(Compression::SNAPPY)
2406            .set_max_row_group_row_count(Some(2))
2407            .build();
2408        let file = File::create(&file_path).unwrap();
2409        let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2410        writer.write(&batch).unwrap();
2411        writer.close().unwrap();
2412
2413        // A byte range starting just past row group 0 prunes it (its midpoint falls below
2414        // `start`) while keeping groups 1 and 2 (physical rows 2..6).
2415        let metadata = SerializedFileReader::new(File::open(&file_path).unwrap())
2416            .unwrap()
2417            .metadata()
2418            .clone();
2419        assert_eq!(metadata.num_row_groups(), 3);
2420        let start = 4 + metadata.row_group(0).compressed_size() as u64;
2421        let file_size = std::fs::metadata(&file_path).unwrap().len();
2422
2423        let task = row_id_task_with_options(file_path, Some(100), start, file_size - start, vec![]);
2424
2425        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2426        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2427        let batches: Vec<RecordBatch> = reader
2428            .read(tasks)
2429            .unwrap()
2430            .stream()
2431            .try_collect()
2432            .await
2433            .unwrap();
2434
2435        // Groups 1 and 2 survive at physical positions 2..6 -> _row_id 102..106, not a
2436        // per-group restart at 100.
2437        assert_row_id_column(&batches, &[Some(102), Some(103), Some(104), Some(105)]);
2438    }
2439
2440    /// Writes a positional delete file (`file_path` + `pos` reserved columns) marking the
2441    /// given `positions` of `data_file_path` as deleted.
2442    fn write_positional_delete(
2443        dir: &str,
2444        name: &str,
2445        data_file_path: &str,
2446        positions: &[i64],
2447    ) -> String {
2448        use arrow_array::StringArray;
2449
2450        let file_path_field =
2451            Field::new("file_path", DataType::Utf8, false).with_metadata(HashMap::from([(
2452                PARQUET_FIELD_ID_META_KEY.to_string(),
2453                "2147483546".to_string(),
2454            )]));
2455        let pos_field = Field::new("pos", DataType::Int64, false).with_metadata(HashMap::from([(
2456            PARQUET_FIELD_ID_META_KEY.to_string(),
2457            "2147483545".to_string(),
2458        )]));
2459        let arrow_schema = Arc::new(ArrowSchema::new(vec![file_path_field, pos_field]));
2460        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
2461            Arc::new(StringArray::from(vec![data_file_path; positions.len()])),
2462            Arc::new(Int64Array::from(positions.to_vec())),
2463        ])
2464        .unwrap();
2465
2466        let path = format!("{dir}/{name}");
2467        let file = File::create(&path).unwrap();
2468        let props = WriterProperties::builder()
2469            .set_compression(Compression::SNAPPY)
2470            .build();
2471        let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2472        writer.write(&batch).unwrap();
2473        writer.close().unwrap();
2474        path
2475    }
2476
2477    #[tokio::test]
2478    async fn test_row_id_survives_positional_delete() {
2479        use crate::spec::DataContentType;
2480
2481        let tmp_dir = TempDir::new().unwrap();
2482        let dir = tmp_dir.path().to_str().unwrap();
2483        // id = [1, 2, 3]; a positional delete drops the middle physical row (pos = 1).
2484        let data_path = write_plain_parquet(dir, "row_id_posdel_data.parquet", vec![], vec![]);
2485        let del_path = write_positional_delete(dir, "row_id_posdel.parquet", &data_path, &[1]);
2486
2487        let delete = FileScanTaskDeleteFile::builder()
2488            .with_file_path(del_path.clone())
2489            .with_file_size_in_bytes(std::fs::metadata(&del_path).unwrap().len())
2490            .with_file_type(DataContentType::PositionDeletes)
2491            .with_file_format(DataFileFormat::Parquet)
2492            .with_partition_spec_id(0)
2493            .build()
2494            .unwrap();
2495        let task = row_id_task_with_options(data_path, Some(100), 0, 0, vec![delete]);
2496
2497        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2498        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2499        let batches: Vec<RecordBatch> = reader
2500            .read(tasks)
2501            .unwrap()
2502            .stream()
2503            .try_collect()
2504            .await
2505            .unwrap();
2506
2507        // Survivors keep their ABSOLUTE positions -- [100, 102], not renumbered [100, 101].
2508        // Positional-delete selection reaches the reader via a different path than predicate
2509        // selection, so this covers it explicitly.
2510        assert_row_id_column(&batches, &[Some(100), Some(102)]);
2511    }
2512
2513    #[tokio::test]
2514    async fn test_row_id_name_collision_under_positional_fallback() {
2515        let tmp_dir = TempDir::new().unwrap();
2516        let dir = tmp_dir.path().to_str().unwrap();
2517
2518        // A file with NO embedded field ids (positional fallback) whose columns include one
2519        // literally named `_row_id` (user data). Projecting `_row_id` must NOT be rejected as
2520        // an unthreadable physical metadata column -- under fallback the reserved column is
2521        // synthesized and the same-named user column is just data.
2522        let file_path = format!("{dir}/fallback_row_id_name.parquet");
2523        let arrow_schema = Arc::new(ArrowSchema::new(vec![
2524            Field::new("id", DataType::Int32, false),
2525            Field::new(RESERVED_COL_NAME_ROW_ID, DataType::Int64, true),
2526        ]));
2527        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![
2528            Arc::new(Int32Array::from(vec![1, 2, 3])),
2529            Arc::new(Int64Array::from(vec![7i64, 8, 9])),
2530        ])
2531        .unwrap();
2532        let props = WriterProperties::builder()
2533            .set_compression(Compression::SNAPPY)
2534            .build();
2535        let file = File::create(&file_path).unwrap();
2536        let mut writer = ArrowWriter::try_new(file, arrow_schema, Some(props)).unwrap();
2537        writer.write(&batch).unwrap();
2538        writer.close().unwrap();
2539
2540        let task = row_id_task(file_path, Some(100));
2541        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
2542        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2543        let batches: Vec<RecordBatch> = reader
2544            .read(tasks)
2545            .unwrap()
2546            .stream()
2547            .try_collect()
2548            .await
2549            .unwrap();
2550
2551        // Not rejected; `_row_id` is synthesized as first_row_id + pos.
2552        assert_row_id_column(&batches, &[Some(100), Some(101), Some(102)]);
2553    }
2554
2555    #[tokio::test]
2556    async fn test_read_encrypted_parquet_with_wrong_key_fails() {
2557        let encryption_key = b"0123456789abcdef";
2558        let wrong_key = b"fedcba9876543210";
2559
2560        let schema = Arc::new(
2561            Schema::builder()
2562                .with_schema_id(1)
2563                .with_fields(vec![
2564                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2565                ])
2566                .build()
2567                .unwrap(),
2568        );
2569
2570        let arrow_schema = Arc::new(ArrowSchema::new(vec![
2571            Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2572                PARQUET_FIELD_ID_META_KEY.to_string(),
2573                "1".to_string(),
2574            )])),
2575        ]));
2576
2577        let tmp_dir = TempDir::new().unwrap();
2578        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2579        let file_io = FileIO::new_with_fs();
2580
2581        let id_data = Arc::new(Int32Array::from(vec![1, 2, 3])) as ArrayRef;
2582        let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_data]).unwrap();
2583
2584        let file_path = format!("{table_location}/encrypted_wrong_key.parquet");
2585        write_encrypted_parquet(&file_path, &batch, encryption_key, None);
2586
2587        let wrong_key_metadata = crate::encryption::StandardKeyMetadata::try_new(wrong_key)
2588            .unwrap()
2589            .encode()
2590            .unwrap();
2591
2592        let reader = ArrowReaderBuilder::new(file_io, Runtime::current()).build();
2593
2594        let task = FileScanTask::builder()
2595            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
2596            .with_start(0)
2597            .with_length(0)
2598            .with_data_file_path(file_path)
2599            .with_data_file_format(DataFileFormat::Parquet)
2600            .with_schema(schema)
2601            .with_project_field_ids(vec![1])
2602            .with_case_sensitive(false)
2603            .with_key_metadata(Some(wrong_key_metadata))
2604            .build()
2605            .unwrap();
2606
2607        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
2608        let result: Result<Vec<RecordBatch>, _> =
2609            reader.read(tasks).unwrap().stream().try_collect().await;
2610
2611        let err = result.unwrap_err();
2612        assert_eq!(err.kind(), crate::ErrorKind::Unexpected);
2613        let err_str = format!("{err}");
2614        assert!(
2615            err_str.contains("unable to decrypt parquet footer"),
2616            "Expected error about decryption failure, got: {err_str}"
2617        );
2618    }
2619
2620    /// Test that concurrency=1 reads all files correctly and in deterministic order.
2621    /// This verifies the fast-path optimization for single concurrency.
2622    #[tokio::test]
2623    async fn test_read_with_concurrency_one() {
2624        use arrow_array::Int32Array;
2625
2626        let schema = Arc::new(
2627            Schema::builder()
2628                .with_schema_id(1)
2629                .with_fields(vec![
2630                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2631                    NestedField::required(2, "file_num", Type::Primitive(PrimitiveType::Int))
2632                        .into(),
2633                ])
2634                .build()
2635                .unwrap(),
2636        );
2637
2638        let arrow_schema = Arc::new(ArrowSchema::new(vec![
2639            Field::new("id", DataType::Int32, false).with_metadata(HashMap::from([(
2640                PARQUET_FIELD_ID_META_KEY.to_string(),
2641                "1".to_string(),
2642            )])),
2643            Field::new("file_num", DataType::Int32, false).with_metadata(HashMap::from([(
2644                PARQUET_FIELD_ID_META_KEY.to_string(),
2645                "2".to_string(),
2646            )])),
2647        ]));
2648
2649        let tmp_dir = TempDir::new().unwrap();
2650        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2651        let file_io = FileIO::new_with_fs();
2652
2653        // Create 3 parquet files with different data
2654        let props = WriterProperties::builder()
2655            .set_compression(Compression::SNAPPY)
2656            .build();
2657
2658        for file_num in 0..3 {
2659            let id_data = Arc::new(Int32Array::from_iter_values(
2660                file_num * 10..(file_num + 1) * 10,
2661            )) as ArrayRef;
2662            let file_num_data = Arc::new(Int32Array::from(vec![file_num; 10])) as ArrayRef;
2663
2664            let to_write =
2665                RecordBatch::try_new(arrow_schema.clone(), vec![id_data, file_num_data]).unwrap();
2666
2667            let file = File::create(format!("{table_location}/file_{file_num}.parquet")).unwrap();
2668            let mut writer =
2669                ArrowWriter::try_new(file, to_write.schema(), Some(props.clone())).unwrap();
2670            writer.write(&to_write).expect("Writing batch");
2671            writer.close().unwrap();
2672        }
2673
2674        // Read with concurrency=1 (fast-path)
2675        let reader = ArrowReaderBuilder::new(file_io, Runtime::current())
2676            .with_data_file_concurrency_limit(1)
2677            .build();
2678
2679        // Create tasks in a specific order: file_0, file_1, file_2
2680        let tasks = vec![
2681            Ok(FileScanTask::builder()
2682                .with_file_size_in_bytes(
2683                    std::fs::metadata(format!("{table_location}/file_0.parquet"))
2684                        .unwrap()
2685                        .len(),
2686                )
2687                .with_start(0)
2688                .with_length(0)
2689                .with_data_file_path(format!("{table_location}/file_0.parquet"))
2690                .with_data_file_format(DataFileFormat::Parquet)
2691                .with_schema(schema.clone())
2692                .with_project_field_ids(vec![1, 2])
2693                .with_case_sensitive(false)
2694                .build()
2695                .unwrap()),
2696            Ok(FileScanTask::builder()
2697                .with_file_size_in_bytes(
2698                    std::fs::metadata(format!("{table_location}/file_1.parquet"))
2699                        .unwrap()
2700                        .len(),
2701                )
2702                .with_start(0)
2703                .with_length(0)
2704                .with_data_file_path(format!("{table_location}/file_1.parquet"))
2705                .with_data_file_format(DataFileFormat::Parquet)
2706                .with_schema(schema.clone())
2707                .with_project_field_ids(vec![1, 2])
2708                .with_case_sensitive(false)
2709                .build()
2710                .unwrap()),
2711            Ok(FileScanTask::builder()
2712                .with_file_size_in_bytes(
2713                    std::fs::metadata(format!("{table_location}/file_2.parquet"))
2714                        .unwrap()
2715                        .len(),
2716                )
2717                .with_start(0)
2718                .with_length(0)
2719                .with_data_file_path(format!("{table_location}/file_2.parquet"))
2720                .with_data_file_format(DataFileFormat::Parquet)
2721                .with_schema(schema.clone())
2722                .with_project_field_ids(vec![1, 2])
2723                .with_case_sensitive(false)
2724                .build()
2725                .unwrap()),
2726        ];
2727
2728        let tasks_stream = Box::pin(futures::stream::iter(tasks)) as FileScanTaskStream;
2729
2730        let result = reader
2731            .read(tasks_stream)
2732            .unwrap()
2733            .stream()
2734            .try_collect::<Vec<RecordBatch>>()
2735            .await
2736            .unwrap();
2737
2738        // Verify we got all 30 rows (10 from each file)
2739        let total_rows: usize = result.iter().map(|b| b.num_rows()).sum();
2740        assert_eq!(total_rows, 30, "Should have 30 total rows");
2741
2742        // Collect all ids and file_nums to verify data
2743        let mut all_ids = Vec::new();
2744        let mut all_file_nums = Vec::new();
2745
2746        for batch in &result {
2747            let id_col = batch
2748                .column(0)
2749                .as_primitive::<arrow_array::types::Int32Type>();
2750            let file_num_col = batch
2751                .column(1)
2752                .as_primitive::<arrow_array::types::Int32Type>();
2753
2754            for i in 0..batch.num_rows() {
2755                all_ids.push(id_col.value(i));
2756                all_file_nums.push(file_num_col.value(i));
2757            }
2758        }
2759
2760        assert_eq!(all_ids.len(), 30);
2761        assert_eq!(all_file_nums.len(), 30);
2762
2763        // With concurrency=1 and sequential processing, files should be processed in order
2764        // file_0: ids 0-9, file_num=0
2765        // file_1: ids 10-19, file_num=1
2766        // file_2: ids 20-29, file_num=2
2767        for i in 0..10 {
2768            assert_eq!(all_file_nums[i], 0, "First 10 rows should be from file_0");
2769            assert_eq!(all_ids[i], i as i32, "IDs should be 0-9");
2770        }
2771        for i in 10..20 {
2772            assert_eq!(all_file_nums[i], 1, "Next 10 rows should be from file_1");
2773            assert_eq!(all_ids[i], i as i32, "IDs should be 10-19");
2774        }
2775        for i in 20..30 {
2776            assert_eq!(all_file_nums[i], 2, "Last 10 rows should be from file_2");
2777            assert_eq!(all_ids[i], i as i32, "IDs should be 20-29");
2778        }
2779    }
2780
2781    #[tokio::test]
2782    async fn test_read_int96_timestamps_with_field_ids() {
2783        let schema = Arc::new(
2784            Schema::builder()
2785                .with_schema_id(1)
2786                .with_fields(vec![
2787                    NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2788                        .into(),
2789                    NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2790                ])
2791                .build()
2792                .unwrap(),
2793        );
2794
2795        let tmp_dir = TempDir::new().unwrap();
2796        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2797        let (file_path, expected_micros) =
2798            write_int96_parquet_file(&table_location, "with_ids.parquet", true);
2799
2800        assert_int96_read_matches(&file_path, schema, vec![1, 2], &expected_micros).await;
2801    }
2802
2803    #[tokio::test]
2804    async fn test_read_int96_timestamps_without_field_ids() {
2805        let schema = Arc::new(
2806            Schema::builder()
2807                .with_schema_id(1)
2808                .with_fields(vec![
2809                    NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2810                        .into(),
2811                    NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2812                ])
2813                .build()
2814                .unwrap(),
2815        );
2816
2817        let tmp_dir = TempDir::new().unwrap();
2818        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2819        let (file_path, expected_micros) =
2820            write_int96_parquet_file(&table_location, "no_ids.parquet", false);
2821
2822        assert_int96_read_matches(&file_path, schema, vec![1, 2], &expected_micros).await;
2823    }
2824
2825    #[tokio::test]
2826    async fn test_read_int96_timestamps_with_fallback_ids_and_pos() {
2827        use arrow_array::TimestampMicrosecondArray;
2828
2829        // Regression test for the combined path this refactor introduced: a field-id-less
2830        // file (positional fallback IDs) with an INT96 column and a `_pos` projection.
2831        // All three transforms -- field-ID assignment, INT96 coercion, and the row-number
2832        // virtual column -- apply in the single ArrowReaderMetadata rebuild.
2833        let schema = Arc::new(
2834            Schema::builder()
2835                .with_schema_id(1)
2836                .with_fields(vec![
2837                    NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2838                        .into(),
2839                    NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2840                ])
2841                .build()
2842                .unwrap(),
2843        );
2844
2845        let tmp_dir = TempDir::new().unwrap();
2846        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2847        let (file_path, expected_micros) =
2848            write_int96_parquet_file(&table_location, "no_ids_with_pos.parquet", false);
2849
2850        let batches =
2851            read_int96_batches(&file_path, schema, vec![1, 2, RESERVED_FIELD_ID_POS]).await;
2852
2853        assert_eq!(batches.len(), 1);
2854        // The INT96 timestamps are coerced to micros...
2855        let ts_col = batches[0]
2856            .column_by_name("ts")
2857            .expect("ts column should be present")
2858            .as_any()
2859            .downcast_ref::<TimestampMicrosecondArray>()
2860            .expect("Expected TimestampMicrosecondArray");
2861        for (i, expected) in expected_micros.iter().enumerate() {
2862            assert_eq!(ts_col.value(i), *expected, "Row {i}");
2863        }
2864        // ...and `_pos`, materialized by the RowNumber virtual column, counts rows 0,1,2.
2865        let pos_col = batches[0]
2866            .column_by_name(RESERVED_COL_NAME_POS)
2867            .expect("_pos column should be present")
2868            .as_primitive::<arrow_array::types::Int64Type>();
2869        assert_eq!(pos_col.values(), &[0, 1, 2]);
2870    }
2871
2872    #[tokio::test]
2873    async fn test_read_int96_timestamps_with_field_ids_and_pos() {
2874        use arrow_array::TimestampMicrosecondArray;
2875
2876        let schema = Arc::new(
2877            Schema::builder()
2878                .with_schema_id(1)
2879                .with_fields(vec![
2880                    NestedField::optional(1, "ts", Type::Primitive(PrimitiveType::Timestamp))
2881                        .into(),
2882                    NestedField::required(2, "id", Type::Primitive(PrimitiveType::Int)).into(),
2883                ])
2884                .build()
2885                .unwrap(),
2886        );
2887
2888        let tmp_dir = TempDir::new().unwrap();
2889        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2890        let (file_path, expected_micros) =
2891            write_int96_parquet_file(&table_location, "with_ids_with_pos.parquet", true);
2892
2893        let batches =
2894            read_int96_batches(&file_path, schema, vec![1, 2, RESERVED_FIELD_ID_POS]).await;
2895
2896        assert_eq!(batches.len(), 1);
2897        let ts_col = batches[0]
2898            .column_by_name("ts")
2899            .expect("ts column should be present")
2900            .as_any()
2901            .downcast_ref::<TimestampMicrosecondArray>()
2902            .expect("Expected TimestampMicrosecondArray");
2903
2904        for (i, expected) in expected_micros.iter().enumerate() {
2905            assert_eq!(ts_col.value(i), *expected, "Row {i}");
2906        }
2907
2908        let pos_col = batches[0]
2909            .column_by_name(RESERVED_COL_NAME_POS)
2910            .expect("_pos column should be present")
2911            .as_primitive::<arrow_array::types::Int64Type>();
2912
2913        assert_eq!(pos_col.values(), &[0, 1, 2]);
2914    }
2915
2916    #[tokio::test]
2917    async fn test_read_int96_timestamps_in_struct() {
2918        use arrow_array::{StructArray, TimestampMicrosecondArray};
2919        use parquet::basic::{Repetition, Type as PhysicalType};
2920        use parquet::data_type::Int96Type;
2921        use parquet::file::writer::SerializedFileWriter;
2922        use parquet::schema::types::Type as SchemaType;
2923
2924        let tmp_dir = TempDir::new().unwrap();
2925        let table_location = tmp_dir.path().to_str().unwrap().to_string();
2926        let file_path = format!("{table_location}/struct_int96.parquet");
2927
2928        let ts_type = SchemaType::primitive_type_builder("ts", PhysicalType::INT96)
2929            .with_repetition(Repetition::OPTIONAL)
2930            .with_id(Some(2))
2931            .build()
2932            .unwrap();
2933
2934        let struct_type = SchemaType::group_type_builder("data")
2935            .with_repetition(Repetition::REQUIRED)
2936            .with_id(Some(1))
2937            .with_fields(vec![Arc::new(ts_type)])
2938            .build()
2939            .unwrap();
2940
2941        let parquet_schema = SchemaType::group_type_builder("schema")
2942            .with_fields(vec![Arc::new(struct_type)])
2943            .build()
2944            .unwrap();
2945
2946        let (int96_val, expected_micros) = make_int96_test_value();
2947
2948        let file = File::create(&file_path).unwrap();
2949        let mut writer =
2950            SerializedFileWriter::new(file, Arc::new(parquet_schema), Default::default()).unwrap();
2951
2952        // def=1: struct is REQUIRED so no level, ts is OPTIONAL and present (1).
2953        // No repetition levels needed (no repeated groups).
2954        let mut row_group = writer.next_row_group().unwrap();
2955        {
2956            let mut col = row_group.next_column().unwrap().unwrap();
2957            col.typed::<Int96Type>()
2958                .write_batch(&[int96_val], Some(&[1]), None)
2959                .unwrap();
2960            col.close().unwrap();
2961        }
2962        row_group.close().unwrap();
2963        writer.close().unwrap();
2964
2965        let iceberg_schema = Arc::new(
2966            Schema::builder()
2967                .with_schema_id(1)
2968                .with_fields(vec![
2969                    NestedField::required(
2970                        1,
2971                        "data",
2972                        Type::Struct(crate::spec::StructType::new(vec![
2973                            NestedField::optional(
2974                                2,
2975                                "ts",
2976                                Type::Primitive(PrimitiveType::Timestamp),
2977                            )
2978                            .into(),
2979                        ])),
2980                    )
2981                    .into(),
2982                ])
2983                .build()
2984                .unwrap(),
2985        );
2986
2987        let batches = read_int96_batches(&file_path, iceberg_schema, vec![1]).await;
2988
2989        assert_eq!(batches.len(), 1);
2990        let struct_array = batches[0]
2991            .column(0)
2992            .as_any()
2993            .downcast_ref::<StructArray>()
2994            .expect("Expected StructArray");
2995        let ts_array = struct_array
2996            .column(0)
2997            .as_any()
2998            .downcast_ref::<TimestampMicrosecondArray>()
2999            .expect("Expected TimestampMicrosecondArray inside struct");
3000
3001        assert_eq!(
3002            ts_array.value(0),
3003            expected_micros,
3004            "INT96 in struct: got {}, expected {expected_micros}",
3005            ts_array.value(0)
3006        );
3007    }
3008
3009    #[tokio::test]
3010    async fn test_read_int96_timestamps_in_list() {
3011        use arrow_array::{ListArray, TimestampMicrosecondArray};
3012        use parquet::basic::{Repetition, Type as PhysicalType};
3013        use parquet::data_type::Int96Type;
3014        use parquet::file::writer::SerializedFileWriter;
3015        use parquet::schema::types::Type as SchemaType;
3016
3017        let tmp_dir = TempDir::new().unwrap();
3018        let table_location = tmp_dir.path().to_str().unwrap().to_string();
3019        let file_path = format!("{table_location}/list_int96.parquet");
3020
3021        // 3-level LIST encoding:
3022        //   optional group timestamps (LIST) {
3023        //     repeated group list {
3024        //       optional int96 element;
3025        //     }
3026        //   }
3027        let element_type = SchemaType::primitive_type_builder("element", PhysicalType::INT96)
3028            .with_repetition(Repetition::OPTIONAL)
3029            .with_id(Some(2))
3030            .build()
3031            .unwrap();
3032
3033        let list_group = SchemaType::group_type_builder("list")
3034            .with_repetition(Repetition::REPEATED)
3035            .with_fields(vec![Arc::new(element_type)])
3036            .build()
3037            .unwrap();
3038
3039        let list_type = SchemaType::group_type_builder("timestamps")
3040            .with_repetition(Repetition::OPTIONAL)
3041            .with_id(Some(1))
3042            .with_logical_type(Some(parquet::basic::LogicalType::List))
3043            .with_fields(vec![Arc::new(list_group)])
3044            .build()
3045            .unwrap();
3046
3047        let parquet_schema = SchemaType::group_type_builder("schema")
3048            .with_fields(vec![Arc::new(list_type)])
3049            .build()
3050            .unwrap();
3051
3052        let (int96_val, expected_micros) = make_int96_test_value();
3053
3054        let file = File::create(&file_path).unwrap();
3055        let mut writer =
3056            SerializedFileWriter::new(file, Arc::new(parquet_schema), Default::default()).unwrap();
3057
3058        // Write a single row with a list containing one INT96 element.
3059        // def=3: list present (1) + repeated group (2) + element present (3)
3060        // rep=0: start of a new list
3061        let mut row_group = writer.next_row_group().unwrap();
3062        {
3063            let mut col = row_group.next_column().unwrap().unwrap();
3064            col.typed::<Int96Type>()
3065                .write_batch(&[int96_val], Some(&[3]), Some(&[0]))
3066                .unwrap();
3067            col.close().unwrap();
3068        }
3069        row_group.close().unwrap();
3070        writer.close().unwrap();
3071
3072        let iceberg_schema = Arc::new(
3073            Schema::builder()
3074                .with_schema_id(1)
3075                .with_fields(vec![
3076                    NestedField::optional(
3077                        1,
3078                        "timestamps",
3079                        Type::List(crate::spec::ListType {
3080                            element_field: NestedField::optional(
3081                                2,
3082                                "element",
3083                                Type::Primitive(PrimitiveType::Timestamp),
3084                            )
3085                            .into(),
3086                        }),
3087                    )
3088                    .into(),
3089                ])
3090                .build()
3091                .unwrap(),
3092        );
3093
3094        let batches = read_int96_batches(&file_path, iceberg_schema, vec![1]).await;
3095
3096        assert_eq!(batches.len(), 1);
3097        let list_array = batches[0]
3098            .column(0)
3099            .as_any()
3100            .downcast_ref::<ListArray>()
3101            .expect("Expected ListArray");
3102        let ts_array = list_array
3103            .values()
3104            .as_any()
3105            .downcast_ref::<TimestampMicrosecondArray>()
3106            .expect("Expected TimestampMicrosecondArray inside list");
3107
3108        assert_eq!(
3109            ts_array.value(0),
3110            expected_micros,
3111            "INT96 in list: got {}, expected {expected_micros}",
3112            ts_array.value(0)
3113        );
3114    }
3115
3116    #[tokio::test]
3117    async fn test_read_int96_timestamps_in_map() {
3118        use arrow_array::{MapArray, TimestampMicrosecondArray};
3119        use parquet::basic::{Repetition, Type as PhysicalType};
3120        use parquet::data_type::{ByteArrayType, Int96Type};
3121        use parquet::file::writer::SerializedFileWriter;
3122        use parquet::schema::types::Type as SchemaType;
3123
3124        let tmp_dir = TempDir::new().unwrap();
3125        let table_location = tmp_dir.path().to_str().unwrap().to_string();
3126        let file_path = format!("{table_location}/map_int96.parquet");
3127
3128        // MAP encoding:
3129        //   optional group ts_map (MAP) {
3130        //     repeated group key_value {
3131        //       required binary key (UTF8);
3132        //       optional int96 value;
3133        //     }
3134        //   }
3135        let key_type = SchemaType::primitive_type_builder("key", PhysicalType::BYTE_ARRAY)
3136            .with_repetition(Repetition::REQUIRED)
3137            .with_logical_type(Some(parquet::basic::LogicalType::String))
3138            .with_id(Some(2))
3139            .build()
3140            .unwrap();
3141
3142        let value_type = SchemaType::primitive_type_builder("value", PhysicalType::INT96)
3143            .with_repetition(Repetition::OPTIONAL)
3144            .with_id(Some(3))
3145            .build()
3146            .unwrap();
3147
3148        let key_value_group = SchemaType::group_type_builder("key_value")
3149            .with_repetition(Repetition::REPEATED)
3150            .with_fields(vec![Arc::new(key_type), Arc::new(value_type)])
3151            .build()
3152            .unwrap();
3153
3154        let map_type = SchemaType::group_type_builder("ts_map")
3155            .with_repetition(Repetition::OPTIONAL)
3156            .with_id(Some(1))
3157            .with_logical_type(Some(parquet::basic::LogicalType::Map))
3158            .with_fields(vec![Arc::new(key_value_group)])
3159            .build()
3160            .unwrap();
3161
3162        let parquet_schema = SchemaType::group_type_builder("schema")
3163            .with_fields(vec![Arc::new(map_type)])
3164            .build()
3165            .unwrap();
3166
3167        let (int96_val, expected_micros) = make_int96_test_value();
3168
3169        let file = File::create(&file_path).unwrap();
3170        let mut writer =
3171            SerializedFileWriter::new(file, Arc::new(parquet_schema), Default::default()).unwrap();
3172
3173        // Write a single row with a map containing one key-value pair.
3174        // rep=0 for both columns: start of a new map.
3175        // key def=2: map present (1) + key_value entry present (2), key is REQUIRED.
3176        // value def=3: map present (1) + key_value entry present (2) + value present (3).
3177        let mut row_group = writer.next_row_group().unwrap();
3178        {
3179            let mut col = row_group.next_column().unwrap().unwrap();
3180            col.typed::<ByteArrayType>()
3181                .write_batch(
3182                    &[parquet::data_type::ByteArray::from("event_time")],
3183                    Some(&[2]),
3184                    Some(&[0]),
3185                )
3186                .unwrap();
3187            col.close().unwrap();
3188        }
3189        {
3190            let mut col = row_group.next_column().unwrap().unwrap();
3191            col.typed::<Int96Type>()
3192                .write_batch(&[int96_val], Some(&[3]), Some(&[0]))
3193                .unwrap();
3194            col.close().unwrap();
3195        }
3196        row_group.close().unwrap();
3197        writer.close().unwrap();
3198
3199        let iceberg_schema = Arc::new(
3200            Schema::builder()
3201                .with_schema_id(1)
3202                .with_fields(vec![
3203                    NestedField::optional(
3204                        1,
3205                        "ts_map",
3206                        Type::Map(crate::spec::MapType {
3207                            key_field: NestedField::required(
3208                                2,
3209                                "key",
3210                                Type::Primitive(PrimitiveType::String),
3211                            )
3212                            .into(),
3213                            value_field: NestedField::optional(
3214                                3,
3215                                "value",
3216                                Type::Primitive(PrimitiveType::Timestamp),
3217                            )
3218                            .into(),
3219                        }),
3220                    )
3221                    .into(),
3222                ])
3223                .build()
3224                .unwrap(),
3225        );
3226
3227        let batches = read_int96_batches(&file_path, iceberg_schema, vec![1]).await;
3228
3229        assert_eq!(batches.len(), 1);
3230        let map_array = batches[0]
3231            .column(0)
3232            .as_any()
3233            .downcast_ref::<MapArray>()
3234            .expect("Expected MapArray");
3235        let ts_array = map_array
3236            .values()
3237            .as_any()
3238            .downcast_ref::<TimestampMicrosecondArray>()
3239            .expect("Expected TimestampMicrosecondArray as map values");
3240
3241        assert_eq!(
3242            ts_array.value(0),
3243            expected_micros,
3244            "INT96 in map: got {}, expected {expected_micros}",
3245            ts_array.value(0)
3246        );
3247    }
3248
3249    /// Writes `id` (Int32) plus a wide string column (field id 2) whose bytes dominate
3250    /// the file, so that reading it is visible in `bytes_read`.
3251    ///
3252    /// `extra_fields`/`extra_columns` (e.g. a physical metadata leaf) are appended after
3253    /// the `id` and wide columns, mirroring `write_plain_parquet`'s shape.
3254    fn write_parquet_with_wide_column(
3255        dir: &str,
3256        name: &str,
3257        extra_fields: Vec<Field>,
3258        extra_columns: Vec<ArrayRef>,
3259    ) -> String {
3260        let wide_field =
3261            Field::new("wide", DataType::Utf8, false).with_metadata(HashMap::from([(
3262                PARQUET_FIELD_ID_META_KEY.to_string(),
3263                "2".to_string(),
3264            )]));
3265        // Varied bytes so the column chunk does not compress away under SNAPPY, keeping
3266        // the `bytes_read` difference between projecting it and not unambiguous.
3267        let wide_values: Vec<String> = (0..3)
3268            .map(|i| {
3269                (0..2048)
3270                    .map(|j| ((i * 2048 + j) % 251) as u8 as char)
3271                    .collect()
3272            })
3273            .collect();
3274
3275        let mut fields = vec![wide_field];
3276        fields.extend(extra_fields);
3277        let mut columns: Vec<ArrayRef> = vec![Arc::new(StringArray::from(wide_values))];
3278        columns.extend(extra_columns);
3279        write_plain_parquet(dir, name, fields, columns)
3280    }
3281
3282    /// Schema with `id` (field 1, Int) and `wide` (field 2, String), matching
3283    /// `write_parquet_with_wide_column`.
3284    fn id_and_wide_schema() -> SchemaRef {
3285        Arc::new(
3286            Schema::builder()
3287                .with_schema_id(1)
3288                .with_fields(vec![
3289                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
3290                    NestedField::required(2, "wide", Type::Primitive(PrimitiveType::String)).into(),
3291                ])
3292                .build()
3293                .unwrap(),
3294        )
3295    }
3296
3297    /// Builds a scan task over `file_path` projecting `project_field_ids`.
3298    fn metadata_projection_task(
3299        file_path: String,
3300        schema: SchemaRef,
3301        project_field_ids: Vec<i32>,
3302    ) -> FileScanTask {
3303        metadata_projection_task_with_first_row_id(file_path, schema, project_field_ids, None)
3304    }
3305
3306    fn metadata_projection_task_with_first_row_id(
3307        file_path: String,
3308        schema: SchemaRef,
3309        project_field_ids: Vec<i32>,
3310        first_row_id: Option<i64>,
3311    ) -> FileScanTask {
3312        FileScanTask::builder()
3313            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3314            .with_start(0)
3315            .with_length(0)
3316            .with_data_file_path(file_path)
3317            .with_data_file_format(DataFileFormat::Parquet)
3318            .with_schema(schema)
3319            .with_project_field_ids(project_field_ids)
3320            .with_first_row_id(first_row_id)
3321            .with_case_sensitive(false)
3322            .build()
3323            .unwrap()
3324    }
3325
3326    /// Runs a single-task scan and returns the batches plus the bytes read from storage.
3327    async fn scan_task(task: FileScanTask) -> (Vec<RecordBatch>, u64) {
3328        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current()).build();
3329        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
3330        let scan = reader.read(tasks).unwrap();
3331        let metrics = scan.metrics().clone();
3332        let batches = scan.stream().try_collect().await.unwrap();
3333        (batches, metrics.bytes_read())
3334    }
3335
3336    #[tokio::test]
3337    async fn test_pos_only_projection_reads_no_data_columns() {
3338        let tmp_dir = TempDir::new().unwrap();
3339        let dir = tmp_dir.path().to_str().unwrap();
3340
3341        let pos_only = metadata_projection_task(
3342            write_parquet_with_wide_column(dir, "pos_only.parquet", vec![], vec![]),
3343            id_and_wide_schema(),
3344            vec![RESERVED_FIELD_ID_POS],
3345        );
3346        let (batches, pos_only_bytes) = scan_task(pos_only).await;
3347
3348        // Only `_pos` is materialized -- no data columns.
3349        assert_eq!(batches[0].num_columns(), 1);
3350        let pos_col = batches[0]
3351            .column_by_name(RESERVED_COL_NAME_POS)
3352            .expect("_pos column should be present")
3353            .as_primitive::<arrow_array::types::Int64Type>();
3354        assert_eq!(pos_col.values(), &[0, 1, 2]);
3355
3356        // A scan of the same-shaped file that also projects the wide data column must read
3357        // materially more, proving the wide column chunk was not fetched above.
3358        let with_data = metadata_projection_task(
3359            write_parquet_with_wide_column(dir, "pos_only_ref.parquet", vec![], vec![]),
3360            id_and_wide_schema(),
3361            vec![2, RESERVED_FIELD_ID_POS],
3362        );
3363        let (_, with_data_bytes) = scan_task(with_data).await;
3364
3365        assert!(
3366            pos_only_bytes < with_data_bytes,
3367            "_pos-only scan should read fewer bytes than a scan of the wide column: \
3368             {pos_only_bytes} vs {with_data_bytes}"
3369        );
3370    }
3371
3372    #[tokio::test]
3373    async fn test_pos_only_projection_keeps_absolute_pos_under_predicate() {
3374        use crate::expr::{Bind, Reference};
3375        use crate::spec::Datum;
3376
3377        let tmp_dir = TempDir::new().unwrap();
3378        let dir = tmp_dir.path().to_str().unwrap();
3379        // id = [1, 2, 3]; drop the middle physical row via a predicate + row selection.
3380        let file_path = write_plain_parquet(dir, "pos_only_predicate.parquet", vec![], vec![]);
3381
3382        let schema = Arc::new(
3383            Schema::builder()
3384                .with_schema_id(1)
3385                .with_fields(vec![
3386                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
3387                ])
3388                .build()
3389                .unwrap(),
3390        );
3391        let bound = Reference::new("id")
3392            .not_equal_to(Datum::int(2))
3393            .bind(Arc::clone(&schema), false)
3394            .unwrap();
3395        let task = FileScanTask::builder()
3396            .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3397            .with_start(0)
3398            .with_length(0)
3399            .with_data_file_path(file_path)
3400            .with_data_file_format(DataFileFormat::Parquet)
3401            .with_schema(schema)
3402            .with_project_field_ids(vec![RESERVED_FIELD_ID_POS])
3403            .with_predicate(Some(bound))
3404            .with_case_sensitive(false)
3405            .build()
3406            .unwrap();
3407
3408        // Row selection must be enabled for the predicate to filter rows. The row filter
3409        // reads `id` for its own evaluation even though `id` is not projected; the surviving
3410        // rows must keep their ABSOLUTE positions (0 and 2), not renumbered (0 and 1).
3411        let reader = ArrowReaderBuilder::new(FileIO::new_with_fs(), Runtime::current())
3412            .with_row_selection_enabled(true)
3413            .build();
3414        let tasks = Box::pin(futures::stream::iter(vec![Ok(task)])) as FileScanTaskStream;
3415        let batches: Vec<RecordBatch> = reader
3416            .read(tasks)
3417            .unwrap()
3418            .stream()
3419            .try_collect()
3420            .await
3421            .unwrap();
3422
3423        let pos: Vec<i64> = batches
3424            .iter()
3425            .flat_map(|b| {
3426                b.column_by_name(RESERVED_COL_NAME_POS)
3427                    .expect("_pos column should be present")
3428                    .as_primitive::<arrow_array::types::Int64Type>()
3429                    .values()
3430                    .to_vec()
3431            })
3432            .collect();
3433        assert_eq!(pos, vec![0, 2]);
3434    }
3435
3436    #[tokio::test]
3437    async fn test_pos_and_file_projection() {
3438        use crate::metadata_columns::RESERVED_COL_NAME_FILE;
3439
3440        let tmp_dir = TempDir::new().unwrap();
3441        let dir = tmp_dir.path().to_str().unwrap();
3442        // The motivating row-lineage shape: a synthesized position column (mask -> none)
3443        // alongside a materialized per-file constant.
3444        let file_path = write_parquet_with_wide_column(dir, "pos_and_file.parquet", vec![], vec![]);
3445        let task = metadata_projection_task(file_path.clone(), id_and_wide_schema(), vec![
3446            RESERVED_FIELD_ID_POS,
3447            RESERVED_FIELD_ID_FILE,
3448        ]);
3449        let (batches, _) = scan_task(task).await;
3450
3451        // Both metadata columns materialize; no data column is read.
3452        assert_eq!(batches[0].num_columns(), 2);
3453        let pos_col = batches[0]
3454            .column_by_name(RESERVED_COL_NAME_POS)
3455            .expect("_pos column should be present")
3456            .as_primitive::<arrow_array::types::Int64Type>();
3457        assert_eq!(pos_col.values(), &[0, 1, 2]);
3458        let file_col = batches[0]
3459            .column_by_name(RESERVED_COL_NAME_FILE)
3460            .expect("_file column should be present");
3461        let file_col = cast(file_col, &DataType::Utf8).unwrap();
3462        let file_col = file_col.as_any().downcast_ref::<StringArray>().unwrap();
3463        assert_eq!(file_col.value(0), file_path);
3464    }
3465
3466    #[tokio::test]
3467    async fn test_pos_and_physical_seq_projection_reads_only_the_leaf() {
3468        use crate::metadata_columns::{
3469            RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
3470            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3471        };
3472
3473        // A v3 rewrite that carried rows forward stores `_last_updated_sequence_number`
3474        // per-row. Projecting only `_pos` + the sequence column must read just that one
3475        // physical leaf, not every data column.
3476        let tmp_dir = TempDir::new().unwrap();
3477        let dir = tmp_dir.path().to_str().unwrap();
3478
3479        // File: id (1), wide data column (2), physical _last_updated_sequence_number.
3480        let write = |name: &str| {
3481            write_parquet_with_wide_column(
3482                dir,
3483                name,
3484                vec![physical_last_updated_seq_field()],
3485                vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
3486            )
3487        };
3488
3489        let seq_task = |path: String, ids: Vec<i32>| {
3490            FileScanTask::builder()
3491                .with_file_size_in_bytes(std::fs::metadata(&path).unwrap().len())
3492                .with_start(0)
3493                .with_length(0)
3494                .with_data_file_path(path)
3495                .with_data_file_format(DataFileFormat::Parquet)
3496                .with_schema(id_and_wide_schema())
3497                .with_project_field_ids(ids)
3498                .with_first_row_id(Some(100))
3499                .with_data_sequence_number(Some(9))
3500                .with_case_sensitive(false)
3501                .build()
3502                .unwrap()
3503        };
3504
3505        let meta_only = seq_task(write("pos_seq.parquet"), vec![
3506            RESERVED_FIELD_ID_POS,
3507            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3508        ]);
3509        let (batches, meta_only_bytes) = scan_task(meta_only).await;
3510
3511        // `_pos` and the coalesced sequence column materialize; the wide column does not.
3512        let pos_col = batches[0]
3513            .column_by_name(RESERVED_COL_NAME_POS)
3514            .expect("_pos column should be present")
3515            .as_primitive::<arrow_array::types::Int64Type>();
3516        assert_eq!(pos_col.values(), &[0, 1, 2]);
3517        let seq_col = batches[0]
3518            .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
3519            .expect("_last_updated_sequence_number column should be present");
3520        let seq_col = cast(seq_col, &DataType::Int64).unwrap();
3521        let seq_col = seq_col.as_any().downcast_ref::<Int64Array>().unwrap();
3522        // Per-row stored value where non-null, else the data sequence number (9).
3523        assert_eq!(seq_col.value(0), 5);
3524        assert_eq!(seq_col.value(1), 9);
3525        assert_eq!(seq_col.value(2), 8);
3526        assert!(batches[0].column_by_name("wide").is_none());
3527
3528        // A scan that also projects the wide data column must read materially more,
3529        // proving the metadata-only scan pruned to just the sequence leaf.
3530        let with_data = seq_task(write("pos_seq_ref.parquet"), vec![
3531            2,
3532            RESERVED_FIELD_ID_POS,
3533            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3534        ]);
3535        let (_, with_data_bytes) = scan_task(with_data).await;
3536
3537        assert!(
3538            meta_only_bytes < with_data_bytes,
3539            "_pos + physical sequence scan should read fewer bytes than one that also \
3540             reads the wide column: {meta_only_bytes} vs {with_data_bytes}"
3541        );
3542    }
3543
3544    #[tokio::test]
3545    async fn test_seq_only_projection_reads_only_the_leaf() {
3546        use crate::metadata_columns::{
3547            RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
3548            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3549        };
3550
3551        // `_last_updated_sequence_number` alone (no `_pos`, no data column). This is a
3552        // metadata-only projection, so RowNumber drives the `none()` downgrade and the
3553        // physical leaf is unioned back in, pruning the read to just that leaf.
3554        let tmp_dir = TempDir::new().unwrap();
3555        let dir = tmp_dir.path().to_str().unwrap();
3556        let write = |name: &str| {
3557            write_parquet_with_wide_column(
3558                dir,
3559                name,
3560                vec![physical_last_updated_seq_field()],
3561                vec![Arc::new(Int64Array::from(vec![Some(5), None, Some(8)])) as ArrayRef],
3562            )
3563        };
3564        let seq_task = |path: String, ids: Vec<i32>| {
3565            FileScanTask::builder()
3566                .with_file_size_in_bytes(std::fs::metadata(&path).unwrap().len())
3567                .with_start(0)
3568                .with_length(0)
3569                .with_data_file_path(path)
3570                .with_data_file_format(DataFileFormat::Parquet)
3571                .with_schema(id_and_wide_schema())
3572                .with_project_field_ids(ids)
3573                .with_first_row_id(Some(100))
3574                .with_data_sequence_number(Some(9))
3575                .with_case_sensitive(false)
3576                .build()
3577                .unwrap()
3578        };
3579
3580        let meta_only = seq_task(write("seq_only.parquet"), vec![
3581            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3582        ]);
3583        let (batches, meta_only_bytes) = scan_task(meta_only).await;
3584
3585        let seq_col = batches[0]
3586            .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
3587            .expect("_last_updated_sequence_number column should be present");
3588        let seq_col = cast(seq_col, &DataType::Int64).unwrap();
3589        let seq_col = seq_col.as_any().downcast_ref::<Int64Array>().unwrap();
3590        assert_eq!(seq_col.value(0), 5);
3591        assert_eq!(seq_col.value(1), 9);
3592        assert_eq!(seq_col.value(2), 8);
3593        let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3594        assert_eq!(total_rows, 3);
3595        assert!(batches[0].column_by_name("wide").is_none());
3596
3597        let with_data = seq_task(write("seq_only_ref.parquet"), vec![
3598            2,
3599            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3600        ]);
3601        let (_, with_data_bytes) = scan_task(with_data).await;
3602
3603        assert!(
3604            meta_only_bytes < with_data_bytes,
3605            "seq-only scan should read fewer bytes than one that also reads the wide \
3606             column: {meta_only_bytes} vs {with_data_bytes}"
3607        );
3608    }
3609
3610    #[tokio::test]
3611    async fn test_seq_only_null_first_row_id_reads_no_data_columns() {
3612        use crate::metadata_columns::{
3613            RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
3614            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3615        };
3616
3617        // Seq-only projection with a null first_row_id: the column is nulled and the
3618        // physical leaf is NOT read (the gated `coalesce_last_updated_seq_leaf` is None).
3619        // This is a metadata-only projection, so RowNumber supplies the row count and the
3620        // data columns are pruned; the row count is still 3 and the values all null.
3621        let tmp_dir = TempDir::new().unwrap();
3622        let dir = tmp_dir.path().to_str().unwrap();
3623        let seq_task = |name: &str, ids: Vec<i32>| {
3624            let file_path = write_parquet_with_wide_column(
3625                dir,
3626                name,
3627                vec![physical_last_updated_seq_field()],
3628                vec![Arc::new(Int64Array::from(vec![Some(5), Some(6), Some(7)])) as ArrayRef],
3629            );
3630            FileScanTask::builder()
3631                .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3632                .with_start(0)
3633                .with_length(0)
3634                .with_data_file_path(file_path)
3635                .with_data_file_format(DataFileFormat::Parquet)
3636                .with_schema(id_and_wide_schema())
3637                .with_project_field_ids(ids)
3638                .with_first_row_id(None)
3639                .with_data_sequence_number(Some(9))
3640                .with_case_sensitive(false)
3641                .build()
3642                .unwrap()
3643        };
3644
3645        let meta_only = seq_task("seq_only_null_first.parquet", vec![
3646            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3647        ]);
3648        let (batches, meta_only_bytes) = scan_task(meta_only).await;
3649
3650        assert_eq!(batches[0].num_columns(), 1);
3651        let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3652        assert_eq!(total_rows, 3);
3653        let seq_col = batches[0]
3654            .column_by_name(RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER)
3655            .expect("_last_updated_sequence_number column should be present");
3656        let seq_col = cast(seq_col, &DataType::Int64).unwrap();
3657        let seq_col = seq_col.as_any().downcast_ref::<Int64Array>().unwrap();
3658        assert!((0..3).all(|i| seq_col.is_null(i)));
3659
3660        let with_data = seq_task("seq_only_null_first_ref.parquet", vec![
3661            2,
3662            RESERVED_FIELD_ID_LAST_UPDATED_SEQUENCE_NUMBER,
3663        ]);
3664        let (_, with_data_bytes) = scan_task(with_data).await;
3665
3666        assert!(
3667            meta_only_bytes < with_data_bytes,
3668            "seq-only scan with a null first_row_id should read fewer bytes than a scan of \
3669             the wide column: {meta_only_bytes} vs {with_data_bytes}"
3670        );
3671    }
3672
3673    #[tokio::test]
3674    async fn test_file_only_reads_no_data_columns() {
3675        use crate::metadata_columns::RESERVED_COL_NAME_FILE;
3676
3677        // `SELECT _file` materializes a per-file string constant with no physical column to
3678        // read. The metadata-only RowNumber counter sizes it, so the scan prunes the data
3679        // columns rather than reading them all just to recover the row count.
3680        let tmp_dir = TempDir::new().unwrap();
3681        let dir = tmp_dir.path().to_str().unwrap();
3682
3683        let file_path = write_parquet_with_wide_column(dir, "file_only.parquet", vec![], vec![]);
3684        let meta_only = metadata_projection_task(file_path.clone(), id_and_wide_schema(), vec![
3685            RESERVED_FIELD_ID_FILE,
3686        ]);
3687        let (batches, meta_only_bytes) = scan_task(meta_only).await;
3688
3689        assert_eq!(batches[0].num_columns(), 1);
3690        let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3691        assert_eq!(total_rows, 3);
3692        let file_col = batches[0]
3693            .column_by_name(RESERVED_COL_NAME_FILE)
3694            .expect("_file column should be present");
3695        let file_col = cast(file_col, &DataType::Utf8).unwrap();
3696        let file_col = file_col.as_any().downcast_ref::<StringArray>().unwrap();
3697        assert_eq!(file_col.value(0), file_path);
3698
3699        let with_data = metadata_projection_task(
3700            write_parquet_with_wide_column(dir, "file_only_ref.parquet", vec![], vec![]),
3701            id_and_wide_schema(),
3702            vec![2, RESERVED_FIELD_ID_FILE],
3703        );
3704        let (_, with_data_bytes) = scan_task(with_data).await;
3705
3706        assert!(
3707            meta_only_bytes < with_data_bytes,
3708            "_file-only scan should read fewer bytes than a scan of the wide column: \
3709             {meta_only_bytes} vs {with_data_bytes}"
3710        );
3711    }
3712
3713    #[tokio::test]
3714    async fn test_spec_id_only_reads_no_data_columns() {
3715        use crate::metadata_columns::{RESERVED_COL_NAME_SPEC_ID, RESERVED_FIELD_ID_SPEC_ID};
3716        use crate::spec::{Literal, PartitionSpec, Struct, Transform};
3717
3718        // `SELECT _spec_id` materializes an int constant from the task's partition spec id,
3719        // with no physical column to read. The RowNumber counter sizes it, so the scan
3720        // prunes the data columns rather than reading them all to recover the row count.
3721        let tmp_dir = TempDir::new().unwrap();
3722        let dir = tmp_dir.path().to_str().unwrap();
3723        let schema = id_and_wide_schema();
3724        let spec = Arc::new(
3725            PartitionSpec::builder(schema.clone())
3726                .with_spec_id(7)
3727                .add_partition_field("id", "id", Transform::Identity)
3728                .unwrap()
3729                .build()
3730                .unwrap(),
3731        );
3732
3733        let spec_id_task = |file_path: String, project_field_ids: Vec<i32>| {
3734            FileScanTask::builder()
3735                .with_file_size_in_bytes(std::fs::metadata(&file_path).unwrap().len())
3736                .with_start(0)
3737                .with_length(0)
3738                .with_data_file_path(file_path)
3739                .with_data_file_format(DataFileFormat::Parquet)
3740                .with_schema(schema.clone())
3741                .with_project_field_ids(project_field_ids)
3742                .with_partition(Some(Struct::from_iter([Some(Literal::int(42))])))
3743                .with_partition_spec(Some(spec.clone()))
3744                .with_case_sensitive(false)
3745                .build()
3746                .unwrap()
3747        };
3748
3749        let meta_only = spec_id_task(
3750            write_parquet_with_wide_column(dir, "spec_id_only.parquet", vec![], vec![]),
3751            vec![RESERVED_FIELD_ID_SPEC_ID],
3752        );
3753        let (batches, meta_only_bytes) = scan_task(meta_only).await;
3754
3755        assert_eq!(batches[0].num_columns(), 1);
3756        let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3757        assert_eq!(total_rows, 3);
3758        let spec_id_col = batches[0]
3759            .column_by_name(RESERVED_COL_NAME_SPEC_ID)
3760            .expect("_spec_id column should be present");
3761        let spec_id_col = cast(spec_id_col, &DataType::Int32).unwrap();
3762        let spec_id_col = spec_id_col.as_any().downcast_ref::<Int32Array>().unwrap();
3763        assert_eq!(spec_id_col.values(), &[7, 7, 7]);
3764
3765        let with_data = spec_id_task(
3766            write_parquet_with_wide_column(dir, "spec_id_only_ref.parquet", vec![], vec![]),
3767            vec![2, RESERVED_FIELD_ID_SPEC_ID],
3768        );
3769        let (_, with_data_bytes) = scan_task(with_data).await;
3770
3771        assert!(
3772            meta_only_bytes < with_data_bytes,
3773            "_spec_id-only scan should read fewer bytes than a scan of the wide column: \
3774             {meta_only_bytes} vs {with_data_bytes}"
3775        );
3776    }
3777
3778    #[tokio::test]
3779    async fn test_empty_projection_preserves_row_count() {
3780        let tmp_dir = TempDir::new().unwrap();
3781        let dir = tmp_dir.path().to_str().unwrap();
3782        let file_path = write_plain_parquet(dir, "empty_projection.parquet", vec![], vec![]);
3783        let schema = Arc::new(
3784            Schema::builder()
3785                .with_schema_id(1)
3786                .with_fields(vec![
3787                    NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
3788                ])
3789                .build()
3790                .unwrap(),
3791        );
3792        let task = metadata_projection_task(file_path, schema, vec![]);
3793        let (batches, _) = scan_task(task).await;
3794
3795        // A bare COUNT(*)-style empty projection must still report the row count.
3796        let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
3797        assert_eq!(total_rows, 3);
3798    }
3799}