Skip to main content

iceberg/scan/
mod.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//! Table scan api.
19
20mod cache;
21use cache::*;
22mod context;
23use context::*;
24mod task;
25
26use std::collections::HashMap;
27use std::future::Future;
28use std::sync::Arc;
29
30use arrow_array::RecordBatch;
31use futures::channel::mpsc::{Sender, channel};
32use futures::stream::BoxStream;
33use futures::{SinkExt, StreamExt, TryStreamExt};
34pub use task::*;
35
36use crate::arrow::ArrowReaderBuilder;
37pub use crate::arrow::{ScanMetrics, ScanResult};
38use crate::delete_file_index::DeleteFileIndex;
39use crate::error::invalid_data;
40use crate::expr::visitors::inclusive_metrics_evaluator::InclusiveMetricsEvaluator;
41use crate::expr::{Bind, BoundPredicate, Predicate};
42use crate::io::FileIO;
43use crate::metadata_columns::{
44    RESERVED_FIELD_ID_PARTITION, get_metadata_field_id, is_metadata_column_name,
45};
46use crate::runtime::Runtime;
47use crate::spec::{DataContentType, Schema, SchemaRef, SnapshotRef, SortOrderRef, StructType};
48use crate::table::Table;
49use crate::util::available_parallelism;
50use crate::{Error, ErrorKind, Result};
51
52/// A stream of arrow [`RecordBatch`]es.
53pub type ArrowRecordBatchStream = BoxStream<'static, Result<RecordBatch>>;
54
55/// Resolves a column name to its field ID, honouring the scan's case sensitivity.
56fn resolve_field_id(schema: &Schema, column_name: &str, case_sensitive: bool) -> Option<i32> {
57    if case_sensitive {
58        schema.field_id_by_name(column_name)
59    } else {
60        schema
61            .field_by_name_case_insensitive(column_name)
62            .map(|field| field.id)
63    }
64}
65
66fn collect_scan_field_ids(
67    schema: &Schema,
68    column_names: Option<&[String]>,
69    case_sensitive: bool,
70) -> Result<Vec<i32>> {
71    let Some(column_names) = column_names else {
72        return Ok(schema.as_struct().fields().iter().map(|f| f.id).collect());
73    };
74
75    column_names
76        .iter()
77        .map(|column_name| {
78            if is_metadata_column_name(column_name) {
79                return get_metadata_field_id(column_name);
80            }
81
82            let field_id = resolve_field_id(schema, column_name, case_sensitive).ok_or_else(|| {
83                invalid_data!("Column {column_name} not found in table. Schema: {schema}")
84            })?;
85
86            schema
87                .as_struct()
88                .field_by_id(field_id)
89                .ok_or_else(|| {
90                    Error::new(
91                        ErrorKind::FeatureUnsupported,
92                        format!(
93                            "Column {column_name} is not a direct child of schema but a nested field, which is not supported now. Schema: {schema}"
94                        ),
95                    )
96                })?;
97
98            Ok(field_id)
99        })
100        .collect()
101}
102
103fn bind_scan_predicate(
104    schema: &SchemaRef,
105    predicate: Option<&Predicate>,
106    case_sensitive: bool,
107) -> Result<Option<Arc<BoundPredicate>>> {
108    predicate
109        .map(|predicate| predicate.bind(schema.clone(), case_sensitive))
110        .transpose()
111        .map(|predicate| predicate.map(Arc::new))
112}
113
114fn projected_partition_type(
115    table: &Table,
116    schema: &Schema,
117    field_ids: &[i32],
118) -> Result<Option<Arc<StructType>>> {
119    if !field_ids.contains(&RESERVED_FIELD_ID_PARTITION) {
120        return Ok(None);
121    }
122
123    table
124        .metadata()
125        .unified_partition_type(schema)
126        .map(Arc::new)
127        .map(Some)
128}
129
130/// Builder to create table scan.
131pub struct TableScanBuilder<'a> {
132    table: &'a Table,
133    // Defaults to none which means select all columns
134    column_names: Option<Vec<String>>,
135    snapshot_id: Option<i64>,
136    batch_size: Option<usize>,
137    case_sensitive: bool,
138    filter: Option<Predicate>,
139    concurrency_limit_data_files: usize,
140    concurrency_limit_manifest_entries: usize,
141    concurrency_limit_manifest_files: usize,
142    row_group_filtering_enabled: bool,
143    row_selection_enabled: bool,
144    bloom_filter_enabled: bool,
145}
146
147impl<'a> TableScanBuilder<'a> {
148    pub(crate) fn new(table: &'a Table) -> Self {
149        let num_cpus = available_parallelism().get();
150
151        Self {
152            table,
153            column_names: None,
154            snapshot_id: None,
155            batch_size: None,
156            case_sensitive: true,
157            filter: None,
158            concurrency_limit_data_files: num_cpus,
159            concurrency_limit_manifest_entries: num_cpus,
160            concurrency_limit_manifest_files: num_cpus,
161            row_group_filtering_enabled: true,
162            row_selection_enabled: false,
163            bloom_filter_enabled: false,
164        }
165    }
166
167    /// Sets the desired size of batches in the response
168    /// to something other than the default
169    pub fn with_batch_size(mut self, batch_size: Option<usize>) -> Self {
170        self.batch_size = batch_size;
171        self
172    }
173
174    /// Sets the scan's case sensitivity
175    pub fn with_case_sensitive(mut self, case_sensitive: bool) -> Self {
176        self.case_sensitive = case_sensitive;
177        self
178    }
179
180    /// Specifies a predicate to use as a filter
181    pub fn with_filter(mut self, predicate: Predicate) -> Self {
182        // calls rewrite_not to remove Not nodes, which must be absent
183        // when applying the manifest evaluator
184        self.filter = Some(predicate.rewrite_not());
185        self
186    }
187
188    /// Select all columns.
189    pub fn select_all(mut self) -> Self {
190        self.column_names = None;
191        self
192    }
193
194    /// Select empty columns.
195    pub fn select_empty(mut self) -> Self {
196        self.column_names = Some(vec![]);
197        self
198    }
199
200    /// Select some columns of the table.
201    pub fn select(mut self, column_names: impl IntoIterator<Item = impl ToString>) -> Self {
202        self.column_names = Some(
203            column_names
204                .into_iter()
205                .map(|item| item.to_string())
206                .collect(),
207        );
208        self
209    }
210
211    /// Set the snapshot to scan. When not set, it uses current snapshot.
212    pub fn snapshot_id(mut self, snapshot_id: i64) -> Self {
213        self.snapshot_id = Some(snapshot_id);
214        self
215    }
216
217    /// Sets the concurrency limit for both manifest files and manifest
218    /// entries for this scan
219    pub fn with_concurrency_limit(mut self, limit: usize) -> Self {
220        self.concurrency_limit_manifest_files = limit;
221        self.concurrency_limit_manifest_entries = limit;
222        self.concurrency_limit_data_files = limit;
223        self
224    }
225
226    /// Sets the data file concurrency limit for this scan
227    pub fn with_data_file_concurrency_limit(mut self, limit: usize) -> Self {
228        self.concurrency_limit_data_files = limit;
229        self
230    }
231
232    /// Sets the manifest entry concurrency limit for this scan
233    pub fn with_manifest_entry_concurrency_limit(mut self, limit: usize) -> Self {
234        self.concurrency_limit_manifest_entries = limit;
235        self
236    }
237
238    /// Determines whether to enable row group filtering.
239    /// When enabled, if a read is performed with a filter predicate,
240    /// then the metadata for each row group in the parquet file is
241    /// evaluated against the filter predicate and row groups
242    /// that cant contain matching rows will be skipped entirely.
243    ///
244    /// Defaults to enabled, as it generally improves performance or
245    /// keeps it the same, with performance degradation unlikely.
246    pub fn with_row_group_filtering_enabled(mut self, row_group_filtering_enabled: bool) -> Self {
247        self.row_group_filtering_enabled = row_group_filtering_enabled;
248        self
249    }
250
251    /// Determines whether to enable row selection.
252    /// When enabled, if a read is performed with a filter predicate,
253    /// then (for row groups that have not been skipped) the page index
254    /// for each row group in a parquet file is parsed and evaluated
255    /// against the filter predicate to determine if ranges of rows
256    /// within a row group can be skipped, based upon the page-level
257    /// statistics for each column.
258    ///
259    /// Defaults to being disabled. Enabling requires parsing the parquet page
260    /// index, which can be slow enough that parsing the page index outweighs any
261    /// gains from the reduced number of rows that need scanning.
262    /// It is recommended to experiment with partitioning, sorting, row group size,
263    /// page size, and page row limit Iceberg settings on the table being scanned in
264    /// order to get the best performance from using row selection.
265    pub fn with_row_selection_enabled(mut self, row_selection_enabled: bool) -> Self {
266        self.row_selection_enabled = row_selection_enabled;
267        self
268    }
269
270    /// Determines whether to enable bloom filter-based row group filtering.
271    ///
272    /// When enabled, if a read is performed with an equality or IN predicate,
273    /// the bloom filter for relevant columns in each row group is read and
274    /// checked. Row groups where the bloom filter proves the value is absent
275    /// are skipped entirely.
276    ///
277    /// Defaults to disabled, as reading bloom filters requires additional I/O
278    /// per column per row group.
279    pub fn with_bloom_filter_enabled(mut self, bloom_filter_enabled: bool) -> Self {
280        self.bloom_filter_enabled = bloom_filter_enabled;
281        self
282    }
283
284    /// Build the table scan.
285    pub fn build(self) -> Result<TableScan> {
286        let snapshot = match self.snapshot_id {
287            Some(snapshot_id) => self
288                .table
289                .metadata()
290                .snapshot_by_id(snapshot_id)
291                .ok_or_else(|| invalid_data!("Snapshot with id {snapshot_id} not found"))?
292                .clone(),
293            None => {
294                let Some(current_snapshot_id) = self.table.metadata().current_snapshot() else {
295                    return Ok(TableScan {
296                        batch_size: self.batch_size,
297                        column_names: self.column_names,
298                        file_io: self.table.file_io().clone(),
299                        plan_context: None,
300                        concurrency_limit_data_files: self.concurrency_limit_data_files,
301                        concurrency_limit_manifest_entries: self.concurrency_limit_manifest_entries,
302                        concurrency_limit_manifest_files: self.concurrency_limit_manifest_files,
303                        row_group_filtering_enabled: self.row_group_filtering_enabled,
304                        row_selection_enabled: self.row_selection_enabled,
305                        bloom_filter_enabled: self.bloom_filter_enabled,
306                        runtime: self.table.runtime().clone(),
307                    });
308                };
309                current_snapshot_id.clone()
310            }
311        };
312
313        let schema = snapshot.schema(self.table.metadata())?;
314        let field_ids =
315            collect_scan_field_ids(&schema, self.column_names.as_deref(), self.case_sensitive)?;
316        let snapshot_bound_predicate =
317            bind_scan_predicate(&schema, self.filter.as_ref(), self.case_sensitive)?;
318        let name_mapping = self
319            .table
320            .metadata()
321            .table_properties()
322            .default_name_mapping()?
323            .map(Arc::new);
324        let unified_partition_type = projected_partition_type(self.table, &schema, &field_ids)?;
325
326        // Precompute the table's sort orders once, keyed by id, so each manifest-file
327        // context carries only this narrow map instead of the full table metadata.
328        let sort_orders = Arc::new(
329            self.table
330                .metadata()
331                .sort_orders_iter()
332                .map(|order| (order.order_id, order.clone()))
333                .collect::<HashMap<i64, SortOrderRef>>(),
334        );
335
336        let plan_context = PlanContext {
337            snapshot,
338            table_metadata: self.table.metadata_ref(),
339            snapshot_schema: schema,
340            case_sensitive: self.case_sensitive,
341            predicate: self.filter.map(Arc::new),
342            snapshot_bound_predicate,
343            object_cache: self.table.object_cache(),
344            field_ids: Arc::new(field_ids),
345            name_mapping,
346            partition_filter_cache: Arc::new(PartitionFilterCache::new()),
347            manifest_evaluator_cache: Arc::new(ManifestEvaluatorCache::new()),
348            expression_evaluator_cache: Arc::new(ExpressionEvaluatorCache::new()),
349            unified_partition_type,
350            sort_orders,
351        };
352
353        Ok(TableScan {
354            batch_size: self.batch_size,
355            column_names: self.column_names,
356            file_io: self.table.file_io().clone(),
357            plan_context: Some(plan_context),
358            concurrency_limit_data_files: self.concurrency_limit_data_files,
359            concurrency_limit_manifest_entries: self.concurrency_limit_manifest_entries,
360            concurrency_limit_manifest_files: self.concurrency_limit_manifest_files,
361            row_group_filtering_enabled: self.row_group_filtering_enabled,
362            row_selection_enabled: self.row_selection_enabled,
363            bloom_filter_enabled: self.bloom_filter_enabled,
364            runtime: self.table.runtime().clone(),
365        })
366    }
367}
368
369/// Table scan.
370#[derive(Debug)]
371pub struct TableScan {
372    /// A [PlanContext], if this table has at least one snapshot, otherwise None.
373    ///
374    /// If this is None, then the scan contains no rows.
375    plan_context: Option<PlanContext>,
376    batch_size: Option<usize>,
377    file_io: FileIO,
378    column_names: Option<Vec<String>>,
379    /// The maximum number of manifest files that will be
380    /// retrieved from [`FileIO`] concurrently
381    concurrency_limit_manifest_files: usize,
382
383    /// The maximum number of [`ManifestEntry`]s that will
384    /// be processed in parallel
385    concurrency_limit_manifest_entries: usize,
386
387    /// The maximum number of [`ManifestEntry`]s that will
388    /// be processed in parallel
389    concurrency_limit_data_files: usize,
390
391    row_group_filtering_enabled: bool,
392    row_selection_enabled: bool,
393    bloom_filter_enabled: bool,
394
395    runtime: Runtime,
396}
397
398impl TableScan {
399    /// Returns a stream of [`FileScanTask`]s.
400    pub async fn plan_files(&self) -> Result<FileScanTaskStream> {
401        self.plan_data_files(|ctx| ctx.into_file_scan_task()).await
402    }
403
404    pub(crate) async fn plan_cow_rewrite_files(
405        &self,
406    ) -> Result<BoxStream<'static, Result<crate::cow_rewrite::CowRewriteFile>>> {
407        self.plan_data_files(|ctx| ctx.into_cow_rewrite_file())
408            .await
409    }
410
411    async fn plan_data_files<T, F, Fut>(
412        &self,
413        build_result: F,
414    ) -> Result<BoxStream<'static, Result<T>>>
415    where
416        T: Send + 'static,
417        F: Fn(ManifestEntryContext) -> Fut + Copy + Send + Sync + 'static,
418        Fut: Future<Output = Result<T>> + Send + 'static,
419    {
420        let Some(plan_context) = self.plan_context.as_ref() else {
421            return Ok(Box::pin(futures::stream::empty()));
422        };
423
424        let concurrency_limit_manifest_files = self.concurrency_limit_manifest_files;
425        let concurrency_limit_manifest_entries = self.concurrency_limit_manifest_entries;
426
427        // used to stream ManifestEntryContexts between stages of the file plan operation
428        let (manifest_entry_data_ctx_tx, manifest_entry_data_ctx_rx) =
429            channel(concurrency_limit_manifest_files);
430        let (manifest_entry_delete_ctx_tx, manifest_entry_delete_ctx_rx) =
431            channel(concurrency_limit_manifest_files);
432
433        // used to stream the planned data file results back to the caller
434        let (planned_file_tx, planned_file_rx) = channel(concurrency_limit_manifest_entries);
435
436        let (delete_file_idx, delete_file_tx) = DeleteFileIndex::new(self.runtime.clone());
437
438        let manifest_list = plan_context.get_manifest_list().await?;
439
440        // get the [`ManifestFile`]s from the [`ManifestList`], filtering out any
441        // whose partitions cannot match this
442        // scan's filter
443        let manifest_file_contexts = plan_context.build_manifest_file_contexts(
444            manifest_list,
445            manifest_entry_data_ctx_tx,
446            delete_file_idx.clone(),
447            manifest_entry_delete_ctx_tx,
448        )?;
449
450        let mut channel_for_manifest_error = planned_file_tx.clone();
451        let mut channel_for_data_manifest_entry_error = planned_file_tx.clone();
452        let mut channel_for_delete_manifest_entry_error = planned_file_tx.clone();
453
454        let rt = self.runtime.clone();
455
456        // Concurrently load all [`Manifest`]s and stream their [`ManifestEntry`]s
457        rt.io().spawn(async move {
458            let result = futures::stream::iter(manifest_file_contexts)
459                .try_for_each_concurrent(concurrency_limit_manifest_files, |ctx| async move {
460                    ctx.fetch_manifest_and_stream_manifest_entries().await
461                })
462                .await;
463
464            if let Err(error) = result {
465                let _ = channel_for_manifest_error.send(Err(error)).await;
466            }
467        });
468
469        // Process the delete file [`ManifestEntry`] stream in parallel
470        {
471            let rt = rt.clone();
472            let rt_inner = rt.clone();
473            rt.cpu().spawn(async move {
474                let result = manifest_entry_delete_ctx_rx
475                    .map(|me_ctx| Ok((me_ctx, delete_file_tx.clone())))
476                    .try_for_each_concurrent(
477                        concurrency_limit_manifest_entries,
478                        |(manifest_entry_context, tx)| {
479                            let rt_inner = rt_inner.clone();
480                            async move {
481                                rt_inner
482                                    .cpu()
483                                    .spawn(async move {
484                                        Self::process_delete_manifest_entry(
485                                            manifest_entry_context,
486                                            tx,
487                                        )
488                                        .await
489                                    })
490                                    .await?
491                            }
492                        },
493                    )
494                    .await;
495
496                if let Err(error) = result {
497                    let _ = channel_for_delete_manifest_entry_error
498                        .send(Err(error))
499                        .await;
500                }
501            });
502        }
503
504        // Process the data file [`ManifestEntry`] stream in parallel
505        {
506            let rt_inner = rt.clone();
507            rt.cpu().spawn(async move {
508                let result = manifest_entry_data_ctx_rx
509                    .map(|me_ctx| Ok((me_ctx, planned_file_tx.clone())))
510                    .try_for_each_concurrent(
511                        concurrency_limit_manifest_entries,
512                        |(manifest_entry_context, tx)| {
513                            let rt_inner = rt_inner.clone();
514                            async move {
515                                rt_inner
516                                    .cpu()
517                                    .spawn(async move {
518                                        Self::process_data_manifest_entry(
519                                            manifest_entry_context,
520                                            tx,
521                                            build_result,
522                                        )
523                                        .await
524                                    })
525                                    .await?
526                            }
527                        },
528                    )
529                    .await;
530
531                if let Err(error) = result {
532                    let _ = channel_for_data_manifest_entry_error.send(Err(error)).await;
533                }
534            });
535        }
536
537        Ok(planned_file_rx.boxed())
538    }
539
540    /// Returns an [`ArrowRecordBatchStream`].
541    pub async fn to_arrow(&self) -> Result<ArrowRecordBatchStream> {
542        let mut arrow_reader_builder =
543            ArrowReaderBuilder::new(self.file_io.clone(), self.runtime.clone())
544                .with_data_file_concurrency_limit(self.concurrency_limit_data_files)
545                .with_row_group_filtering_enabled(self.row_group_filtering_enabled)
546                .with_row_selection_enabled(self.row_selection_enabled)
547                .with_bloom_filter_enabled(self.bloom_filter_enabled);
548
549        if let Some(batch_size) = self.batch_size {
550            arrow_reader_builder = arrow_reader_builder.with_batch_size(batch_size);
551        }
552
553        arrow_reader_builder
554            .build()
555            .read(self.plan_files().await?)
556            .map(|result| result.stream())
557    }
558
559    /// Returns a reference to the column names of the table scan.
560    pub fn column_names(&self) -> Option<&[String]> {
561        self.column_names.as_deref()
562    }
563
564    /// Returns a reference to the snapshot of the table scan.
565    pub fn snapshot(&self) -> Option<&SnapshotRef> {
566        self.plan_context.as_ref().map(|x| &x.snapshot)
567    }
568
569    async fn process_data_manifest_entry<T, F, Fut>(
570        manifest_entry_context: ManifestEntryContext,
571        mut planned_file_tx: Sender<Result<T>>,
572        build_result: F,
573    ) -> Result<()>
574    where
575        F: Fn(ManifestEntryContext) -> Fut + Send + Sync + 'static,
576        Fut: Future<Output = Result<T>> + Send + 'static,
577    {
578        // skip processing this manifest entry if it has been marked as deleted
579        if !manifest_entry_context.manifest_entry.is_alive() {
580            return Ok(());
581        }
582
583        // abort the plan if we encounter a manifest entry for a delete file
584        if manifest_entry_context.manifest_entry.content_type() != DataContentType::Data {
585            return Err(Error::new(
586                ErrorKind::FeatureUnsupported,
587                "Encountered an entry for a delete file in a data file manifest",
588            ));
589        }
590
591        if let Some(ref bound_predicates) = manifest_entry_context.bound_predicates {
592            let BoundPredicates {
593                snapshot_bound_predicate,
594                partition_bound_predicate,
595            } = bound_predicates.as_ref();
596
597            let expression_evaluator_cache =
598                manifest_entry_context.expression_evaluator_cache.as_ref();
599
600            let expression_evaluator = expression_evaluator_cache.get(
601                manifest_entry_context.partition_spec_id,
602                partition_bound_predicate,
603            )?;
604
605            // skip any data file whose partition data indicates that it can't contain
606            // any data that matches this scan's filter
607            if !expression_evaluator.eval(manifest_entry_context.manifest_entry.data_file())? {
608                return Ok(());
609            }
610
611            // skip any data file whose metrics don't match this scan's filter
612            if !InclusiveMetricsEvaluator::eval(
613                snapshot_bound_predicate,
614                manifest_entry_context.manifest_entry.data_file(),
615                false,
616            )? {
617                return Ok(());
618            }
619        }
620
621        // congratulations! the manifest entry has made its way through the
622        // entire plan without getting filtered out. Build the planned file before
623        // sending so delete-file lookup preserves the original scan timing.
624        planned_file_tx
625            .send(Ok(build_result(manifest_entry_context).await?))
626            .await?;
627
628        Ok(())
629    }
630
631    async fn process_delete_manifest_entry(
632        manifest_entry_context: ManifestEntryContext,
633        mut delete_file_ctx_tx: Sender<DeleteFileContext>,
634    ) -> Result<()> {
635        // skip processing this manifest entry if it has been marked as deleted
636        if !manifest_entry_context.manifest_entry.is_alive() {
637            return Ok(());
638        }
639
640        // abort the plan if we encounter a manifest entry that is not for a delete file
641        if manifest_entry_context.manifest_entry.content_type() == DataContentType::Data {
642            return Err(Error::new(
643                ErrorKind::FeatureUnsupported,
644                "Encountered an entry for a data file in a delete manifest",
645            ));
646        }
647
648        if let Some(ref bound_predicates) = manifest_entry_context.bound_predicates {
649            let expression_evaluator_cache =
650                manifest_entry_context.expression_evaluator_cache.as_ref();
651
652            let expression_evaluator = expression_evaluator_cache.get(
653                manifest_entry_context.partition_spec_id,
654                &bound_predicates.partition_bound_predicate,
655            )?;
656
657            // skip any data file whose partition data indicates that it can't contain
658            // any data that matches this scan's filter
659            if !expression_evaluator.eval(manifest_entry_context.manifest_entry.data_file())? {
660                return Ok(());
661            }
662        }
663
664        delete_file_ctx_tx
665            .send(DeleteFileContext {
666                manifest_entry: manifest_entry_context.manifest_entry.clone(),
667                partition_spec_id: manifest_entry_context.partition_spec_id,
668            })
669            .await?;
670
671        Ok(())
672    }
673}
674
675pub(crate) struct BoundPredicates {
676    partition_bound_predicate: BoundPredicate,
677    snapshot_bound_predicate: BoundPredicate,
678}
679
680#[cfg(test)]
681mod tests {
682    //! shared tests for the table scan API
683
684    use std::collections::HashMap;
685    use std::sync::Arc;
686
687    use arrow_array::cast::AsArray;
688    use arrow_array::types::Int32Type;
689    use arrow_array::{
690        Array, ArrayRef, BooleanArray, Float64Array, Int32Array, Int64Array, RunArray, StringArray,
691    };
692    use futures::{TryStreamExt, stream};
693    use uuid::Uuid;
694
695    use crate::arrow::ArrowReaderBuilder;
696    use crate::expr::{BoundPredicate, Reference};
697    use crate::io::FileIO;
698    use crate::metadata_columns::{
699        RESERVED_COL_NAME_FILE, RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER,
700        RESERVED_COL_NAME_POS, RESERVED_COL_NAME_SPEC_ID, RESERVED_FIELD_ID_POS,
701    };
702    use crate::scan::{FileScanTask, FileScanTaskDeleteFile};
703    use crate::spec::{
704        DataContentType, DataFileFormat, Datum, Literal, MAIN_BRANCH, MappedField, NameMapping,
705        NestedField, NullOrder, Operation, PartitionSpec, PrimitiveType, Schema, Snapshot,
706        SortDirection, SortField, SortOrder, Struct, Summary, TableMetadataBuilder,
707        TableProperties, Transform, Type, UnboundPartitionSpec,
708    };
709    use crate::table::Table;
710    use crate::test_utils::scan::{TableTestFixture, assert_last_updated_seq_all};
711    use crate::test_utils::test_runtime;
712    use crate::{ErrorKind, TableIdent};
713
714    #[tokio::test]
715    async fn test_table_scan_columns() {
716        let table = TableTestFixture::new().table;
717
718        let table_scan = table.scan().select(["x", "y"]).build().unwrap();
719        assert_eq!(
720            Some(vec!["x".to_string(), "y".to_string()]),
721            table_scan.column_names
722        );
723
724        let table_scan = table
725            .scan()
726            .select(["x", "y"])
727            .select(["z"])
728            .build()
729            .unwrap();
730        assert_eq!(Some(vec!["z".to_string()]), table_scan.column_names);
731    }
732
733    #[tokio::test]
734    async fn test_select_all() {
735        let table = TableTestFixture::new().table;
736
737        let table_scan = table.scan().select_all().build().unwrap();
738        assert!(table_scan.column_names.is_none());
739    }
740
741    #[test]
742    fn test_select_no_exist_column() {
743        let table = TableTestFixture::new().table;
744
745        let table_scan = table.scan().select(["x", "y", "z", "a", "b"]).build();
746        assert!(table_scan.is_err());
747    }
748
749    #[test]
750    fn test_case_sensitive_scan_rejects_mismatched_column_case() {
751        let table = TableTestFixture::new().table;
752
753        // Case sensitivity defaults to true, so "X" must not resolve to "x".
754        assert!(table.scan().select(["X"]).build().is_err());
755        assert!(
756            table
757                .scan()
758                .with_filter(Reference::new("X").greater_than(Datum::long(1)))
759                .build()
760                .is_err()
761        );
762    }
763
764    #[tokio::test]
765    async fn test_case_insensitive_scan_resolves_mismatched_column_case() {
766        let mut fixture = TableTestFixture::new();
767        fixture.setup_manifest_files().await;
768
769        // The schema declares lowercase "x" and "z"; a case-insensitive scan must
770        // resolve the upper-cased names in both the projection and the filter.
771        let table_scan = fixture
772            .table
773            .scan()
774            .with_case_sensitive(false)
775            .select(["X", "Z"])
776            .with_filter(Reference::new("Y").greater_than(Datum::long(1)))
777            .build()
778            .unwrap();
779
780        let batches: Vec<_> = table_scan
781            .to_arrow()
782            .await
783            .unwrap()
784            .try_collect()
785            .await
786            .unwrap();
787
788        assert_eq!(batches[0].num_columns(), 2);
789        assert_eq!(
790            batches[0]
791                .column_by_name("x")
792                .unwrap()
793                .as_any()
794                .downcast_ref::<Int64Array>()
795                .unwrap()
796                .value(0),
797            1
798        );
799        assert_eq!(
800            batches[0]
801                .column_by_name("z")
802                .unwrap()
803                .as_any()
804                .downcast_ref::<Int64Array>()
805                .unwrap()
806                .value(0),
807            3
808        );
809    }
810
811    #[tokio::test]
812    async fn test_table_scan_default_snapshot_id() {
813        let table = TableTestFixture::new().table;
814
815        let table_scan = table.scan().build().unwrap();
816        assert_eq!(
817            table.metadata().current_snapshot().unwrap().snapshot_id(),
818            table_scan.snapshot().unwrap().snapshot_id()
819        );
820    }
821
822    #[test]
823    fn test_table_scan_non_exist_snapshot_id() {
824        let table = TableTestFixture::new().table;
825
826        let table_scan = table.scan().snapshot_id(1024).build();
827        assert!(table_scan.is_err());
828    }
829
830    #[tokio::test]
831    async fn test_table_scan_with_snapshot_id() {
832        let table = TableTestFixture::new().table;
833
834        let table_scan = table
835            .scan()
836            .snapshot_id(3051729675574597004)
837            .with_row_selection_enabled(true)
838            .build()
839            .unwrap();
840        assert_eq!(
841            table_scan.snapshot().unwrap().snapshot_id(),
842            3051729675574597004
843        );
844    }
845
846    fn table_with_property(key: &str, value: &str) -> Table {
847        let fixture = TableTestFixture::new();
848        let mut metadata = fixture.table.metadata().clone();
849        metadata
850            .properties
851            .insert(key.to_string(), value.to_string());
852        Table::builder()
853            .metadata(metadata)
854            .identifier(fixture.table.identifier().clone())
855            .file_io(fixture.table.file_io().clone())
856            .metadata_location(fixture.table.metadata_location().unwrap().to_string())
857            .runtime(test_runtime())
858            .build()
859            .unwrap()
860    }
861
862    #[test]
863    fn test_table_scan_without_name_mapping_property() {
864        let table = TableTestFixture::new().table;
865
866        let table_scan = table.scan().build().unwrap();
867        assert!(
868            table_scan
869                .plan_context
870                .as_ref()
871                .unwrap()
872                .name_mapping
873                .is_none()
874        );
875    }
876
877    #[test]
878    fn test_table_scan_with_name_mapping_property() {
879        let mapping_json = r#"[{"field-id":1,"names":["id","record_id"]}]"#;
880        let table =
881            table_with_property(TableProperties::PROPERTY_DEFAULT_NAME_MAPPING, mapping_json);
882
883        let table_scan = table.scan().build().unwrap();
884        let mapping = table_scan
885            .plan_context
886            .as_ref()
887            .unwrap()
888            .name_mapping
889            .as_ref()
890            .expect("name_mapping should be parsed from the table property");
891        let fields = mapping.fields();
892        assert_eq!(fields.len(), 1);
893        assert_eq!(fields[0].field_id(), Some(1));
894        assert_eq!(fields[0].names(), &[
895            "id".to_string(),
896            "record_id".to_string()
897        ]);
898    }
899
900    #[test]
901    fn test_table_scan_with_malformed_name_mapping_property() {
902        let table = table_with_property(
903            TableProperties::PROPERTY_DEFAULT_NAME_MAPPING,
904            "{ not valid json",
905        );
906
907        let err = table
908            .scan()
909            .build()
910            .expect_err("malformed name mapping should fail to parse");
911        assert_eq!(err.kind(), ErrorKind::DataInvalid);
912    }
913
914    #[tokio::test]
915    async fn test_plan_files_carries_name_mapping_into_file_scan_task() {
916        let mut fixture = TableTestFixture::new();
917        fixture.setup_manifest_files().await;
918
919        let mapping_json = r#"[{"field-id":1,"names":["id","record_id"]}]"#;
920        let mut metadata = fixture.table.metadata().clone();
921        metadata.properties.insert(
922            TableProperties::PROPERTY_DEFAULT_NAME_MAPPING.to_string(),
923            mapping_json.to_string(),
924        );
925        let table = Table::builder()
926            .metadata(metadata)
927            .identifier(fixture.table.identifier().clone())
928            .file_io(fixture.table.file_io().clone())
929            .metadata_location(fixture.table.metadata_location().unwrap().to_string())
930            .runtime(test_runtime())
931            .build()
932            .unwrap();
933
934        let tasks: Vec<_> = table
935            .scan()
936            .build()
937            .unwrap()
938            .plan_files()
939            .await
940            .unwrap()
941            .try_collect()
942            .await
943            .unwrap();
944
945        assert!(!tasks.is_empty(), "expected at least one FileScanTask");
946        for task in &tasks {
947            let mapping = task
948                .name_mapping()
949                .expect("name_mapping should reach the FileScanTask");
950            assert_eq!(mapping.fields().len(), 1);
951            assert_eq!(mapping.fields()[0].field_id(), Some(1));
952        }
953    }
954
955    #[tokio::test]
956    async fn test_plan_files_carries_sort_order_into_file_scan_task() {
957        let mut fixture = TableTestFixture::new();
958
959        // Inject the reserved unsorted order (id 0) inline rather than editing the shared
960        // testdata fixture, so the id-0 file below exercises the `!is_unsorted()` filter
961        // branch instead of the `and_then` short-circuit an absent entry would take.
962        let mut metadata = fixture.table.metadata().clone();
963        metadata
964            .sort_orders
965            .insert(0, Arc::new(SortOrder::unsorted_order()));
966        fixture.table = fixture.table.with_metadata(Arc::new(metadata));
967
968        let expected_sort_order = fixture
969            .table
970            .metadata()
971            .sort_order_by_id(3)
972            .unwrap()
973            .clone();
974
975        // sort_order_ids: resolvable (3), absent, unresolvable (99), reserved unsorted (0).
976        fixture
977            .setup_manifest_files_with_sort_order_ids([Some(3), None, Some(99), Some(0)])
978            .await;
979
980        let tasks: Vec<_> = fixture
981            .table
982            .scan()
983            .build()
984            .unwrap()
985            .plan_files()
986            .await
987            .unwrap()
988            .try_collect()
989            .await
990            .unwrap();
991
992        assert_eq!(tasks.len(), 4, "expected all four FileScanTasks");
993
994        // Aggregates catch a systemic regression (every entry resolving to id 3, or
995        // resolution dropping entirely) that the per-file checks below would each still pass.
996        assert_eq!(
997            tasks.iter().filter(|t| t.sort_order().is_some()).count(),
998            1,
999            "exactly one file resolves to a sort order"
1000        );
1001        assert_eq!(
1002            tasks.iter().filter(|t| t.sort_order_id().is_some()).count(),
1003            3,
1004            "three files carry a raw sort_order_id"
1005        );
1006
1007        let resolved = tasks
1008            .iter()
1009            .find(|t| t.data_file_path().ends_with("1.parquet"))
1010            .unwrap();
1011        assert_eq!(resolved.sort_order_id(), Some(3));
1012        assert_eq!(
1013            resolved.sort_order(),
1014            Some(&expected_sort_order),
1015            "sort_order_id 3 should resolve to the table's sort order at id 3"
1016        );
1017
1018        let missing = tasks
1019            .iter()
1020            .find(|t| t.data_file_path().ends_with("2.parquet"))
1021            .unwrap();
1022        assert_eq!(missing.sort_order_id(), None);
1023        assert!(
1024            missing.sort_order().is_none(),
1025            "a file with no sort_order_id carries no sort_order"
1026        );
1027
1028        let unresolvable = tasks
1029            .iter()
1030            .find(|t| t.data_file_path().ends_with("3.parquet"))
1031            .unwrap();
1032        assert_eq!(
1033            unresolvable.sort_order_id(),
1034            Some(99),
1035            "the raw id is preserved even when it does not resolve"
1036        );
1037        assert!(
1038            unresolvable.sort_order().is_none(),
1039            "an unresolvable sort_order_id resolves to no sort_order"
1040        );
1041
1042        let unsorted = tasks
1043            .iter()
1044            .find(|t| t.data_file_path().ends_with("4.parquet"))
1045            .unwrap();
1046        assert_eq!(unsorted.sort_order_id(), Some(0));
1047        assert!(
1048            unsorted.sort_order().is_none(),
1049            "the reserved unsorted order (id 0) resolves to no sort_order"
1050        );
1051    }
1052
1053    #[tokio::test]
1054    async fn test_plan_files_on_table_without_any_snapshots() {
1055        let table = TableTestFixture::new_empty().table;
1056        let batch_stream = table.scan().build().unwrap().to_arrow().await.unwrap();
1057        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1058        assert!(batches.is_empty());
1059    }
1060
1061    #[tokio::test]
1062    async fn test_plan_files_no_deletions() {
1063        let mut fixture = TableTestFixture::new();
1064        fixture.setup_manifest_files().await;
1065
1066        // Create table scan for current snapshot and plan files
1067        let table_scan = fixture
1068            .table
1069            .scan()
1070            .with_row_selection_enabled(true)
1071            .build()
1072            .unwrap();
1073
1074        let mut tasks = table_scan
1075            .plan_files()
1076            .await
1077            .unwrap()
1078            .try_fold(vec![], |mut acc, task| async move {
1079                acc.push(task);
1080                Ok(acc)
1081            })
1082            .await
1083            .unwrap();
1084
1085        assert_eq!(tasks.len(), 2);
1086
1087        tasks.sort_by_key(|t| t.data_file_path().to_string());
1088
1089        // Check first task is added data file
1090        assert_eq!(
1091            tasks[0].data_file_path(),
1092            format!("{}/1.parquet", &fixture.table_location)
1093        );
1094
1095        // Check second task is existing data file
1096        assert_eq!(
1097            tasks[1].data_file_path(),
1098            format!("{}/3.parquet", &fixture.table_location)
1099        );
1100    }
1101
1102    #[tokio::test]
1103    async fn test_plan_files_carries_row_lineage_into_file_scan_task() {
1104        let mut fixture = TableTestFixture::new();
1105        fixture.setup_manifest_files().await;
1106
1107        let mut tasks: Vec<_> = fixture
1108            .table
1109            .scan()
1110            .build()
1111            .unwrap()
1112            .plan_files()
1113            .await
1114            .unwrap()
1115            .try_collect()
1116            .await
1117            .unwrap();
1118
1119        assert_eq!(tasks.len(), 2);
1120        tasks.sort_by_key(|task| task.data_file_path().to_string());
1121
1122        // The added file inherits the current snapshot's data sequence number,
1123        // the existing file keeps the one it was written with.
1124        assert_eq!(
1125            tasks[0].data_file_path(),
1126            format!("{}/1.parquet", &fixture.table_location)
1127        );
1128        assert_eq!(tasks[0].data_sequence_number(), Some(1));
1129        assert_eq!(
1130            tasks[1].data_file_path(),
1131            format!("{}/3.parquet", &fixture.table_location)
1132        );
1133        assert_eq!(tasks[1].data_sequence_number(), Some(0));
1134
1135        // first_row_id is a v3 concept; a v2 manifest carries none.
1136        assert!(tasks.iter().all(|task| task.first_row_id().is_none()));
1137    }
1138
1139    #[tokio::test]
1140    async fn test_plan_files_carries_row_lineage_from_v3_manifest() {
1141        let mut fixture = TableTestFixture::new();
1142        fixture.setup_v3_manifest_files().await;
1143
1144        let task = fixture
1145            .table
1146            .scan()
1147            .build()
1148            .unwrap()
1149            .plan_files()
1150            .await
1151            .unwrap()
1152            .try_collect::<Vec<_>>()
1153            .await
1154            .unwrap()
1155            .into_iter()
1156            .next()
1157            .expect("expected one FileScanTask");
1158
1159        // The manifest-level first_row_id (42) is inherited onto the entry on
1160        // read, then carried onto the task.
1161        assert_eq!(task.first_row_id(), Some(42));
1162        // The data sequence number is threaded through the same v3 read path.
1163        assert_eq!(task.data_sequence_number(), Some(1));
1164    }
1165
1166    #[tokio::test]
1167    async fn test_filtered_scan_with_dropped_partition_source_column() {
1168        let mut fixture = TableTestFixture::new();
1169        fixture.setup_manifest_files().await;
1170
1171        // Baseline: the same filtered scan against the table before evolution.
1172        let baseline = scan_y_gte_5(&fixture.table).await;
1173        assert!(!baseline.is_empty());
1174        assert!(baseline.iter().all(|y| *y >= 5));
1175
1176        // Evolve the table so that the manifests reference a historical spec whose source
1177        // column is no longer in the current schema: make an unpartitioned spec the
1178        // default, then drop the original spec's source column from the schema.
1179        let current_schema = fixture.table.metadata().current_schema();
1180        let evolved_schema = Schema::builder()
1181            .with_fields(
1182                current_schema
1183                    .as_struct()
1184                    .fields()
1185                    .iter()
1186                    .filter(|field| field.id != 1)
1187                    .cloned(),
1188            )
1189            .with_identifier_field_ids(vec![2])
1190            .build()
1191            .unwrap();
1192        let evolved =
1193            TableMetadataBuilder::new_from_metadata(fixture.table.metadata().clone(), None)
1194                .add_default_partition_spec(UnboundPartitionSpec::builder().build())
1195                .unwrap()
1196                .add_current_schema(evolved_schema)
1197                .unwrap()
1198                .build()
1199                .unwrap()
1200                .metadata;
1201
1202        // a commit after the evolution carries the previous manifests forward: the new
1203        // snapshot uses the evolved schema while its manifests still use historical spec 0
1204        let parent = evolved.current_snapshot().unwrap().clone();
1205        let snapshot = Snapshot::builder()
1206            .with_snapshot_id(parent.snapshot_id() + 1)
1207            .with_parent_snapshot_id(Some(parent.snapshot_id()))
1208            .with_sequence_number(evolved.last_sequence_number() + 1)
1209            .with_timestamp_ms(evolved.last_updated_ms + 1)
1210            .with_schema_id(evolved.current_schema_id())
1211            .with_manifest_list(parent.manifest_list())
1212            .with_summary(Summary {
1213                operation: Operation::Append,
1214                additional_properties: HashMap::new(),
1215            })
1216            .build();
1217        let metadata = TableMetadataBuilder::new_from_metadata(evolved, None)
1218            .set_branch_snapshot(snapshot, MAIN_BRANCH)
1219            .unwrap()
1220            .build()
1221            .unwrap()
1222            .metadata;
1223        let table = fixture.table.clone().with_metadata(Arc::new(metadata));
1224
1225        // Planning and reading must succeed, and the results must match the table before
1226        // evolution: no rows wrongly pruned and none returned unfiltered.
1227        let evolved = scan_y_gte_5(&table).await;
1228        assert_eq!(evolved, baseline);
1229    }
1230
1231    async fn scan_y_gte_5(table: &Table) -> Vec<i64> {
1232        let table_scan = table
1233            .scan()
1234            .select(["y"])
1235            .with_filter(Reference::new("y").greater_than_or_equal_to(Datum::long(5)))
1236            .build()
1237            .unwrap();
1238        let batches: Vec<_> = table_scan
1239            .to_arrow()
1240            .await
1241            .unwrap()
1242            .try_collect()
1243            .await
1244            .unwrap();
1245
1246        let mut values: Vec<i64> = batches
1247            .iter()
1248            .flat_map(|batch| {
1249                let col = batch.column_by_name("y").unwrap();
1250                let arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1251                (0..arr.len()).map(|i| arr.value(i)).collect::<Vec<_>>()
1252            })
1253            .collect();
1254        values.sort_unstable();
1255        values
1256    }
1257
1258    #[tokio::test]
1259    async fn test_open_parquet_no_deletions() {
1260        let mut fixture = TableTestFixture::new();
1261        fixture.setup_manifest_files().await;
1262
1263        // Create table scan for current snapshot and plan files
1264        let table_scan = fixture
1265            .table
1266            .scan()
1267            .with_row_selection_enabled(true)
1268            .build()
1269            .unwrap();
1270
1271        let batch_stream = table_scan.to_arrow().await.unwrap();
1272
1273        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1274
1275        let col = batches[0].column_by_name("x").unwrap();
1276
1277        let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1278        assert_eq!(int64_arr.value(0), 1);
1279    }
1280
1281    #[tokio::test]
1282    async fn test_open_parquet_no_deletions_by_separate_reader() {
1283        let mut fixture = TableTestFixture::new();
1284        fixture.setup_manifest_files().await;
1285
1286        // Create table scan for current snapshot and plan files
1287        let table_scan = fixture
1288            .table
1289            .scan()
1290            .with_row_selection_enabled(true)
1291            .build()
1292            .unwrap();
1293
1294        let mut plan_task: Vec<_> = table_scan
1295            .plan_files()
1296            .await
1297            .unwrap()
1298            .try_collect()
1299            .await
1300            .unwrap();
1301        assert_eq!(plan_task.len(), 2);
1302
1303        let reader = ArrowReaderBuilder::new(
1304            fixture.table.file_io().clone(),
1305            fixture.table.runtime().clone(),
1306        )
1307        .build();
1308        let batch_stream = reader
1309            .clone()
1310            .read(Box::pin(stream::iter(vec![Ok(plan_task.remove(0))])))
1311            .unwrap()
1312            .stream();
1313        let batch_1: Vec<_> = batch_stream.try_collect().await.unwrap();
1314
1315        let reader = ArrowReaderBuilder::new(
1316            fixture.table.file_io().clone(),
1317            fixture.table.runtime().clone(),
1318        )
1319        .build();
1320        let batch_stream = reader
1321            .read(Box::pin(stream::iter(vec![Ok(plan_task.remove(0))])))
1322            .unwrap()
1323            .stream();
1324        let batch_2: Vec<_> = batch_stream.try_collect().await.unwrap();
1325
1326        assert_eq!(batch_1, batch_2);
1327    }
1328
1329    #[tokio::test]
1330    async fn test_open_parquet_with_projection() {
1331        let mut fixture = TableTestFixture::new();
1332        fixture.setup_manifest_files().await;
1333
1334        // Create table scan for current snapshot and plan files
1335        let table_scan = fixture
1336            .table
1337            .scan()
1338            .select(["x", "z"])
1339            .with_row_selection_enabled(true)
1340            .build()
1341            .unwrap();
1342
1343        let batch_stream = table_scan.to_arrow().await.unwrap();
1344
1345        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1346
1347        assert_eq!(batches[0].num_columns(), 2);
1348
1349        let col1 = batches[0].column_by_name("x").unwrap();
1350        let int64_arr = col1.as_any().downcast_ref::<Int64Array>().unwrap();
1351        assert_eq!(int64_arr.value(0), 1);
1352
1353        let col2 = batches[0].column_by_name("z").unwrap();
1354        let int64_arr = col2.as_any().downcast_ref::<Int64Array>().unwrap();
1355        assert_eq!(int64_arr.value(0), 3);
1356
1357        // test empty scan
1358        let table_scan = fixture.table.scan().select_empty().build().unwrap();
1359        let batch_stream = table_scan.to_arrow().await.unwrap();
1360        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1361
1362        assert_eq!(batches[0].num_columns(), 0);
1363        assert_eq!(batches[0].num_rows(), 1024);
1364    }
1365
1366    #[tokio::test]
1367    async fn test_filter_on_arrow_lt() {
1368        let mut fixture = TableTestFixture::new();
1369        fixture.setup_manifest_files().await;
1370
1371        // Filter: y < 3
1372        let mut builder = fixture.table.scan();
1373        let predicate = Reference::new("y").less_than(Datum::long(3));
1374        builder = builder
1375            .with_filter(predicate)
1376            .with_row_selection_enabled(true);
1377        let table_scan = builder.build().unwrap();
1378
1379        let batch_stream = table_scan.to_arrow().await.unwrap();
1380
1381        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1382
1383        assert_eq!(batches[0].num_rows(), 512);
1384
1385        let col = batches[0].column_by_name("x").unwrap();
1386        let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1387        assert_eq!(int64_arr.value(0), 1);
1388
1389        let col = batches[0].column_by_name("y").unwrap();
1390        let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1391        assert_eq!(int64_arr.value(0), 2);
1392    }
1393
1394    #[tokio::test]
1395    async fn test_filter_on_arrow_gt_eq() {
1396        let mut fixture = TableTestFixture::new();
1397        fixture.setup_manifest_files().await;
1398
1399        // Filter: y >= 5
1400        let mut builder = fixture.table.scan();
1401        let predicate = Reference::new("y").greater_than_or_equal_to(Datum::long(5));
1402        builder = builder
1403            .with_filter(predicate)
1404            .with_row_selection_enabled(true);
1405        let table_scan = builder.build().unwrap();
1406
1407        let batch_stream = table_scan.to_arrow().await.unwrap();
1408
1409        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1410
1411        assert_eq!(batches[0].num_rows(), 12);
1412
1413        let col = batches[0].column_by_name("x").unwrap();
1414        let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1415        assert_eq!(int64_arr.value(0), 1);
1416
1417        let col = batches[0].column_by_name("y").unwrap();
1418        let int64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1419        assert_eq!(int64_arr.value(0), 5);
1420    }
1421
1422    #[tokio::test]
1423    async fn test_filter_double_eq() {
1424        let mut fixture = TableTestFixture::new();
1425        fixture.setup_manifest_files().await;
1426
1427        // Filter: dbl == 150.0
1428        let mut builder = fixture.table.scan();
1429        let predicate = Reference::new("dbl").equal_to(Datum::double(150.0f64));
1430        builder = builder
1431            .with_filter(predicate)
1432            .with_row_selection_enabled(true);
1433        let table_scan = builder.build().unwrap();
1434
1435        let batch_stream = table_scan.to_arrow().await.unwrap();
1436
1437        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1438
1439        assert_eq!(batches.len(), 2);
1440        assert_eq!(batches[0].num_rows(), 12);
1441
1442        let col = batches[0].column_by_name("dbl").unwrap();
1443        let f64_arr = col.as_any().downcast_ref::<Float64Array>().unwrap();
1444        assert_eq!(f64_arr.value(1), 150.0f64);
1445    }
1446
1447    #[tokio::test]
1448    async fn test_filter_int_eq() {
1449        let mut fixture = TableTestFixture::new();
1450        fixture.setup_manifest_files().await;
1451
1452        // Filter: i32 == 150
1453        let mut builder = fixture.table.scan();
1454        let predicate = Reference::new("i32").equal_to(Datum::int(150i32));
1455        builder = builder
1456            .with_filter(predicate)
1457            .with_row_selection_enabled(true);
1458        let table_scan = builder.build().unwrap();
1459
1460        let batch_stream = table_scan.to_arrow().await.unwrap();
1461
1462        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1463
1464        assert_eq!(batches.len(), 2);
1465        assert_eq!(batches[0].num_rows(), 12);
1466
1467        let col = batches[0].column_by_name("i32").unwrap();
1468        let i32_arr = col.as_any().downcast_ref::<Int32Array>().unwrap();
1469        assert_eq!(i32_arr.value(1), 150i32);
1470    }
1471
1472    #[tokio::test]
1473    async fn test_filter_long_eq() {
1474        let mut fixture = TableTestFixture::new();
1475        fixture.setup_manifest_files().await;
1476
1477        // Filter: i64 == 150
1478        let mut builder = fixture.table.scan();
1479        let predicate = Reference::new("i64").equal_to(Datum::long(150i64));
1480        builder = builder
1481            .with_filter(predicate)
1482            .with_row_selection_enabled(true);
1483        let table_scan = builder.build().unwrap();
1484
1485        let batch_stream = table_scan.to_arrow().await.unwrap();
1486
1487        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1488
1489        assert_eq!(batches.len(), 2);
1490        assert_eq!(batches[0].num_rows(), 12);
1491
1492        let col = batches[0].column_by_name("i64").unwrap();
1493        let i64_arr = col.as_any().downcast_ref::<Int64Array>().unwrap();
1494        assert_eq!(i64_arr.value(1), 150i64);
1495    }
1496
1497    #[tokio::test]
1498    async fn test_filter_bool_eq() {
1499        let mut fixture = TableTestFixture::new();
1500        fixture.setup_manifest_files().await;
1501
1502        // Filter: bool == true
1503        let mut builder = fixture.table.scan();
1504        let predicate = Reference::new("bool").equal_to(Datum::bool(true));
1505        builder = builder
1506            .with_filter(predicate)
1507            .with_row_selection_enabled(true);
1508        let table_scan = builder.build().unwrap();
1509
1510        let batch_stream = table_scan.to_arrow().await.unwrap();
1511
1512        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1513
1514        assert_eq!(batches.len(), 2);
1515        assert_eq!(batches[0].num_rows(), 512);
1516
1517        let col = batches[0].column_by_name("bool").unwrap();
1518        let bool_arr = col.as_any().downcast_ref::<BooleanArray>().unwrap();
1519        assert!(bool_arr.value(1));
1520    }
1521
1522    #[tokio::test]
1523    async fn test_filter_on_arrow_is_null() {
1524        let mut fixture = TableTestFixture::new();
1525        fixture.setup_manifest_files().await;
1526
1527        // Filter: y is null
1528        let mut builder = fixture.table.scan();
1529        let predicate = Reference::new("y").is_null();
1530        builder = builder
1531            .with_filter(predicate)
1532            .with_row_selection_enabled(true);
1533        let table_scan = builder.build().unwrap();
1534
1535        let batch_stream = table_scan.to_arrow().await.unwrap();
1536
1537        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1538        assert_eq!(batches.len(), 0);
1539    }
1540
1541    #[tokio::test]
1542    async fn test_filter_on_arrow_is_not_null() {
1543        let mut fixture = TableTestFixture::new();
1544        fixture.setup_manifest_files().await;
1545
1546        // Filter: y is not null
1547        let mut builder = fixture.table.scan();
1548        let predicate = Reference::new("y").is_not_null();
1549        builder = builder
1550            .with_filter(predicate)
1551            .with_row_selection_enabled(true);
1552        let table_scan = builder.build().unwrap();
1553
1554        let batch_stream = table_scan.to_arrow().await.unwrap();
1555
1556        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1557        assert_eq!(batches[0].num_rows(), 1024);
1558    }
1559
1560    #[tokio::test]
1561    async fn test_filter_on_arrow_lt_and_gt() {
1562        let mut fixture = TableTestFixture::new();
1563        fixture.setup_manifest_files().await;
1564
1565        // Filter: y < 5 AND z >= 4
1566        let mut builder = fixture.table.scan();
1567        let predicate = Reference::new("y")
1568            .less_than(Datum::long(5))
1569            .and(Reference::new("z").greater_than_or_equal_to(Datum::long(4)));
1570        builder = builder
1571            .with_filter(predicate)
1572            .with_row_selection_enabled(true);
1573        let table_scan = builder.build().unwrap();
1574
1575        let batch_stream = table_scan.to_arrow().await.unwrap();
1576
1577        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1578        assert_eq!(batches[0].num_rows(), 500);
1579
1580        let col = batches[0].column_by_name("x").unwrap();
1581        let expected_x = Arc::new(Int64Array::from_iter_values(vec![1; 500])) as ArrayRef;
1582        assert_eq!(col, &expected_x);
1583
1584        let col = batches[0].column_by_name("y").unwrap();
1585        let mut values = vec![];
1586        values.append(vec![3; 200].as_mut());
1587        values.append(vec![4; 300].as_mut());
1588        let expected_y = Arc::new(Int64Array::from_iter_values(values)) as ArrayRef;
1589        assert_eq!(col, &expected_y);
1590
1591        let col = batches[0].column_by_name("z").unwrap();
1592        let expected_z = Arc::new(Int64Array::from_iter_values(vec![4; 500])) as ArrayRef;
1593        assert_eq!(col, &expected_z);
1594    }
1595
1596    #[tokio::test]
1597    async fn test_filter_on_arrow_lt_or_gt() {
1598        let mut fixture = TableTestFixture::new();
1599        fixture.setup_manifest_files().await;
1600
1601        // Filter: y < 5 AND z >= 4
1602        let mut builder = fixture.table.scan();
1603        let predicate = Reference::new("y")
1604            .less_than(Datum::long(5))
1605            .or(Reference::new("z").greater_than_or_equal_to(Datum::long(4)));
1606        builder = builder
1607            .with_filter(predicate)
1608            .with_row_selection_enabled(true);
1609        let table_scan = builder.build().unwrap();
1610
1611        let batch_stream = table_scan.to_arrow().await.unwrap();
1612
1613        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1614        assert_eq!(batches[0].num_rows(), 1024);
1615
1616        let col = batches[0].column_by_name("x").unwrap();
1617        let expected_x = Arc::new(Int64Array::from_iter_values(vec![1; 1024])) as ArrayRef;
1618        assert_eq!(col, &expected_x);
1619
1620        let col = batches[0].column_by_name("y").unwrap();
1621        let mut values = vec![2; 512];
1622        values.append(vec![3; 200].as_mut());
1623        values.append(vec![4; 300].as_mut());
1624        values.append(vec![5; 12].as_mut());
1625        let expected_y = Arc::new(Int64Array::from_iter_values(values)) as ArrayRef;
1626        assert_eq!(col, &expected_y);
1627
1628        let col = batches[0].column_by_name("z").unwrap();
1629        let mut values = vec![3; 512];
1630        values.append(vec![4; 512].as_mut());
1631        let expected_z = Arc::new(Int64Array::from_iter_values(values)) as ArrayRef;
1632        assert_eq!(col, &expected_z);
1633    }
1634
1635    #[tokio::test]
1636    async fn test_filter_on_arrow_startswith() {
1637        let mut fixture = TableTestFixture::new();
1638        fixture.setup_manifest_files().await;
1639
1640        // Filter: a STARTSWITH "Ice"
1641        let mut builder = fixture.table.scan();
1642        let predicate = Reference::new("a").starts_with(Datum::string("Ice"));
1643        builder = builder
1644            .with_filter(predicate)
1645            .with_row_selection_enabled(true);
1646        let table_scan = builder.build().unwrap();
1647
1648        let batch_stream = table_scan.to_arrow().await.unwrap();
1649
1650        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1651
1652        assert_eq!(batches[0].num_rows(), 512);
1653
1654        let col = batches[0].column_by_name("a").unwrap();
1655        let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1656        assert_eq!(string_arr.value(0), "Iceberg");
1657    }
1658
1659    #[tokio::test]
1660    async fn test_filter_on_arrow_not_startswith() {
1661        let mut fixture = TableTestFixture::new();
1662        fixture.setup_manifest_files().await;
1663
1664        // Filter: a NOT STARTSWITH "Ice"
1665        let mut builder = fixture.table.scan();
1666        let predicate = Reference::new("a").not_starts_with(Datum::string("Ice"));
1667        builder = builder
1668            .with_filter(predicate)
1669            .with_row_selection_enabled(true);
1670        let table_scan = builder.build().unwrap();
1671
1672        let batch_stream = table_scan.to_arrow().await.unwrap();
1673
1674        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1675
1676        assert_eq!(batches[0].num_rows(), 512);
1677
1678        let col = batches[0].column_by_name("a").unwrap();
1679        let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1680        assert_eq!(string_arr.value(0), "Apache");
1681    }
1682
1683    #[tokio::test]
1684    async fn test_filter_on_arrow_in() {
1685        let mut fixture = TableTestFixture::new();
1686        fixture.setup_manifest_files().await;
1687
1688        // Filter: a IN ("Sioux", "Iceberg")
1689        let mut builder = fixture.table.scan();
1690        let predicate =
1691            Reference::new("a").is_in([Datum::string("Sioux"), Datum::string("Iceberg")]);
1692        builder = builder
1693            .with_filter(predicate)
1694            .with_row_selection_enabled(true);
1695        let table_scan = builder.build().unwrap();
1696
1697        let batch_stream = table_scan.to_arrow().await.unwrap();
1698
1699        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1700
1701        assert_eq!(batches[0].num_rows(), 512);
1702
1703        let col = batches[0].column_by_name("a").unwrap();
1704        let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1705        assert_eq!(string_arr.value(0), "Iceberg");
1706    }
1707
1708    #[tokio::test]
1709    async fn test_filter_on_arrow_not_in() {
1710        let mut fixture = TableTestFixture::new();
1711        fixture.setup_manifest_files().await;
1712
1713        // Filter: a NOT IN ("Sioux", "Iceberg")
1714        let mut builder = fixture.table.scan();
1715        let predicate =
1716            Reference::new("a").is_not_in([Datum::string("Sioux"), Datum::string("Iceberg")]);
1717        builder = builder
1718            .with_filter(predicate)
1719            .with_row_selection_enabled(true);
1720        let table_scan = builder.build().unwrap();
1721
1722        let batch_stream = table_scan.to_arrow().await.unwrap();
1723
1724        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1725
1726        assert_eq!(batches[0].num_rows(), 512);
1727
1728        let col = batches[0].column_by_name("a").unwrap();
1729        let string_arr = col.as_any().downcast_ref::<StringArray>().unwrap();
1730        assert_eq!(string_arr.value(0), "Apache");
1731    }
1732
1733    fn file_scan_task_test_schema(primitive_type: PrimitiveType) -> Arc<Schema> {
1734        Arc::new(
1735            Schema::builder()
1736                .with_fields(vec![Arc::new(NestedField::required(
1737                    1,
1738                    "x",
1739                    Type::Primitive(primitive_type),
1740                ))])
1741                .build()
1742                .unwrap(),
1743        )
1744    }
1745
1746    fn assert_file_scan_task_serde_round_trip(task: FileScanTask) {
1747        // Regression test for https://github.com/apache/iceberg-rust/issues/3089.
1748        let serialized = serde_json::to_string(&task).unwrap();
1749        let deserialized: FileScanTask = serde_json::from_str(&serialized).unwrap();
1750
1751        assert_eq!(task, deserialized);
1752    }
1753
1754    fn file_scan_task_with_partition(
1755        primitive_type: PrimitiveType,
1756        transform: Transform,
1757        partition_value: Literal,
1758    ) -> FileScanTask {
1759        let schema = file_scan_task_test_schema(primitive_type);
1760        let partition_spec = Arc::new(
1761            PartitionSpec::builder(schema.clone())
1762                .add_partition_field("x", "x_partition", transform)
1763                .unwrap()
1764                .build()
1765                .unwrap(),
1766        );
1767        FileScanTask::builder()
1768            .with_data_file_path("data_file_path".to_string())
1769            .with_file_size_in_bytes(123)
1770            .with_start(10)
1771            .with_length(100)
1772            .with_project_field_ids(vec![1])
1773            .with_schema(schema)
1774            .with_data_file_format(DataFileFormat::Parquet)
1775            .with_partition(Some(Struct::from_iter([Some(partition_value)])))
1776            .with_partition_spec(Some(partition_spec))
1777            .with_case_sensitive(true)
1778            .build()
1779            .unwrap()
1780    }
1781
1782    #[test]
1783    fn test_file_scan_task_serde_without_predicate() {
1784        let task = FileScanTask::builder()
1785            .with_data_file_path("data_file_path".to_string())
1786            .with_file_size_in_bytes(0)
1787            .with_start(0)
1788            .with_length(100)
1789            .with_project_field_ids(vec![1, 2, 3])
1790            .with_schema(file_scan_task_test_schema(PrimitiveType::Binary))
1791            .with_record_count(Some(100))
1792            .with_first_row_id(Some(1000))
1793            .with_data_sequence_number(Some(5))
1794            .with_data_file_format(DataFileFormat::Parquet)
1795            .with_case_sensitive(false)
1796            .build()
1797            .unwrap();
1798        assert_file_scan_task_serde_round_trip(task);
1799    }
1800
1801    #[test]
1802    fn test_file_scan_task_serde_with_predicate() {
1803        let task = FileScanTask::builder()
1804            .with_data_file_path("data_file_path".to_string())
1805            .with_file_size_in_bytes(0)
1806            .with_start(0)
1807            .with_length(100)
1808            .with_project_field_ids(vec![1, 2, 3])
1809            .with_predicate(Some(BoundPredicate::AlwaysTrue))
1810            .with_schema(file_scan_task_test_schema(PrimitiveType::Binary))
1811            .with_data_file_format(DataFileFormat::Avro)
1812            .with_case_sensitive(false)
1813            .build()
1814            .unwrap();
1815
1816        let serialized = serde_json::to_value(&task).unwrap();
1817        assert!(serialized.get("record_count").is_none());
1818        assert_file_scan_task_serde_round_trip(task);
1819    }
1820
1821    #[test]
1822    fn test_unpartitioned_file_scan_task_serde() {
1823        let task = FileScanTask::builder()
1824            .with_data_file_path("data_file_path".to_string())
1825            .with_file_size_in_bytes(0)
1826            .with_start(0)
1827            .with_length(100)
1828            .with_project_field_ids(vec![1, 2, 3])
1829            .with_schema(file_scan_task_test_schema(PrimitiveType::Binary))
1830            .with_data_file_format(DataFileFormat::Parquet)
1831            .with_partition(Some(Struct::empty()))
1832            .with_case_sensitive(false)
1833            .build()
1834            .unwrap();
1835        assert_file_scan_task_serde_round_trip(task);
1836    }
1837
1838    #[test]
1839    fn test_file_scan_task_serde_with_all_optional_fields() {
1840        let schema = file_scan_task_test_schema(PrimitiveType::Long);
1841        let partition_spec = Arc::new(
1842            PartitionSpec::builder(schema.clone())
1843                .add_partition_field("x", "x", Transform::Identity)
1844                .unwrap()
1845                .build()
1846                .unwrap(),
1847        );
1848        let unified_partition_type = Arc::new(partition_spec.partition_type(&schema).unwrap());
1849        let sort_order = Arc::new(
1850            SortOrder::builder()
1851                .with_order_id(1)
1852                .with_sort_field(
1853                    SortField::builder()
1854                        .source_id(1)
1855                        .transform(Transform::Identity)
1856                        .direction(SortDirection::Ascending)
1857                        .null_order(NullOrder::First)
1858                        .build(),
1859                )
1860                .build(&schema)
1861                .unwrap(),
1862        );
1863        let task = FileScanTask::builder()
1864            .with_data_file_path("data_file_path".to_string())
1865            .with_file_size_in_bytes(123)
1866            .with_start(10)
1867            .with_length(100)
1868            .with_project_field_ids(vec![1])
1869            .with_schema(schema)
1870            .with_data_file_format(DataFileFormat::Parquet)
1871            .with_deletes(vec![
1872                FileScanTaskDeleteFile::builder()
1873                    .with_file_path("delete_file_path".to_string())
1874                    .with_file_size_in_bytes(23)
1875                    .with_file_type(DataContentType::EqualityDeletes)
1876                    .with_file_format(DataFileFormat::Parquet)
1877                    .with_partition_spec_id(0)
1878                    .with_equality_ids(Some(vec![1]))
1879                    .with_referenced_data_file(Some("data_file_path".to_string()))
1880                    .with_content_offset(Some(12))
1881                    .with_content_size_in_bytes(Some(34))
1882                    .with_record_count(Some(5))
1883                    .with_key_metadata(Some(vec![4, 5, 6].into_boxed_slice()))
1884                    .build()
1885                    .unwrap(),
1886            ])
1887            .with_partition(Some(Struct::from_iter([Some(Literal::long(42))])))
1888            .with_partition_spec(Some(partition_spec))
1889            .with_name_mapping(Some(Arc::new(NameMapping::new(vec![MappedField::new(
1890                Some(1),
1891                vec!["x".to_string()],
1892                vec![],
1893            )]))))
1894            .with_unified_partition_type(Some(unified_partition_type))
1895            .with_sort_order_id(Some(1))
1896            .with_sort_order(Some(sort_order))
1897            .with_case_sensitive(true)
1898            .with_key_metadata(Some(vec![1, 2, 3].into_boxed_slice()))
1899            .build()
1900            .unwrap();
1901        assert_file_scan_task_serde_round_trip(task);
1902    }
1903
1904    #[test]
1905    fn test_file_scan_task_serde_with_date_partition() {
1906        let task = file_scan_task_with_partition(
1907            PrimitiveType::Date,
1908            Transform::Identity,
1909            Literal::date(19_000),
1910        );
1911        assert_file_scan_task_serde_round_trip(task);
1912    }
1913
1914    #[test]
1915    fn test_file_scan_task_serde_with_timestamp_ns_partition() {
1916        let task = file_scan_task_with_partition(
1917            PrimitiveType::TimestampNs,
1918            Transform::Identity,
1919            Literal::timestamp_nano(1_510_871_468_123_456_789),
1920        );
1921        assert_file_scan_task_serde_round_trip(task);
1922    }
1923
1924    #[test]
1925    fn test_file_scan_task_serde_with_timestamptz_ns_partition() {
1926        let task = file_scan_task_with_partition(
1927            PrimitiveType::TimestamptzNs,
1928            Transform::Identity,
1929            Literal::timestamptz_nano(1_510_871_468_123_456_789),
1930        );
1931        assert_file_scan_task_serde_round_trip(task);
1932    }
1933
1934    #[test]
1935    fn test_file_scan_task_serde_with_decimal_partition() {
1936        let task = file_scan_task_with_partition(
1937            PrimitiveType::Decimal {
1938                precision: 9,
1939                scale: 2,
1940            },
1941            Transform::Identity,
1942            Literal::decimal(12_345),
1943        );
1944        assert_file_scan_task_serde_round_trip(task);
1945    }
1946
1947    #[test]
1948    fn test_file_scan_task_serde_with_uuid_partition() {
1949        let task = file_scan_task_with_partition(
1950            PrimitiveType::Uuid,
1951            Transform::Identity,
1952            Literal::uuid(Uuid::from_u128(0x12345678_90ab_cdef_1234_567890abcdef)),
1953        );
1954        assert_file_scan_task_serde_round_trip(task);
1955    }
1956
1957    #[test]
1958    fn test_file_scan_task_serde_with_fixed_partition() {
1959        let task = file_scan_task_with_partition(
1960            PrimitiveType::Fixed(4),
1961            Transform::Identity,
1962            Literal::fixed([1, 2, 3, 4]),
1963        );
1964        assert_file_scan_task_serde_round_trip(task);
1965    }
1966
1967    #[test]
1968    fn test_file_scan_task_serde_with_bucket_partition() {
1969        let task = file_scan_task_with_partition(
1970            PrimitiveType::String,
1971            Transform::Bucket(4),
1972            Literal::int(2),
1973        );
1974        assert_file_scan_task_serde_round_trip(task);
1975    }
1976
1977    #[tokio::test]
1978    async fn test_select_with_file_column() {
1979        let mut fixture = TableTestFixture::new();
1980        fixture.setup_manifest_files().await;
1981
1982        // Select regular columns plus the _file column
1983        let table_scan = fixture
1984            .table
1985            .scan()
1986            .select(["x", RESERVED_COL_NAME_FILE])
1987            .with_row_selection_enabled(true)
1988            .build()
1989            .unwrap();
1990
1991        let batch_stream = table_scan.to_arrow().await.unwrap();
1992        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
1993
1994        // Verify we have 2 columns: x and _file
1995        assert_eq!(batches[0].num_columns(), 2);
1996
1997        // Verify the x column exists and has correct data
1998        let x_col = batches[0].column_by_name("x").unwrap();
1999        let x_arr = x_col.as_primitive::<arrow_array::types::Int64Type>();
2000        assert_eq!(x_arr.value(0), 1);
2001
2002        // Verify the _file column exists
2003        let file_col = batches[0].column_by_name(RESERVED_COL_NAME_FILE);
2004        assert!(
2005            file_col.is_some(),
2006            "_file column should be present in the batch"
2007        );
2008
2009        // Verify the _file column contains a file path
2010        let file_col = file_col.unwrap();
2011        assert!(
2012            matches!(
2013                file_col.data_type(),
2014                arrow_schema::DataType::RunEndEncoded(_, _)
2015            ),
2016            "_file column should use RunEndEncoded type"
2017        );
2018
2019        // Decode the RunArray to verify it contains the file path
2020        let run_array = file_col
2021            .as_any()
2022            .downcast_ref::<RunArray<Int32Type>>()
2023            .expect("_file column should be a RunArray");
2024
2025        let values = run_array.values();
2026        let string_values = values.as_string::<i32>();
2027        assert_eq!(string_values.len(), 1, "Should have a single file path");
2028
2029        let file_path = string_values.value(0);
2030        assert!(
2031            file_path.ends_with(".parquet"),
2032            "File path should end with .parquet, got: {file_path}"
2033        );
2034    }
2035
2036    #[tokio::test]
2037    async fn test_select_file_column_position() {
2038        let mut fixture = TableTestFixture::new();
2039        fixture.setup_manifest_files().await;
2040
2041        // Select columns in specific order: x, _file, z
2042        let table_scan = fixture
2043            .table
2044            .scan()
2045            .select(["x", RESERVED_COL_NAME_FILE, "z"])
2046            .with_row_selection_enabled(true)
2047            .build()
2048            .unwrap();
2049
2050        let batch_stream = table_scan.to_arrow().await.unwrap();
2051        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2052
2053        assert_eq!(batches[0].num_columns(), 3);
2054
2055        // Verify column order: x at position 0, _file at position 1, z at position 2
2056        let schema = batches[0].schema();
2057        assert_eq!(schema.field(0).name(), "x");
2058        assert_eq!(schema.field(1).name(), RESERVED_COL_NAME_FILE);
2059        assert_eq!(schema.field(2).name(), "z");
2060
2061        // Verify columns by name also works
2062        assert!(batches[0].column_by_name("x").is_some());
2063        assert!(batches[0].column_by_name(RESERVED_COL_NAME_FILE).is_some());
2064        assert!(batches[0].column_by_name("z").is_some());
2065    }
2066
2067    #[tokio::test]
2068    async fn test_select_file_column_only() {
2069        let mut fixture = TableTestFixture::new();
2070        fixture.setup_manifest_files().await;
2071
2072        // Select only the _file column
2073        let table_scan = fixture
2074            .table
2075            .scan()
2076            .select([RESERVED_COL_NAME_FILE])
2077            .with_row_selection_enabled(true)
2078            .build()
2079            .unwrap();
2080
2081        let batch_stream = table_scan.to_arrow().await.unwrap();
2082        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2083
2084        // Should have exactly 1 column
2085        assert_eq!(batches[0].num_columns(), 1);
2086
2087        // Verify it's the _file column
2088        let schema = batches[0].schema();
2089        assert_eq!(schema.field(0).name(), RESERVED_COL_NAME_FILE);
2090
2091        // Verify the batch has the correct number of rows
2092        // The scan reads files 1.parquet and 3.parquet (2.parquet is deleted)
2093        // Each file has 1024 rows, so total is 2048 rows
2094        let total_rows: usize = batches.iter().map(|b| b.num_rows()).sum();
2095        assert_eq!(total_rows, 2048);
2096    }
2097
2098    #[tokio::test]
2099    async fn test_file_column_with_multiple_files() {
2100        use std::collections::HashSet;
2101
2102        let mut fixture = TableTestFixture::new();
2103        fixture.setup_manifest_files().await;
2104
2105        // Select x and _file columns
2106        let table_scan = fixture
2107            .table
2108            .scan()
2109            .select(["x", RESERVED_COL_NAME_FILE])
2110            .with_row_selection_enabled(true)
2111            .build()
2112            .unwrap();
2113
2114        let batch_stream = table_scan.to_arrow().await.unwrap();
2115        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2116
2117        // Collect all unique file paths from the batches
2118        let mut file_paths = HashSet::new();
2119        for batch in &batches {
2120            let file_col = batch.column_by_name(RESERVED_COL_NAME_FILE).unwrap();
2121            let run_array = file_col
2122                .as_any()
2123                .downcast_ref::<RunArray<Int32Type>>()
2124                .expect("_file column should be a RunArray");
2125
2126            let values = run_array.values();
2127            let string_values = values.as_string::<i32>();
2128            for i in 0..string_values.len() {
2129                file_paths.insert(string_values.value(i).to_string());
2130            }
2131        }
2132
2133        // We should have multiple files (the test creates 1.parquet and 3.parquet)
2134        assert!(!file_paths.is_empty(), "Should have at least one file path");
2135
2136        // All paths should end with .parquet
2137        for path in &file_paths {
2138            assert!(
2139                path.ends_with(".parquet"),
2140                "All file paths should end with .parquet, got: {path}"
2141            );
2142        }
2143    }
2144
2145    #[tokio::test]
2146    async fn test_file_column_at_start() {
2147        let mut fixture = TableTestFixture::new();
2148        fixture.setup_manifest_files().await;
2149
2150        // Select _file at the start
2151        let table_scan = fixture
2152            .table
2153            .scan()
2154            .select([RESERVED_COL_NAME_FILE, "x", "y"])
2155            .with_row_selection_enabled(true)
2156            .build()
2157            .unwrap();
2158
2159        let batch_stream = table_scan.to_arrow().await.unwrap();
2160        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2161
2162        assert_eq!(batches[0].num_columns(), 3);
2163
2164        // Verify _file is at position 0
2165        let schema = batches[0].schema();
2166        assert_eq!(schema.field(0).name(), RESERVED_COL_NAME_FILE);
2167        assert_eq!(schema.field(1).name(), "x");
2168        assert_eq!(schema.field(2).name(), "y");
2169    }
2170
2171    #[tokio::test]
2172    async fn test_file_column_at_end() {
2173        let mut fixture = TableTestFixture::new();
2174        fixture.setup_manifest_files().await;
2175
2176        // Select _file at the end
2177        let table_scan = fixture
2178            .table
2179            .scan()
2180            .select(["x", "y", RESERVED_COL_NAME_FILE])
2181            .with_row_selection_enabled(true)
2182            .build()
2183            .unwrap();
2184
2185        let batch_stream = table_scan.to_arrow().await.unwrap();
2186        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2187
2188        assert_eq!(batches[0].num_columns(), 3);
2189
2190        // Verify _file is at position 2 (the end)
2191        let schema = batches[0].schema();
2192        assert_eq!(schema.field(0).name(), "x");
2193        assert_eq!(schema.field(1).name(), "y");
2194        assert_eq!(schema.field(2).name(), RESERVED_COL_NAME_FILE);
2195    }
2196
2197    #[tokio::test]
2198    async fn test_select_with_repeated_column_names() {
2199        let mut fixture = TableTestFixture::new();
2200        fixture.setup_manifest_files().await;
2201
2202        // Select with repeated column names - both regular columns and virtual columns
2203        // Repeated columns should appear multiple times in the result (duplicates are allowed)
2204        let table_scan = fixture
2205            .table
2206            .scan()
2207            .select([
2208                "x",
2209                RESERVED_COL_NAME_FILE,
2210                "x", // x repeated
2211                "y",
2212                RESERVED_COL_NAME_FILE, // _file repeated
2213                "y",                    // y repeated
2214            ])
2215            .with_row_selection_enabled(true)
2216            .build()
2217            .unwrap();
2218
2219        let batch_stream = table_scan.to_arrow().await.unwrap();
2220        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2221
2222        // Verify we have exactly 6 columns (duplicates are allowed and preserved)
2223        assert_eq!(
2224            batches[0].num_columns(),
2225            6,
2226            "Should have exactly 6 columns with duplicates"
2227        );
2228
2229        let schema = batches[0].schema();
2230
2231        // Verify columns appear in the exact order requested: x, _file, x, y, _file, y
2232        assert_eq!(schema.field(0).name(), "x", "Column 0 should be x");
2233        assert_eq!(
2234            schema.field(1).name(),
2235            RESERVED_COL_NAME_FILE,
2236            "Column 1 should be _file"
2237        );
2238        assert_eq!(
2239            schema.field(2).name(),
2240            "x",
2241            "Column 2 should be x (duplicate)"
2242        );
2243        assert_eq!(schema.field(3).name(), "y", "Column 3 should be y");
2244        assert_eq!(
2245            schema.field(4).name(),
2246            RESERVED_COL_NAME_FILE,
2247            "Column 4 should be _file (duplicate)"
2248        );
2249        assert_eq!(
2250            schema.field(5).name(),
2251            "y",
2252            "Column 5 should be y (duplicate)"
2253        );
2254
2255        // Verify all columns have correct data types
2256        assert!(
2257            matches!(schema.field(0).data_type(), arrow_schema::DataType::Int64),
2258            "Column x should be Int64"
2259        );
2260        assert!(
2261            matches!(schema.field(2).data_type(), arrow_schema::DataType::Int64),
2262            "Column x (duplicate) should be Int64"
2263        );
2264        assert!(
2265            matches!(schema.field(3).data_type(), arrow_schema::DataType::Int64),
2266            "Column y should be Int64"
2267        );
2268        assert!(
2269            matches!(schema.field(5).data_type(), arrow_schema::DataType::Int64),
2270            "Column y (duplicate) should be Int64"
2271        );
2272        assert!(
2273            matches!(
2274                schema.field(1).data_type(),
2275                arrow_schema::DataType::RunEndEncoded(_, _)
2276            ),
2277            "_file column should use RunEndEncoded type"
2278        );
2279        assert!(
2280            matches!(
2281                schema.field(4).data_type(),
2282                arrow_schema::DataType::RunEndEncoded(_, _)
2283            ),
2284            "_file column (duplicate) should use RunEndEncoded type"
2285        );
2286    }
2287
2288    /// Builds a minimal single-snapshot table (no manifests on disk) whose schema
2289    /// contains a column with the given name, so scan planning resolves column names
2290    /// against a real schema. `TableScan::build()` only reads metadata, not manifest
2291    /// files, so a snapshot pointing at a dummy manifest list is enough. Used to
2292    /// reproduce issue #2837.
2293    fn table_with_data_column(column_name: &str) -> Table {
2294        use crate::spec::{
2295            FormatVersion, MAIN_BRANCH, Operation, Snapshot, SnapshotReference, SnapshotRetention,
2296            SortOrder, Summary, TableMetadataBuilder, UnboundPartitionSpec,
2297        };
2298
2299        let schema = Schema::builder()
2300            .with_schema_id(0)
2301            .with_fields(vec![
2302                NestedField::required(1, "id", Type::Primitive(PrimitiveType::Int)).into(),
2303                NestedField::required(2, column_name, Type::Primitive(PrimitiveType::Int)).into(),
2304            ])
2305            .build()
2306            .unwrap();
2307
2308        let snapshot = Snapshot::builder()
2309            .with_snapshot_id(1)
2310            .with_timestamp_ms(1)
2311            .with_sequence_number(0)
2312            .with_schema_id(0)
2313            .with_manifest_list("/snap-1.avro")
2314            .with_summary(Summary {
2315                operation: Operation::Append,
2316                additional_properties: HashMap::new(),
2317            })
2318            .build();
2319
2320        let metadata = TableMetadataBuilder::new(
2321            schema,
2322            UnboundPartitionSpec::builder().with_spec_id(0).build(),
2323            SortOrder::unsorted_order(),
2324            "s3://bucket/table".to_string(),
2325            FormatVersion::V2,
2326            HashMap::new(),
2327        )
2328        .unwrap()
2329        .add_snapshot(snapshot)
2330        .unwrap()
2331        .set_ref(MAIN_BRANCH, SnapshotReference {
2332            snapshot_id: 1,
2333            retention: SnapshotRetention::Branch {
2334                min_snapshots_to_keep: None,
2335                max_snapshot_age_ms: None,
2336                max_ref_age_ms: None,
2337            },
2338        })
2339        .unwrap()
2340        .build()
2341        .unwrap()
2342        .metadata;
2343
2344        Table::builder()
2345            .metadata(metadata)
2346            .identifier(TableIdent::from_strs(["db", "table1"]).unwrap())
2347            .file_io(FileIO::new_with_fs())
2348            .runtime(test_runtime())
2349            .build()
2350            .unwrap()
2351    }
2352
2353    /// A user data column named `pos` (a delete-file internal column name that is not a
2354    /// data-table metadata column) must be projectable rather than shadowed. Regression
2355    /// test for issue #2837.
2356    #[test]
2357    fn test_scan_projects_data_column_named_like_delete_file_column() {
2358        for column_name in ["pos", "file_path"] {
2359            let table = table_with_data_column(column_name);
2360
2361            // Projecting the data column must succeed and resolve to its real field id (2),
2362            // not the reserved delete-file field id.
2363            let table_scan = table
2364                .scan()
2365                .select([column_name])
2366                .build()
2367                .unwrap_or_else(|e| panic!("scan of data column `{column_name}` failed: {e}"));
2368
2369            assert_eq!(
2370                table_scan.plan_context.as_ref().unwrap().field_ids.as_ref(),
2371                &[2]
2372            );
2373
2374            // The default projection (all columns) must resolve to the real field ids
2375            // too, not shadow the data column with a reserved delete-file id.
2376            let default_scan = table.scan().build().unwrap();
2377            assert_eq!(
2378                default_scan
2379                    .plan_context
2380                    .as_ref()
2381                    .unwrap()
2382                    .field_ids
2383                    .as_ref(),
2384                &[1, 2]
2385            );
2386        }
2387    }
2388
2389    /// Projecting a genuinely absent column still fails with a clear "not found" error
2390    /// rather than being silently accepted as a metadata column.
2391    #[test]
2392    fn test_scan_rejects_unknown_column_named_like_delete_file_column() {
2393        // This table has no `pos` column (only `id` and `file_path`).
2394        let table = table_with_data_column("file_path");
2395
2396        let err = table
2397            .scan()
2398            .select(["pos"])
2399            .build()
2400            .expect_err("projecting an absent column should fail");
2401        assert_eq!(err.kind(), ErrorKind::DataInvalid);
2402        assert!(err.to_string().contains("not found"));
2403    }
2404
2405    #[tokio::test]
2406    async fn test_scan_deadlock() {
2407        let mut fixture = TableTestFixture::new();
2408        fixture.setup_deadlock_manifests().await;
2409
2410        // Create table scan with concurrency limit 1
2411        // This sets channel size to 1.
2412        // Data manifest has 10 entries -> will block producer.
2413        // Delete manifest is 2nd in list -> won't be processed.
2414        // Consumer 2 (Data) not started -> blocked.
2415        // Consumer 1 (Delete) waiting -> blocked.
2416        let table_scan = fixture
2417            .table
2418            .scan()
2419            .with_concurrency_limit(1)
2420            .build()
2421            .unwrap();
2422
2423        // This should timeout/hang if deadlock exists
2424        // We can use tokio::time::timeout
2425        let result = tokio::time::timeout(std::time::Duration::from_secs(5), async {
2426            table_scan
2427                .plan_files()
2428                .await
2429                .unwrap()
2430                .try_collect::<Vec<_>>()
2431                .await
2432        })
2433        .await;
2434
2435        // Assert it finished (didn't timeout)
2436        assert!(result.is_ok(), "Scan timed out - deadlock detected");
2437    }
2438
2439    #[tokio::test]
2440    async fn test_select_with_spec_id_column() {
2441        let mut fixture = TableTestFixture::new();
2442        fixture.setup_manifest_files().await;
2443
2444        // Select regular columns plus the _spec_id column
2445        let table_scan = fixture
2446            .table
2447            .scan()
2448            .select(["x", RESERVED_COL_NAME_SPEC_ID, "z"])
2449            .with_row_selection_enabled(true)
2450            .build()
2451            .unwrap();
2452
2453        let batch_stream = table_scan.to_arrow().await.unwrap();
2454        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2455
2456        // Verify we have 3 columns: x, _spec_id, and z
2457        assert_eq!(batches[0].num_columns(), 3);
2458
2459        // Verify the x column exists and has correct data
2460        let col1 = batches[0].column_by_name("x").unwrap();
2461        let int64_arr = col1.as_any().downcast_ref::<Int64Array>().unwrap();
2462        assert_eq!(int64_arr.value(0), 1);
2463
2464        // Verify the _spec_id column exists
2465        let spec_id_col = batches[0].column_by_name(RESERVED_COL_NAME_SPEC_ID);
2466        assert!(
2467            spec_id_col.is_some(),
2468            "_spec_id column should be present in the batch"
2469        );
2470
2471        // Verify the _spec_id data type
2472        let spec_id_col = spec_id_col.unwrap();
2473        assert!(
2474            matches!(
2475                spec_id_col.data_type(),
2476                arrow_schema::DataType::RunEndEncoded(_, _)
2477            ),
2478            "_spec_id column should use RunEndEncoded type"
2479        );
2480
2481        // Decode the RunArray to verify it contains the spec id
2482        let run_array = spec_id_col
2483            .as_any()
2484            .downcast_ref::<RunArray<Int32Type>>()
2485            .expect("_spec_id column should be a RunArray");
2486
2487        let values = run_array.values();
2488        let int_values = values.as_primitive::<Int32Type>();
2489        assert_eq!(int_values.len(), 1, "Should have a single _spec_id");
2490
2491        let spec_id = int_values.value(0);
2492        assert_eq!(spec_id, 0, "_spec_id should be 0, got: {spec_id}");
2493
2494        // Verify 'z' column exists
2495        assert!(batches[0].column_by_name("z").is_some());
2496    }
2497
2498    #[tokio::test]
2499    async fn test_select_with_last_updated_sequence_number_column() {
2500        // A v2 fixture: data files have a null first_row_id. Per the spec's Row
2501        // Lineage read rules, a file with a null first_row_id produces a null
2502        // _last_updated_sequence_number for all rows; both lineage columns are
2503        // gated on first_row_id.
2504        let mut fixture = TableTestFixture::new();
2505        fixture.setup_manifest_files().await;
2506
2507        let table_scan = fixture
2508            .table
2509            .scan()
2510            .select(["x", RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER])
2511            .with_row_selection_enabled(true)
2512            .build()
2513            .unwrap();
2514
2515        let batches: Vec<_> = table_scan
2516            .to_arrow()
2517            .await
2518            .unwrap()
2519            .try_collect()
2520            .await
2521            .unwrap();
2522
2523        // Every row's value is null (v2 files have no first_row_id).
2524        assert_last_updated_seq_all(&batches, None);
2525    }
2526
2527    #[tokio::test]
2528    async fn test_select_with_last_updated_sequence_number_column_v3() {
2529        // A v3 fixture: the data file inherits a first_row_id and a data sequence
2530        // number through manifest read. End to end, the projected
2531        // _last_updated_sequence_number materializes to the file's data sequence
2532        // number for every row, exercising the full inherit, populate and
2533        // materialize wiring, not just a hand-built task.
2534        let mut fixture = TableTestFixture::new();
2535        fixture.setup_v3_manifest_files().await;
2536
2537        // The added file inherits the current snapshot's sequence number; derive it
2538        // from the fixture rather than hardcoding so the assertion tracks the fixture.
2539        let expected_seq = fixture
2540            .table
2541            .metadata()
2542            .current_snapshot()
2543            .unwrap()
2544            .sequence_number();
2545
2546        let table_scan = fixture
2547            .table
2548            .scan()
2549            .select(["x", RESERVED_COL_NAME_LAST_UPDATED_SEQUENCE_NUMBER])
2550            .with_row_selection_enabled(true)
2551            .build()
2552            .unwrap();
2553
2554        let batches: Vec<_> = table_scan
2555            .to_arrow()
2556            .await
2557            .unwrap()
2558            .try_collect()
2559            .await
2560            .unwrap();
2561
2562        assert_last_updated_seq_all(&batches, Some(expected_seq));
2563    }
2564
2565    #[tokio::test]
2566    async fn test_select_with_spec_id_column_from_unpartitioned_table() {
2567        let mut fixture = TableTestFixture::new_unpartitioned();
2568        fixture.setup_unpartitioned_manifest_files().await;
2569
2570        // Select regular columns plus the _spec_id column
2571        let table_scan = fixture
2572            .table
2573            .scan()
2574            .select(["x", RESERVED_COL_NAME_SPEC_ID])
2575            .with_row_selection_enabled(true)
2576            .build()
2577            .unwrap();
2578
2579        let batch_stream = table_scan.to_arrow().await.unwrap();
2580        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2581
2582        // Verify we have 2 columns: x and _spec_id
2583        assert_eq!(batches[0].num_columns(), 2);
2584
2585        // Verify the _spec_id column exists
2586        let spec_id_col = batches[0].column_by_name(RESERVED_COL_NAME_SPEC_ID);
2587        assert!(
2588            spec_id_col.is_some(),
2589            "_spec_id column should be present in the batch"
2590        );
2591
2592        // Verify the _spec_id data type
2593        let spec_id_col = spec_id_col.unwrap();
2594        assert!(
2595            matches!(
2596                spec_id_col.data_type(),
2597                arrow_schema::DataType::RunEndEncoded(_, _)
2598            ),
2599            "_spec_id column should use RunEndEncoded type"
2600        );
2601
2602        // Decode the RunArray to verify it contains the spec id
2603        let run_array = spec_id_col
2604            .as_any()
2605            .downcast_ref::<RunArray<Int32Type>>()
2606            .expect("_spec_id column should be a RunArray");
2607
2608        let values = run_array.values();
2609        let int_values = values.as_primitive::<Int32Type>();
2610        assert_eq!(int_values.len(), 1, "Should have a single _spec_id");
2611
2612        let spec_id = int_values.value(0);
2613        assert_eq!(spec_id, 0, "_spec_id should be 0, got: {spec_id}");
2614    }
2615
2616    #[tokio::test]
2617    async fn test_select_with_spec_id_column_with_partition_evolution() {
2618        let mut fixture = TableTestFixture::new_with_partition_evolution();
2619        fixture
2620            .setup_manifest_files_with_partition_evolution()
2621            .await;
2622
2623        // Select regular columns plus the _spec_id column
2624        let table_scan = fixture
2625            .table
2626            .scan()
2627            .select(["x", RESERVED_COL_NAME_SPEC_ID, "z"])
2628            .with_row_selection_enabled(true)
2629            .build()
2630            .unwrap();
2631
2632        let batch_stream = table_scan.to_arrow().await.unwrap();
2633        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2634
2635        // Verify the x column exists and has correct data
2636        let col1 = batches[0].column_by_name("x").unwrap();
2637        let int64_arr = col1.as_any().downcast_ref::<Int64Array>().unwrap();
2638        assert_eq!(int64_arr.value(0), 1);
2639
2640        // Verify the _spec_id column exists
2641        let spec_id_col = batches[0].column_by_name(RESERVED_COL_NAME_SPEC_ID);
2642        assert!(
2643            spec_id_col.is_some(),
2644            "_spec_id column should be present in the batch"
2645        );
2646
2647        // Verify the _spec_id data type
2648        let spec_id_col = spec_id_col.unwrap();
2649        assert!(
2650            matches!(
2651                spec_id_col.data_type(),
2652                arrow_schema::DataType::RunEndEncoded(_, _)
2653            ),
2654            "_spec_id column should use RunEndEncoded type"
2655        );
2656
2657        // Decode the RunArray to verify it contains the spec id
2658        let run_array = spec_id_col
2659            .as_any()
2660            .downcast_ref::<RunArray<Int32Type>>()
2661            .expect("_spec_id column should be a RunArray");
2662
2663        let values = run_array.values();
2664        let int_values = values.as_primitive::<Int32Type>();
2665        assert_eq!(int_values.len(), 1, "Should have a single _spec_id");
2666
2667        let spec_id = int_values.value(0);
2668        assert_eq!(spec_id, 2, "_spec_id should be 2, got: {spec_id}");
2669    }
2670
2671    #[tokio::test]
2672    async fn test_select_with_pos_and_file_columns() {
2673        use arrow_array::cast::AsArray;
2674
2675        let mut fixture = TableTestFixture::new();
2676        fixture.setup_manifest_files().await;
2677
2678        // Select regular columns plus the _pos column
2679        let table_scan = fixture
2680            .table
2681            .scan()
2682            .select(["x", RESERVED_COL_NAME_POS, RESERVED_COL_NAME_FILE])
2683            .with_row_selection_enabled(true)
2684            .build()
2685            .unwrap();
2686
2687        let batch_stream = table_scan.to_arrow().await.unwrap();
2688        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2689        assert_eq!(batches.len(), 2);
2690
2691        // Examine batches are 1.paruqet and 3.parquet, 2.parquet is deleted.
2692        for batch in batches.iter() {
2693            // Verify we have 3 columns: x, _pos and _file
2694            assert_eq!(batch.num_columns(), 3);
2695
2696            // Verify the x column exists and has correct data
2697            let x_col = batch.column_by_name("x").unwrap();
2698            let x_arr = x_col.as_primitive::<arrow_array::types::Int64Type>();
2699            assert_eq!(x_arr.value(0), 1);
2700
2701            // The _pos column exists and verify it is Int64Array with the expected values
2702            let pos_col = batch.column(1);
2703            let pos_array: &Int64Array = pos_col
2704                .as_any()
2705                .downcast_ref::<Int64Array>()
2706                .expect("_pos column should be a Int64Array");
2707            assert_eq!(*pos_array, Int64Array::from_iter_values(0i64..1024));
2708
2709            // Verify the _file column exists
2710            let file_col = batch.column_by_name(RESERVED_COL_NAME_FILE);
2711            assert!(
2712                file_col.is_some(),
2713                "_file column should be present in the batch"
2714            );
2715        }
2716    }
2717
2718    #[tokio::test]
2719    async fn test_pos_column_at_start_with_filters() {
2720        let mut fixture = TableTestFixture::new();
2721        fixture.setup_manifest_files().await;
2722
2723        // y is in [4, 5)
2724        let predicate = Reference::new("y")
2725            .greater_than(Datum::long(4i64))
2726            .and(Reference::new("y").less_than_or_equal_to(Datum::long(5i64)));
2727        // Select _pos at the start
2728        let table_scan = fixture
2729            .table
2730            .scan()
2731            .select([RESERVED_COL_NAME_POS, "x", "y"])
2732            .with_filter(predicate)
2733            .with_row_selection_enabled(true)
2734            .build()
2735            .unwrap();
2736
2737        let batch_stream = table_scan.to_arrow().await.unwrap();
2738        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2739        assert_eq!(batches.len(), 2);
2740
2741        // Examine batches are 1.paruqet and 3.parquet, 2.parquet is deleted.
2742        for batch in batches.iter() {
2743            assert_eq!(batch.num_columns(), 3);
2744            assert_eq!(batch.num_rows(), 12);
2745
2746            // Verify _pos is at position 0
2747            let schema = batch.schema();
2748            assert_eq!(schema.field(0).name(), RESERVED_COL_NAME_POS);
2749            assert_eq!(schema.field(1).name(), "x");
2750            assert_eq!(schema.field(2).name(), "y");
2751
2752            let pos_col = batch.column(0);
2753            let pos_array: &Int64Array = pos_col
2754                .as_any()
2755                .downcast_ref::<Int64Array>()
2756                .expect("_pos column should be a Int64Array");
2757            assert_eq!(*pos_array, Int64Array::from_iter_values(1012i64..1024));
2758        }
2759    }
2760
2761    #[tokio::test]
2762    async fn test_repeated_pos_column_with_filter() {
2763        let mut fixture = TableTestFixture::new();
2764        fixture.setup_manifest_files().await;
2765
2766        // a NOT STARTSWITH "Apa"
2767        let predicate = Reference::new("a").not_starts_with(Datum::string("Apa"));
2768        // Select '_pos' columns twice
2769        let table_scan = fixture
2770            .table
2771            .scan()
2772            .select([RESERVED_COL_NAME_POS, "a", RESERVED_COL_NAME_POS, "x"])
2773            .with_row_selection_enabled(true)
2774            .with_filter(predicate)
2775            .build()
2776            .unwrap();
2777
2778        let batch_stream = table_scan.to_arrow().await.unwrap();
2779
2780        let batches: Vec<_> = batch_stream.try_collect().await.unwrap();
2781        assert_eq!(batches.len(), 2);
2782
2783        // Examine batches are 1.paruqet and 3.parquet, 2.parquet is deleted.
2784        for batch in batches.iter() {
2785            assert_eq!(batch.num_rows(), 512);
2786
2787            // fetch the 1st _pos column by name and verify it is Int64Array with the expected values
2788            let pos_col = batch
2789                .column_by_name("_pos")
2790                .expect("_pos column should be present in the batch");
2791            let pos_array: &Int64Array = pos_col
2792                .as_any()
2793                .downcast_ref::<Int64Array>()
2794                .expect("_pos column should be a Int64Array");
2795            assert_eq!(*pos_array, Int64Array::from_iter_values(512i64..1024));
2796
2797            // fetch the 2nd _pos column by index and verify it is Int64Array with the expected values
2798            let pos_col = batch.column(2);
2799            let pos_array: &Int64Array = pos_col
2800                .as_any()
2801                .downcast_ref::<Int64Array>()
2802                .expect("_pos column should be a Int64Array");
2803            assert_eq!(*pos_array, Int64Array::from_iter_values(512i64..1024));
2804        }
2805    }
2806
2807    /// End-to-end through `TableScan`: a data file with three row groups planned
2808    /// as a single whole-file `FileScanTask` must yield contiguous, file-absolute
2809    /// `_pos` values (0..300) across the row-group boundaries.
2810    #[tokio::test]
2811    async fn test_pos_across_row_groups_via_table_scan() {
2812        let mut fixture = TableTestFixture::new();
2813        fixture.setup_multi_row_group_manifest(&[]).await;
2814
2815        // Planning must produce exactly one whole-file task with _pos projected and
2816        // no delete files, confirming TableScan does not sub-split the file.
2817        let tasks: Vec<_> = fixture
2818            .table
2819            .scan()
2820            .select(["x", RESERVED_COL_NAME_POS])
2821            .build()
2822            .unwrap()
2823            .plan_files()
2824            .await
2825            .unwrap()
2826            .try_collect()
2827            .await
2828            .unwrap();
2829        assert_eq!(tasks.len(), 1, "expected a single FileScanTask");
2830        let task = &tasks[0];
2831        assert!(
2832            task.project_field_ids().contains(&RESERVED_FIELD_ID_POS),
2833            "_pos field id must be projected into the FileScanTask"
2834        );
2835        assert_eq!(task.start(), 0, "TableScan should plan whole-file tasks");
2836        assert_eq!(task.length(), task.file_size_in_bytes());
2837        assert!(task.deletes().is_empty());
2838
2839        // Reading that task yields absolute _pos 0..300 in order.
2840        let batches: Vec<_> = fixture
2841            .table
2842            .scan()
2843            .select(["x", RESERVED_COL_NAME_POS])
2844            .build()
2845            .unwrap()
2846            .to_arrow()
2847            .await
2848            .unwrap()
2849            .try_collect()
2850            .await
2851            .unwrap();
2852
2853        let pos: Vec<i64> = batches
2854            .iter()
2855            .flat_map(|b| {
2856                b.column_by_name(RESERVED_COL_NAME_POS)
2857                    .expect("_pos column should be present")
2858                    .as_any()
2859                    .downcast_ref::<Int64Array>()
2860                    .expect("_pos column should be a Int64Array")
2861                    .values()
2862                    .to_vec()
2863            })
2864            .collect();
2865        assert_eq!(pos, (0..300).collect::<Vec<i64>>());
2866
2867        // Sanity: x == 1000 + _pos, proving _pos aligns with the actual rows read.
2868        let x: Vec<i64> = batches
2869            .iter()
2870            .flat_map(|b| {
2871                b.column_by_name("x")
2872                    .unwrap()
2873                    .as_primitive::<arrow_array::types::Int64Type>()
2874                    .values()
2875                    .to_vec()
2876            })
2877            .collect();
2878        assert_eq!(x, (1000..1300).collect::<Vec<i64>>());
2879    }
2880
2881    /// A positional delete file registered in the manifest must be attached to the planned
2882    /// `FileScanTask` and applied on read, while surviving `_pos` values stay file-absolute.
2883    #[tokio::test]
2884    async fn test_pos_with_positional_deletes_via_table_scan() {
2885        let mut fixture = TableTestFixture::new();
2886        // Delete file-absolute positions 150 (middle row group) and 299 (last row).
2887        fixture.setup_multi_row_group_manifest(&[150, 299]).await;
2888
2889        // Planning must attach the positional delete file to the task.
2890        let tasks: Vec<_> = fixture
2891            .table
2892            .scan()
2893            .select(["x", RESERVED_COL_NAME_POS])
2894            .build()
2895            .unwrap()
2896            .plan_files()
2897            .await
2898            .unwrap()
2899            .try_collect()
2900            .await
2901            .unwrap();
2902        assert_eq!(tasks.len(), 1);
2903        assert_eq!(
2904            tasks[0].deletes().len(),
2905            1,
2906            "positional delete file should be planned into the task"
2907        );
2908        assert_eq!(
2909            tasks[0].deletes()[0].file_type(),
2910            DataContentType::PositionDeletes
2911        );
2912
2913        // Reading applies the deletes; _pos must skip 150 and 299 and stay absolute.
2914        let batches: Vec<_> = fixture
2915            .table
2916            .scan()
2917            .select(["x", RESERVED_COL_NAME_POS])
2918            .build()
2919            .unwrap()
2920            .to_arrow()
2921            .await
2922            .unwrap()
2923            .try_collect()
2924            .await
2925            .unwrap();
2926
2927        let pos: Vec<i64> = batches
2928            .iter()
2929            .flat_map(|b| {
2930                b.column_by_name(RESERVED_COL_NAME_POS)
2931                    .expect("_pos column should be present")
2932                    .as_any()
2933                    .downcast_ref::<Int64Array>()
2934                    .expect("_pos column should be a Int64Array")
2935                    .values()
2936                    .to_vec()
2937            })
2938            .collect();
2939
2940        let total: usize = batches.iter().map(|b| b.num_rows()).sum();
2941        assert_eq!(
2942            total, 298,
2943            "two rows should be removed by positional deletes"
2944        );
2945        assert!(!pos.contains(&150) && !pos.contains(&299), "got {pos:?}");
2946        let expected: Vec<i64> = (0..150).chain(151..299).collect();
2947        assert_eq!(pos, expected);
2948    }
2949
2950    /// A filter that only matches the middle row group (`y` in [1100, 1200)) and
2951    /// prunes the other two row groups by statistics, so only the middle row group
2952    /// is read. `_pos` must report the file-absolute positions 100..200 for those
2953    /// rows, not values reset to 0..100.
2954    ///
2955    /// `y` is a non-partition column, so pruning here is driven purely by Parquet
2956    /// row-group statistics (the TableTestFixture's partition column is `x`).
2957    #[tokio::test]
2958    async fn test_pos_reads_only_middle_row_group_via_filter() {
2959        let mut fixture = TableTestFixture::new();
2960        fixture.setup_multi_row_group_manifest(&[]).await;
2961
2962        // Middle row group holds y = 1100..1200 at file positions 100..200.
2963        let predicate = Reference::new("y")
2964            .greater_than_or_equal_to(Datum::long(1100))
2965            .and(Reference::new("y").less_than(Datum::long(1200)));
2966
2967        let batches: Vec<_> = fixture
2968            .table
2969            .scan()
2970            .select(["y", RESERVED_COL_NAME_POS])
2971            .with_filter(predicate)
2972            .with_row_group_filtering_enabled(true)
2973            .build()
2974            .unwrap()
2975            .to_arrow()
2976            .await
2977            .unwrap()
2978            .try_collect()
2979            .await
2980            .unwrap();
2981
2982        let total: usize = batches.iter().map(|b| b.num_rows()).sum();
2983        assert_eq!(total, 100, "only the middle row group should be read");
2984
2985        let pos: Vec<i64> = batches
2986            .iter()
2987            .flat_map(|b| {
2988                b.column_by_name(RESERVED_COL_NAME_POS)
2989                    .expect("_pos column should be present")
2990                    .as_any()
2991                    .downcast_ref::<Int64Array>()
2992                    .expect("_pos column should be a Int64Array")
2993                    .values()
2994                    .to_vec()
2995            })
2996            .collect();
2997        assert_eq!(
2998            pos,
2999            (100..200).collect::<Vec<i64>>(),
3000            "_pos must be file-absolute for the middle row group"
3001        );
3002
3003        // Cross-check: y == 1000 + _pos for every surviving row.
3004        let y: Vec<i64> = batches
3005            .iter()
3006            .flat_map(|b| {
3007                b.column_by_name("y")
3008                    .unwrap()
3009                    .as_primitive::<arrow_array::types::Int64Type>()
3010                    .values()
3011                    .to_vec()
3012            })
3013            .collect();
3014        assert_eq!(y, (1100..1200).collect::<Vec<i64>>());
3015    }
3016}