Skip to main content

iceberg/scan/
task.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
18use std::sync::Arc;
19
20use futures::stream::BoxStream;
21use serde::{Deserialize, Serialize};
22use typed_builder::TypedBuilder;
23
24use crate::error::invalid_data;
25use crate::expr::BoundPredicate;
26use crate::spec::{
27    DataContentType, DataFileFormat, ManifestEntryRef, NameMapping, PartitionSpec, Schema,
28    SchemaRef, SortOrderRef, Struct, StructType,
29};
30use crate::{Error, Result};
31
32/// A stream of [`FileScanTask`].
33pub type FileScanTaskStream = BoxStream<'static, Result<FileScanTask>>;
34
35/// A task to scan part of file.
36#[derive(Debug, Clone, Deserialize, PartialEq, TypedBuilder)]
37#[serde(try_from = "crate::scan::task::_serde::FileScanTaskSerde")]
38#[builder(
39    field_defaults(setter(prefix = "with_")),
40    build_method(into = Result<FileScanTask>)
41)]
42pub struct FileScanTask {
43    /// The total size of the data file in bytes, from the manifest entry.
44    /// Used to skip a stat/HEAD request when reading Parquet footers.
45    file_size_in_bytes: u64,
46    /// The start offset of the file to scan.
47    start: u64,
48    /// The length of the file to scan.
49    length: u64,
50    /// The number of records in the file to scan.
51    ///
52    /// This is an optional field, and only available if we are
53    /// reading the entire data file.
54    #[builder(default)]
55    record_count: Option<u64>,
56
57    /// The first row id assigned to the data file.
58    ///
59    /// Used to derive the `_row_id` metadata column: for a row without an
60    /// explicit `_row_id`, it is this value plus the row's ordinal position.
61    #[builder(default)]
62    first_row_id: Option<i64>,
63
64    /// The data sequence number of the file, as opposed to its file sequence
65    /// number: the sequence number preserved when a file is carried forward
66    /// across a rewrite. May be null for an existing entry in a malformed
67    /// manifest that lacks one.
68    ///
69    /// Used to derive the `_last_updated_sequence_number` metadata column.
70    #[builder(default)]
71    data_sequence_number: Option<i64>,
72
73    /// The data file path corresponding to the task.
74    data_file_path: String,
75
76    /// The format of the file to scan.
77    data_file_format: DataFileFormat,
78
79    /// The schema of the file to scan.
80    schema: SchemaRef,
81    /// The field ids to project.
82    project_field_ids: Vec<i32>,
83    /// The predicate to filter.
84    #[builder(default)]
85    predicate: Option<BoundPredicate>,
86
87    /// The list of delete files that may need to be applied to this data file
88    #[builder(default)]
89    deletes: Vec<FileScanTaskDeleteFile>,
90
91    /// Partition data from the manifest entry, used to identify which columns can use
92    /// constant values from partition metadata vs. reading from the data file.
93    /// Per the Iceberg spec, only identity-transformed partition fields should use constants.
94    #[builder(default)]
95    partition: Option<Struct>,
96
97    /// The partition spec for this file, used to distinguish identity transforms
98    /// (which use partition metadata constants) from non-identity transforms like
99    /// bucket/truncate (which must read source columns from the data file).
100    #[builder(default)]
101    partition_spec: Option<Arc<PartitionSpec>>,
102
103    /// Name mapping from table metadata (property: schema.name-mapping.default),
104    /// used to resolve field IDs from column names when Parquet files lack field IDs
105    /// or have field ID conflicts.
106    #[builder(default)]
107    name_mapping: Option<Arc<NameMapping>>,
108
109    /// The unified partition type across all specs in the table.
110    /// When `RESERVED_FIELD_ID_PARTITION` is in the projected field IDs, the reader
111    /// uses this type along with the task's partition_spec and partition data to
112    /// materialize the `_partition` struct column at read time.
113    ///
114    /// This is a table-level value (same for all tasks in a scan), stored per-task
115    /// so that readers are self-contained without needing back-pointers to table
116    /// metadata. The cost is one Arc clone per task.
117    #[builder(default)]
118    unified_partition_type: Option<Arc<StructType>>,
119
120    /// The raw `sort_order_id` recorded on the data file (spec field 140), carried through
121    /// unresolved. See [`DataFile::sort_order_id()`](crate::spec::DataFile::sort_order_id).
122    /// `None` if the file has no sort order id.
123    ///
124    /// This preserves the distinction that [`sort_order`](Self::sort_order) collapses: a
125    /// non-zero id that does not resolve against the table's sort orders still marks the file
126    /// as physically sorted, which a consumer such as the DataFusion sorted-scan optimizer
127    /// needs even when the order definition itself is unavailable.
128    #[builder(default)]
129    sort_order_id: Option<i32>,
130
131    /// The sort order that this file's rows are sorted by, resolved from the data file's
132    /// [`sort_order_id`](Self::sort_order_id) against the table's known sort orders. `Some`
133    /// only when the id resolves to an order that has sort fields. `None` if the file has no
134    /// sort order id, the id does not resolve, or it resolves to an order with no sort fields
135    /// (the reserved unsorted order, id 0 per the spec). Use [`sort_order_id`](Self::sort_order_id)
136    /// to tell those cases apart.
137    #[builder(default)]
138    sort_order: Option<SortOrderRef>,
139
140    /// Whether this scan task should treat column names as case-sensitive when binding predicates.
141    case_sensitive: bool,
142
143    /// Key metadata for encrypted data files (Parquet Modular Encryption).
144    /// When present, the reader uses this to build `FileDecryptionProperties`.
145    ///
146    /// Note on the trust boundary: for the standard encryption scheme this
147    /// carries `StandardKeyMetadata`, whose payload is the *plaintext* DEK.
148    /// Because `FileScanTask` implements [`Serialize`], that plaintext DEK is part
149    /// of the serialized scan plan should these tasks ever be serialized and sent
150    /// over the network.
151    #[builder(default)]
152    key_metadata: Option<Box<[u8]>>,
153}
154
155impl FileScanTask {
156    /// Returns the total size of the data file in bytes.
157    pub fn file_size_in_bytes(&self) -> u64 {
158        self.file_size_in_bytes
159    }
160
161    /// Returns the start offset of the file to scan.
162    pub fn start(&self) -> u64 {
163        self.start
164    }
165
166    /// Returns the length of the file to scan.
167    pub fn length(&self) -> u64 {
168        self.length
169    }
170
171    /// Returns the number of records in the file when the whole file is scanned.
172    pub fn record_count(&self) -> Option<u64> {
173        self.record_count
174    }
175
176    /// Returns the first row id assigned to the data file.
177    pub fn first_row_id(&self) -> Option<i64> {
178        self.first_row_id
179    }
180
181    /// Returns the data sequence number of the file.
182    pub fn data_sequence_number(&self) -> Option<i64> {
183        self.data_sequence_number
184    }
185
186    /// Returns the data file path of this file scan task.
187    pub fn data_file_path(&self) -> &str {
188        &self.data_file_path
189    }
190
191    /// Returns the format of the data file.
192    pub fn data_file_format(&self) -> DataFileFormat {
193        self.data_file_format
194    }
195
196    /// Returns the schema of this file scan task as a reference.
197    pub fn schema(&self) -> &Schema {
198        &self.schema
199    }
200
201    /// Returns the schema of this file scan task as a [`SchemaRef`].
202    pub fn schema_ref(&self) -> SchemaRef {
203        self.schema.clone()
204    }
205
206    /// Returns the project field id of this file scan task.
207    pub fn project_field_ids(&self) -> &[i32] {
208        &self.project_field_ids
209    }
210
211    /// Returns the predicate of this file scan task.
212    pub fn predicate(&self) -> Option<&BoundPredicate> {
213        self.predicate.as_ref()
214    }
215
216    /// Clears the row predicate of this file scan task.
217    ///
218    /// The COW rewrite path uses this after candidate selection: candidates are
219    /// chosen with the predicate during planning, but each chosen file must then
220    /// be read in full so the rewrite sees every surviving row. Delete-file
221    /// application is unaffected.
222    pub(crate) fn clear_predicate(&mut self) {
223        self.predicate = None;
224    }
225
226    /// Returns the delete files that may need to be applied to the data file.
227    pub fn deletes(&self) -> &[FileScanTaskDeleteFile] {
228        &self.deletes
229    }
230
231    /// Returns the partition data from the manifest entry.
232    pub fn partition(&self) -> Option<&Struct> {
233        self.partition.as_ref()
234    }
235
236    /// Returns the partition spec for the data file.
237    pub fn partition_spec(&self) -> Option<&Arc<PartitionSpec>> {
238        self.partition_spec.as_ref()
239    }
240
241    /// Returns the name mapping used to resolve field ids.
242    pub fn name_mapping(&self) -> Option<&Arc<NameMapping>> {
243        self.name_mapping.as_ref()
244    }
245
246    /// Returns the unified partition type across all table partition specs.
247    pub fn unified_partition_type(&self) -> Option<&Arc<StructType>> {
248        self.unified_partition_type.as_ref()
249    }
250
251    /// Returns the raw sort order id recorded on the data file, unresolved.
252    pub fn sort_order_id(&self) -> Option<i32> {
253        self.sort_order_id
254    }
255
256    /// Returns the sort order this file's rows are sorted by, if resolved.
257    pub fn sort_order(&self) -> Option<&SortOrderRef> {
258        self.sort_order.as_ref()
259    }
260
261    /// Returns whether names are treated as case-sensitive.
262    pub fn case_sensitive(&self) -> bool {
263        self.case_sensitive
264    }
265
266    /// Returns the key metadata for the encrypted data file.
267    pub fn key_metadata(&self) -> Option<&[u8]> {
268        self.key_metadata.as_deref()
269    }
270
271    fn validate(&self) -> Result<()> {
272        match (self.partition.as_ref(), self.partition_spec.as_deref()) {
273            (None, None) => Ok(()),
274            (None, Some(partition_spec)) if partition_spec.is_unpartitioned() => Ok(()),
275            (None, Some(_)) => Err(invalid_data!(
276                "FileScanTask with a partitioned spec requires partition values"
277            )),
278            (Some(partition), None) if partition.fields().is_empty() => Ok(()),
279            (Some(_), None) => Err(invalid_data!(
280                "Non-empty FileScanTask partition requires a partition spec"
281            )),
282            (Some(partition), Some(partition_spec))
283                if partition.fields().len() != partition_spec.fields().len() =>
284            {
285                Err(invalid_data!(
286                    "FileScanTask partition has {} fields but partition spec has {} fields",
287                    partition.fields().len(),
288                    partition_spec.fields().len()
289                ))
290            }
291            (Some(_), Some(partition_spec)) => {
292                partition_spec.partition_type(&self.schema)?;
293                Ok(())
294            }
295        }
296    }
297}
298
299impl From<FileScanTask> for Result<FileScanTask> {
300    fn from(task: FileScanTask) -> Self {
301        task.validate()?;
302        Ok(task)
303    }
304}
305
306#[derive(Debug)]
307pub(crate) struct DeleteFileContext {
308    pub(crate) manifest_entry: ManifestEntryRef,
309    pub(crate) partition_spec_id: i32,
310}
311
312impl TryFrom<&DeleteFileContext> for FileScanTaskDeleteFile {
313    type Error = Error;
314
315    fn try_from(ctx: &DeleteFileContext) -> Result<Self> {
316        // The manifest stores these as i64. Convert here so a negative byte offset is
317        // rejected once at the boundary, for every delete-file kind rather than only
318        // deletion vectors, instead of being carried inward and re-checked per consumer.
319        let file_path = ctx.manifest_entry.file_path();
320        let to_offset = |value: Option<i64>, field: &str| -> Result<Option<u64>> {
321            value
322                .map(|value| {
323                    u64::try_from(value).map_err(|_| {
324                        invalid_data!("delete file {file_path} has negative {field} {value}")
325                    })
326                })
327                .transpose()
328        };
329
330        FileScanTaskDeleteFile::builder()
331            .with_file_path(ctx.manifest_entry.file_path().to_string())
332            .with_file_size_in_bytes(ctx.manifest_entry.file_size_in_bytes())
333            .with_file_type(ctx.manifest_entry.content_type())
334            .with_file_format(ctx.manifest_entry.data_file().file_format())
335            .with_partition_spec_id(ctx.partition_spec_id)
336            .with_equality_ids(ctx.manifest_entry.data_file.equality_ids.clone())
337            .with_referenced_data_file(ctx.manifest_entry.data_file.referenced_data_file.clone())
338            .with_content_offset(to_offset(
339                ctx.manifest_entry.data_file.content_offset,
340                "content_offset",
341            )?)
342            .with_content_size_in_bytes(to_offset(
343                ctx.manifest_entry.data_file.content_size_in_bytes,
344                "content_size_in_bytes",
345            )?)
346            .with_record_count(Some(ctx.manifest_entry.record_count()))
347            .with_key_metadata(
348                ctx.manifest_entry
349                    .data_file
350                    .key_metadata
351                    .as_deref()
352                    .map(Box::from),
353            )
354            .build()
355    }
356}
357
358/// A task to scan part of file.
359#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, TypedBuilder)]
360#[serde(try_from = "crate::scan::task::_serde::FileScanTaskDeleteFileSerde")]
361#[builder(
362    field_defaults(setter(prefix = "with_")),
363    build_method(into = Result<FileScanTaskDeleteFile>)
364)]
365pub struct FileScanTaskDeleteFile {
366    /// The delete file path
367    file_path: String,
368
369    /// The total size of the delete file in bytes, from the manifest entry.
370    file_size_in_bytes: u64,
371
372    /// delete file type
373    file_type: DataContentType,
374
375    /// The delete file's format, from the manifest entry. A `PositionDeletes` entry written as
376    /// `Puffin` is a V3 deletion vector; one written as `Parquet` is a position delete file.
377    file_format: DataFileFormat,
378
379    /// partition id
380    partition_spec_id: i32,
381
382    /// equality ids for equality deletes (null for anything other than equality-deletes)
383    #[builder(default)]
384    equality_ids: Option<Vec<i32>>,
385
386    /// For a deletion vector, the location of the data file whose rows it deletes. Required for
387    /// deletion vectors, and may also be set on a position delete file scoped to one data file.
388    #[serde(default)]
389    #[serde(skip_serializing_if = "Option::is_none")]
390    #[builder(default)]
391    referenced_data_file: Option<String>,
392
393    /// For a deletion vector, the offset of the blob within its Puffin file. Set only for
394    /// deletion vectors, where it locates the blob for direct access. Unsigned because it is a
395    /// byte offset: a negative value is rejected once, when the task is built from a manifest
396    /// entry, rather than re-checked by every consumer.
397    #[serde(default)]
398    #[serde(skip_serializing_if = "Option::is_none")]
399    #[builder(default)]
400    content_offset: Option<u64>,
401
402    /// For a deletion vector, the length in bytes of the blob within its Puffin file.
403    /// Required together with `content_offset`; both are absent for non-DV delete files.
404    /// Unsigned for the same reason as `content_offset`.
405    #[serde(default)]
406    #[serde(skip_serializing_if = "Option::is_none")]
407    #[builder(default)]
408    content_size_in_bytes: Option<u64>,
409
410    /// The number of records in the delete file, from the manifest entry; for a deletion vector,
411    /// the cardinality of its bitmap. `None` only for a task not built from a manifest entry.
412    #[serde(default)]
413    #[serde(skip_serializing_if = "Option::is_none")]
414    #[builder(default)]
415    record_count: Option<u64>,
416
417    /// Key metadata for an encrypted delete file. When present, the reader uses this to
418    /// decrypt the file: for a Parquet equality or position delete file, this builds
419    /// `FileDecryptionProperties` (Parquet Modular Encryption); for a deletion vector, whose
420    /// Puffin file has no native encryption, this wraps the range read in an
421    /// `EncryptedInputFile` (AGS1 stream encryption).
422    ///
423    /// Same plaintext-DEK trust boundary as [`FileScanTask::key_metadata`]:
424    /// this is serialized into the scan plan and crosses the planner -> worker
425    /// channel in the clear for the standard encryption scheme.
426    #[serde(default)]
427    #[serde(skip_serializing_if = "Option::is_none")]
428    #[builder(default)]
429    key_metadata: Option<Box<[u8]>>,
430}
431
432mod _serde {
433    use std::sync::Arc;
434
435    use serde::{Deserialize, Serialize};
436
437    use super::{FileScanTask, FileScanTaskDeleteFile};
438    use crate::error::invalid_data;
439    use crate::expr::BoundPredicate;
440    use crate::spec::{
441        DataContentType, DataFileFormat, Literal, NameMapping, PartitionSpec, RawLiteral,
442        SchemaRef, SortOrderRef, StructType, Type,
443    };
444    use crate::{Error, Result};
445
446    #[derive(Deserialize)]
447    pub(super) struct FileScanTaskSerde {
448        file_size_in_bytes: u64,
449        start: u64,
450        length: u64,
451        record_count: Option<u64>,
452        first_row_id: Option<i64>,
453        data_sequence_number: Option<i64>,
454        data_file_path: String,
455        data_file_format: DataFileFormat,
456        schema: SchemaRef,
457        project_field_ids: Vec<i32>,
458        predicate: Option<BoundPredicate>,
459        deletes: Vec<FileScanTaskDeleteFile>,
460        #[serde(default)]
461        partition: Option<RawLiteral>,
462        #[serde(default)]
463        partition_spec: Option<Arc<PartitionSpec>>,
464        #[serde(default)]
465        name_mapping: Option<Arc<NameMapping>>,
466        #[serde(default)]
467        unified_partition_type: Option<Arc<StructType>>,
468        #[serde(default)]
469        sort_order_id: Option<i32>,
470        #[serde(default)]
471        sort_order: Option<SortOrderRef>,
472        case_sensitive: bool,
473        #[serde(default)]
474        key_metadata: Option<Box<[u8]>>,
475    }
476
477    #[derive(Serialize)]
478    struct FileScanTaskRefSerde<'a> {
479        file_size_in_bytes: u64,
480        start: u64,
481        length: u64,
482        #[serde(skip_serializing_if = "Option::is_none")]
483        record_count: Option<u64>,
484        #[serde(skip_serializing_if = "Option::is_none")]
485        first_row_id: Option<i64>,
486        #[serde(skip_serializing_if = "Option::is_none")]
487        data_sequence_number: Option<i64>,
488        data_file_path: &'a str,
489        data_file_format: DataFileFormat,
490        schema: &'a SchemaRef,
491        project_field_ids: &'a [i32],
492        #[serde(skip_serializing_if = "Option::is_none")]
493        predicate: Option<&'a BoundPredicate>,
494        deletes: &'a [FileScanTaskDeleteFile],
495        #[serde(skip_serializing_if = "Option::is_none")]
496        partition: Option<RawLiteral>,
497        #[serde(skip_serializing_if = "Option::is_none")]
498        partition_spec: Option<&'a Arc<PartitionSpec>>,
499        #[serde(skip_serializing_if = "Option::is_none")]
500        name_mapping: Option<&'a Arc<NameMapping>>,
501        #[serde(skip_serializing_if = "Option::is_none")]
502        unified_partition_type: Option<&'a Arc<StructType>>,
503        #[serde(skip_serializing_if = "Option::is_none")]
504        sort_order_id: Option<i32>,
505        #[serde(skip_serializing_if = "Option::is_none")]
506        sort_order: Option<&'a SortOrderRef>,
507        case_sensitive: bool,
508        #[serde(skip_serializing_if = "Option::is_none")]
509        key_metadata: Option<&'a [u8]>,
510    }
511
512    fn partition_type(
513        partition_spec: Option<&PartitionSpec>,
514        schema: &crate::spec::Schema,
515    ) -> Result<Type> {
516        let partition_type = match partition_spec {
517            Some(partition_spec) => partition_spec.partition_type(schema)?,
518            None => PartitionSpec::unpartition_spec().partition_type(schema)?,
519        };
520        Ok(Type::Struct(partition_type))
521    }
522
523    impl<'a> TryFrom<&'a FileScanTask> for FileScanTaskRefSerde<'a> {
524        type Error = Error;
525
526        fn try_from(value: &'a FileScanTask) -> Result<Self> {
527            let partition = value
528                .partition
529                .as_ref()
530                .map(|partition| {
531                    let partition_type =
532                        partition_type(value.partition_spec.as_deref(), &value.schema)?;
533                    RawLiteral::try_from(Literal::Struct(partition.clone()), &partition_type)
534                })
535                .transpose()?;
536
537            Ok(Self {
538                file_size_in_bytes: value.file_size_in_bytes,
539                start: value.start,
540                length: value.length,
541                record_count: value.record_count,
542                first_row_id: value.first_row_id,
543                data_sequence_number: value.data_sequence_number,
544                data_file_path: &value.data_file_path,
545                data_file_format: value.data_file_format,
546                schema: &value.schema,
547                project_field_ids: &value.project_field_ids,
548                predicate: value.predicate.as_ref(),
549                deletes: &value.deletes,
550                partition,
551                partition_spec: value.partition_spec.as_ref(),
552                name_mapping: value.name_mapping.as_ref(),
553                unified_partition_type: value.unified_partition_type.as_ref(),
554                sort_order_id: value.sort_order_id,
555                sort_order: value.sort_order.as_ref(),
556                case_sensitive: value.case_sensitive,
557                key_metadata: value.key_metadata.as_deref(),
558            })
559        }
560    }
561
562    impl Serialize for FileScanTask {
563        fn serialize<S>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error>
564        where S: serde::Serializer {
565            FileScanTaskRefSerde::try_from(self)
566                .map_err(serde::ser::Error::custom)?
567                .serialize(serializer)
568        }
569    }
570
571    impl TryFrom<FileScanTaskSerde> for FileScanTask {
572        type Error = Error;
573
574        fn try_from(value: FileScanTaskSerde) -> Result<Self> {
575            let partition = value
576                .partition
577                .map(|partition| {
578                    let partition_type =
579                        partition_type(value.partition_spec.as_deref(), &value.schema)?;
580                    match partition.try_into(&partition_type)? {
581                        Some(Literal::Struct(partition)) => Ok(partition),
582                        _ => Err(invalid_data!("FileScanTask partition must be a struct")),
583                    }
584                })
585                .transpose()?;
586
587            Self::builder()
588                .with_file_size_in_bytes(value.file_size_in_bytes)
589                .with_start(value.start)
590                .with_length(value.length)
591                .with_record_count(value.record_count)
592                .with_first_row_id(value.first_row_id)
593                .with_data_sequence_number(value.data_sequence_number)
594                .with_data_file_path(value.data_file_path)
595                .with_data_file_format(value.data_file_format)
596                .with_schema(value.schema)
597                .with_project_field_ids(value.project_field_ids)
598                .with_predicate(value.predicate)
599                .with_deletes(value.deletes)
600                .with_partition(partition)
601                .with_partition_spec(value.partition_spec)
602                .with_name_mapping(value.name_mapping)
603                .with_unified_partition_type(value.unified_partition_type)
604                .with_sort_order_id(value.sort_order_id)
605                .with_sort_order(value.sort_order)
606                .with_case_sensitive(value.case_sensitive)
607                .with_key_metadata(value.key_metadata)
608                .build()
609        }
610    }
611
612    #[derive(Deserialize)]
613    pub(super) struct FileScanTaskDeleteFileSerde {
614        file_path: String,
615        file_size_in_bytes: u64,
616        file_type: DataContentType,
617        file_format: DataFileFormat,
618        partition_spec_id: i32,
619        equality_ids: Option<Vec<i32>>,
620        #[serde(default)]
621        referenced_data_file: Option<String>,
622        #[serde(default)]
623        content_offset: Option<u64>,
624        #[serde(default)]
625        content_size_in_bytes: Option<u64>,
626        #[serde(default)]
627        record_count: Option<u64>,
628        #[serde(default)]
629        key_metadata: Option<Box<[u8]>>,
630    }
631
632    impl TryFrom<FileScanTaskDeleteFileSerde> for FileScanTaskDeleteFile {
633        type Error = Error;
634
635        fn try_from(value: FileScanTaskDeleteFileSerde) -> Result<Self> {
636            Self::builder()
637                .with_file_path(value.file_path)
638                .with_file_size_in_bytes(value.file_size_in_bytes)
639                .with_file_type(value.file_type)
640                .with_file_format(value.file_format)
641                .with_partition_spec_id(value.partition_spec_id)
642                .with_equality_ids(value.equality_ids)
643                .with_referenced_data_file(value.referenced_data_file)
644                .with_content_offset(value.content_offset)
645                .with_content_size_in_bytes(value.content_size_in_bytes)
646                .with_record_count(value.record_count)
647                .with_key_metadata(value.key_metadata)
648                .build()
649        }
650    }
651}
652
653impl FileScanTaskDeleteFile {
654    /// Returns the delete file path.
655    pub fn file_path(&self) -> &str {
656        &self.file_path
657    }
658
659    /// Returns the total size of the delete file in bytes.
660    pub fn file_size_in_bytes(&self) -> u64 {
661        self.file_size_in_bytes
662    }
663
664    /// Returns the delete file content type.
665    pub fn file_type(&self) -> DataContentType {
666        self.file_type
667    }
668
669    /// Returns the delete file format.
670    pub fn file_format(&self) -> DataFileFormat {
671        self.file_format
672    }
673
674    /// Returns the partition spec id.
675    pub fn partition_spec_id(&self) -> i32 {
676        self.partition_spec_id
677    }
678
679    /// Returns the equality field ids for an equality delete file.
680    pub fn equality_ids(&self) -> Option<&[i32]> {
681        self.equality_ids.as_deref()
682    }
683
684    /// Returns the referenced data file path.
685    pub fn referenced_data_file(&self) -> Option<&str> {
686        self.referenced_data_file.as_deref()
687    }
688
689    /// Returns the deletion vector blob offset.
690    pub fn content_offset(&self) -> Option<u64> {
691        self.content_offset
692    }
693
694    /// Returns the deletion vector blob size in bytes.
695    pub fn content_size_in_bytes(&self) -> Option<u64> {
696        self.content_size_in_bytes
697    }
698
699    /// Returns the number of records in the delete file.
700    pub fn record_count(&self) -> Option<u64> {
701        self.record_count
702    }
703
704    /// Returns the key metadata for the encrypted delete file.
705    pub fn key_metadata(&self) -> Option<&[u8]> {
706        self.key_metadata.as_deref()
707    }
708
709    fn is_deletion_vector(&self) -> bool {
710        self.file_type == DataContentType::PositionDeletes
711            && self.file_format == DataFileFormat::Puffin
712    }
713
714    fn validate(&self) -> Result<()> {
715        // Negative offsets/sizes are unrepresentable: both fields are u64, and the
716        // conversion from the manifest entry's i64 rejects a negative there.
717
718        // Spec: `content_size_in_bytes` is required whenever `content_offset` is
719        // present, for every delete-file kind (not only deletion vectors).
720        if self.content_offset.is_some() && self.content_size_in_bytes.is_none() {
721            return Err(invalid_data!(
722                "delete file {} is missing content_size_in_bytes while content_offset is present",
723                self.file_path
724            ));
725        }
726
727        // Presence of the deletion-vector coordinates is required only for
728        // deletion vectors; other delete files legitimately carry none of these.
729        if !self.is_deletion_vector() {
730            return Ok(());
731        }
732
733        let missing = if self.referenced_data_file.is_none() {
734            Some("referenced_data_file")
735        } else if self.content_offset.is_none() {
736            Some("content_offset")
737        } else if self.content_size_in_bytes.is_none() {
738            Some("content_size_in_bytes")
739        } else if self.record_count.is_none() {
740            Some("record_count")
741        } else {
742            None
743        };
744
745        if let Some(field) = missing {
746            return Err(invalid_data!(
747                "deletion vector {} is missing {field}",
748                self.file_path
749            ));
750        }
751
752        Ok(())
753    }
754}
755
756impl From<FileScanTaskDeleteFile> for Result<FileScanTaskDeleteFile> {
757    fn from(task: FileScanTaskDeleteFile) -> Self {
758        task.validate()?;
759        Ok(task)
760    }
761}
762
763#[cfg(test)]
764mod tests {
765    use std::collections::HashMap;
766
767    use super::*;
768    use crate::ErrorKind;
769    use crate::spec::{
770        DataFileBuilder, Literal, ManifestEntry, ManifestStatus, NestedField, PrimitiveType,
771        Transform, Type,
772    };
773
774    fn build_file_scan_task(
775        schema: SchemaRef,
776        partition: Option<Struct>,
777        partition_spec: Option<Arc<PartitionSpec>>,
778    ) -> Result<FileScanTask> {
779        FileScanTask::builder()
780            .with_file_size_in_bytes(100)
781            .with_start(0)
782            .with_length(100)
783            .with_data_file_path("data_file_path".to_string())
784            .with_data_file_format(DataFileFormat::Parquet)
785            .with_schema(schema)
786            .with_project_field_ids(vec![])
787            .with_partition(partition)
788            .with_partition_spec(partition_spec)
789            .with_case_sensitive(false)
790            .build()
791    }
792
793    fn schema_and_spec(
794        primitive_type: PrimitiveType,
795        transform: Transform,
796    ) -> (SchemaRef, Arc<PartitionSpec>) {
797        let schema = Arc::new(
798            Schema::builder()
799                .with_fields(vec![Arc::new(NestedField::required(
800                    1,
801                    "x",
802                    Type::Primitive(primitive_type),
803                ))])
804                .build()
805                .unwrap(),
806        );
807        let partition_spec = Arc::new(
808            PartitionSpec::builder(schema.clone())
809                .add_partition_field("x", "x_partition", transform)
810                .unwrap()
811                .build()
812                .unwrap(),
813        );
814        (schema, partition_spec)
815    }
816
817    fn build_delete_file_task(
818        file_type: DataContentType,
819        file_format: DataFileFormat,
820    ) -> Result<FileScanTaskDeleteFile> {
821        FileScanTaskDeleteFile::builder()
822            .with_file_path("delete-file".to_string())
823            .with_file_size_in_bytes(100)
824            .with_file_type(file_type)
825            .with_file_format(file_format)
826            .with_partition_spec_id(0)
827            .build()
828    }
829
830    fn assert_delete_file_builder_error(
831        result: Result<FileScanTaskDeleteFile>,
832        expected_message: &str,
833    ) {
834        match result {
835            Ok(task) => panic!(
836                "expected delete file builder to fail with `{expected_message}`, but got Ok({task:?})"
837            ),
838            Err(err) => {
839                assert_eq!(err.kind(), ErrorKind::DataInvalid);
840                assert_eq!(err.message(), expected_message);
841            }
842        }
843    }
844
845    #[test]
846    fn test_file_scan_task_builder_rejects_non_empty_partition_without_spec() {
847        // Regression test for https://github.com/apache/iceberg-rust/issues/3130.
848        let err = build_file_scan_task(
849            Arc::new(Schema::builder().build().unwrap()),
850            Some(Struct::from_iter([Some(Literal::long(42))])),
851            None,
852        )
853        .unwrap_err();
854
855        assert_eq!(err.kind(), ErrorKind::DataInvalid);
856        assert_eq!(
857            err.message(),
858            "Non-empty FileScanTask partition requires a partition spec"
859        );
860    }
861
862    #[test]
863    fn test_file_scan_task_builder_accepts_empty_partition_without_spec() {
864        build_file_scan_task(
865            Arc::new(Schema::builder().build().unwrap()),
866            Some(Struct::empty()),
867            None,
868        )
869        .unwrap();
870    }
871
872    #[test]
873    fn test_file_scan_task_builder_rejects_partitioned_spec_without_partition() {
874        let (schema, partition_spec) = schema_and_spec(PrimitiveType::Long, Transform::Identity);
875
876        let err = build_file_scan_task(schema, None, Some(partition_spec)).unwrap_err();
877
878        assert_eq!(err.kind(), ErrorKind::DataInvalid);
879        assert_eq!(
880            err.message(),
881            "FileScanTask with a partitioned spec requires partition values"
882        );
883    }
884
885    #[test]
886    fn test_file_scan_task_builder_accepts_unpartitioned_spec_without_partition() {
887        build_file_scan_task(
888            Arc::new(Schema::builder().build().unwrap()),
889            None,
890            Some(Arc::new(PartitionSpec::unpartition_spec())),
891        )
892        .unwrap();
893    }
894
895    #[test]
896    fn test_file_scan_task_builder_rejects_partition_arity_mismatch() {
897        let (schema, partition_spec) = schema_and_spec(PrimitiveType::Long, Transform::Identity);
898
899        let err =
900            build_file_scan_task(schema, Some(Struct::empty()), Some(partition_spec)).unwrap_err();
901
902        assert_eq!(err.kind(), ErrorKind::DataInvalid);
903        assert!(err.message().contains("partition has 0 fields"));
904        assert!(err.message().contains("partition spec has 1 fields"));
905    }
906
907    #[test]
908    fn test_file_scan_task_builder_accepts_dropped_partition_source_column() {
909        let (_historical_schema, partition_spec) =
910            schema_and_spec(PrimitiveType::Long, Transform::Identity);
911        let current_schema = Arc::new(
912            Schema::builder()
913                .with_fields(vec![Arc::new(NestedField::required(
914                    2,
915                    "y",
916                    Type::Primitive(PrimitiveType::String),
917                ))])
918                .build()
919                .unwrap(),
920        );
921
922        build_file_scan_task(
923            current_schema,
924            Some(Struct::from_iter([Some(Literal::long(42))])),
925            Some(partition_spec),
926        )
927        .unwrap();
928    }
929
930    #[test]
931    fn test_file_scan_task_builder_rejects_partition_spec_incompatible_with_schema() {
932        let (_historical_schema, partition_spec) =
933            schema_and_spec(PrimitiveType::Timestamp, Transform::Day);
934        let current_schema = Arc::new(
935            Schema::builder()
936                .with_fields(vec![Arc::new(NestedField::required(
937                    1,
938                    "x",
939                    Type::Primitive(PrimitiveType::String),
940                ))])
941                .build()
942                .unwrap(),
943        );
944
945        let err = build_file_scan_task(
946            current_schema,
947            Some(Struct::from_iter([Some(Literal::date(20_000))])),
948            Some(partition_spec),
949        )
950        .unwrap_err();
951
952        assert_eq!(err.kind(), ErrorKind::DataInvalid);
953    }
954
955    #[test]
956    fn test_delete_file_builder_accepts_valid_deletion_vector() {
957        let task = FileScanTaskDeleteFile::builder()
958            .with_file_path("dv.puffin".to_string())
959            .with_file_size_in_bytes(100)
960            .with_file_type(DataContentType::PositionDeletes)
961            .with_file_format(DataFileFormat::Puffin)
962            .with_partition_spec_id(7)
963            .with_referenced_data_file(Some("data.parquet".to_string()))
964            .with_content_offset(Some(11))
965            .with_content_size_in_bytes(Some(13))
966            .with_record_count(Some(3))
967            .with_key_metadata(Some(vec![17, 19].into_boxed_slice()))
968            .build()
969            .unwrap();
970
971        assert_eq!(task.file_path(), "dv.puffin");
972        assert_eq!(task.file_size_in_bytes(), 100);
973        assert_eq!(task.file_type(), DataContentType::PositionDeletes);
974        assert_eq!(task.file_format(), DataFileFormat::Puffin);
975        assert_eq!(task.partition_spec_id(), 7);
976        assert_eq!(task.equality_ids(), None);
977        assert_eq!(task.referenced_data_file(), Some("data.parquet"));
978        assert_eq!(task.content_offset(), Some(11));
979        assert_eq!(task.content_size_in_bytes(), Some(13));
980        assert_eq!(task.record_count(), Some(3));
981        assert_eq!(task.key_metadata(), Some([17, 19].as_slice()));
982    }
983
984    #[test]
985    fn test_delete_file_builder_rejects_dv_missing_referenced_data_file() {
986        assert_delete_file_builder_error(
987            FileScanTaskDeleteFile::builder()
988                .with_file_path("dv.puffin".to_string())
989                .with_file_size_in_bytes(100)
990                .with_file_type(DataContentType::PositionDeletes)
991                .with_file_format(DataFileFormat::Puffin)
992                .with_partition_spec_id(0)
993                .with_content_offset(Some(7))
994                .with_content_size_in_bytes(Some(11))
995                .with_record_count(Some(3))
996                .build(),
997            "deletion vector dv.puffin is missing referenced_data_file",
998        );
999    }
1000
1001    #[test]
1002    fn test_delete_file_builder_rejects_dv_missing_content_offset() {
1003        assert_delete_file_builder_error(
1004            FileScanTaskDeleteFile::builder()
1005                .with_file_path("dv.puffin".to_string())
1006                .with_file_size_in_bytes(100)
1007                .with_file_type(DataContentType::PositionDeletes)
1008                .with_file_format(DataFileFormat::Puffin)
1009                .with_partition_spec_id(0)
1010                .with_referenced_data_file(Some("data.parquet".to_string()))
1011                .with_content_size_in_bytes(Some(11))
1012                .with_record_count(Some(3))
1013                .build(),
1014            "deletion vector dv.puffin is missing content_offset",
1015        );
1016    }
1017
1018    #[test]
1019    fn test_delete_file_builder_rejects_dv_missing_content_size() {
1020        assert_delete_file_builder_error(
1021            FileScanTaskDeleteFile::builder()
1022                .with_file_path("dv.puffin".to_string())
1023                .with_file_size_in_bytes(100)
1024                .with_file_type(DataContentType::PositionDeletes)
1025                .with_file_format(DataFileFormat::Puffin)
1026                .with_partition_spec_id(0)
1027                .with_referenced_data_file(Some("data.parquet".to_string()))
1028                .with_content_offset(Some(7))
1029                .with_record_count(Some(3))
1030                .build(),
1031            "delete file dv.puffin is missing content_size_in_bytes while content_offset is present",
1032        );
1033    }
1034
1035    #[test]
1036    fn test_delete_file_builder_rejects_dv_missing_record_count() {
1037        assert_delete_file_builder_error(
1038            FileScanTaskDeleteFile::builder()
1039                .with_file_path("dv.puffin".to_string())
1040                .with_file_size_in_bytes(100)
1041                .with_file_type(DataContentType::PositionDeletes)
1042                .with_file_format(DataFileFormat::Puffin)
1043                .with_partition_spec_id(0)
1044                .with_referenced_data_file(Some("data.parquet".to_string()))
1045                .with_content_offset(Some(7))
1046                .with_content_size_in_bytes(Some(11))
1047                .build(),
1048            "deletion vector dv.puffin is missing record_count",
1049        );
1050    }
1051
1052    #[test]
1053    fn test_delete_file_builder_rejects_non_dv_offset_without_size() {
1054        // Spec pairing rule applies to every delete-file kind: a Parquet
1055        // position-delete file that carries a content_offset must also carry a
1056        // content_size_in_bytes, even though it is not a deletion vector.
1057        assert_delete_file_builder_error(
1058            FileScanTaskDeleteFile::builder()
1059                .with_file_path("position-deletes.parquet".to_string())
1060                .with_file_size_in_bytes(100)
1061                .with_file_type(DataContentType::PositionDeletes)
1062                .with_file_format(DataFileFormat::Parquet)
1063                .with_partition_spec_id(0)
1064                .with_content_offset(Some(7))
1065                .build(),
1066            "delete file position-deletes.parquet is missing content_size_in_bytes while content_offset is present",
1067        );
1068    }
1069
1070    fn delete_manifest_entry(
1071        file_path: &str,
1072        file_format: DataFileFormat,
1073        content_offset: Option<i64>,
1074        content_size_in_bytes: Option<i64>,
1075    ) -> DeleteFileContext {
1076        let mut data_file = DataFileBuilder::default()
1077            .content(DataContentType::PositionDeletes)
1078            .file_path(file_path.to_string())
1079            .file_format(file_format)
1080            .partition(Struct::empty())
1081            .record_count(3)
1082            .file_size_in_bytes(100)
1083            .column_sizes(HashMap::new())
1084            .value_counts(HashMap::new())
1085            .null_value_counts(HashMap::new())
1086            .partition_spec_id(0)
1087            .referenced_data_file(Some("data.parquet".to_string()))
1088            .build()
1089            .unwrap();
1090        data_file.content_offset = content_offset;
1091        data_file.content_size_in_bytes = content_size_in_bytes;
1092
1093        DeleteFileContext {
1094            manifest_entry: Arc::new(
1095                ManifestEntry::builder()
1096                    .status(ManifestStatus::Added)
1097                    .data_file(data_file)
1098                    .build(),
1099            ),
1100            partition_spec_id: 0,
1101        }
1102    }
1103
1104    fn assert_delete_context_error(ctx: &DeleteFileContext, expected_message: &str) {
1105        match FileScanTaskDeleteFile::try_from(ctx) {
1106            Ok(task) => panic!(
1107                "expected conversion to fail with `{expected_message}`, but got Ok({task:?})"
1108            ),
1109            Err(error) => {
1110                assert_eq!(error.kind(), ErrorKind::DataInvalid);
1111                assert!(
1112                    error.message().contains(expected_message),
1113                    "expected `{expected_message}`, got `{}`",
1114                    error.message()
1115                );
1116            }
1117        }
1118    }
1119
1120    /// The manifest entry supplies `content_offset` / `content_size_in_bytes` as
1121    /// `i64`; the conversion into a task rejects a negative value as it enters,
1122    /// since the task fields are `u64`.
1123    #[test]
1124    fn test_delete_context_rejects_negative_dv_offset() {
1125        assert_delete_context_error(
1126            &delete_manifest_entry("dv.puffin", DataFileFormat::Puffin, Some(-1), Some(11)),
1127            "delete file dv.puffin has negative content_offset -1",
1128        );
1129    }
1130
1131    #[test]
1132    fn test_delete_context_rejects_negative_dv_size() {
1133        assert_delete_context_error(
1134            &delete_manifest_entry("dv.puffin", DataFileFormat::Puffin, Some(7), Some(-1)),
1135            "delete file dv.puffin has negative content_size_in_bytes -1",
1136        );
1137    }
1138
1139    /// The negative-coordinate check applies to every delete-file kind, including
1140    /// a Parquet position-delete file rather than only deletion vectors.
1141    #[test]
1142    fn test_delete_context_rejects_negative_coordinates_for_non_dv() {
1143        assert_delete_context_error(
1144            &delete_manifest_entry(
1145                "position-deletes.parquet",
1146                DataFileFormat::Parquet,
1147                Some(-1),
1148                None,
1149            ),
1150            "delete file position-deletes.parquet has negative content_offset -1",
1151        );
1152        assert_delete_context_error(
1153            &delete_manifest_entry(
1154                "position-deletes.parquet",
1155                DataFileFormat::Parquet,
1156                None,
1157                Some(-1),
1158            ),
1159            "delete file position-deletes.parquet has negative content_size_in_bytes -1",
1160        );
1161    }
1162
1163    #[test]
1164    fn test_delete_file_builder_accepts_non_dv_delete_without_dv_fields() {
1165        build_delete_file_task(DataContentType::PositionDeletes, DataFileFormat::Parquet).unwrap();
1166        build_delete_file_task(DataContentType::EqualityDeletes, DataFileFormat::Parquet).unwrap();
1167    }
1168}