Skip to main content

iceberg/spec/manifest/
writer.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::cmp::min;
19
20use apache_avro::{Writer as AvroWriter, to_value};
21use bytes::Bytes;
22use itertools::Itertools;
23use serde_json::to_vec;
24
25use super::{
26    Datum, FormatVersion, ManifestContentType, PartitionSpec, PrimitiveType,
27    UNASSIGNED_SEQUENCE_NUMBER,
28};
29use crate::encryption::EncryptedOutputFile;
30use crate::error::{Result, invalid_data};
31use crate::io::{FileMetadata, FileWrite, OutputFile};
32use crate::spec::manifest::_serde::{ManifestEntryV1, ManifestEntryV2};
33use crate::spec::manifest::{manifest_schema_v1, manifest_schema_v2};
34use crate::spec::{
35    DataContentType, DataFile, FieldSummary, ManifestEntry, ManifestFile, ManifestMetadata,
36    ManifestStatus, PrimitiveLiteral, SchemaRef, StructType, Type,
37};
38
39/// Placeholder for snapshot ID. The field with this value must be replaced
40/// with the actual snapshot ID before it is committed.
41const UNASSIGNED_SNAPSHOT_ID: i64 = -1;
42
43/// Retains the output until close provides the size required by encrypted key metadata.
44pub(crate) enum ManifestOutput {
45    Plain(OutputFile),
46    Encrypted(EncryptedOutputFile),
47}
48
49impl ManifestOutput {
50    async fn writer(&self) -> Result<Box<dyn FileWrite>> {
51        match self {
52            Self::Plain(output) => output.writer().await,
53            Self::Encrypted(output) => output.writer().await,
54        }
55    }
56
57    fn encoded_key_metadata(&self, file_metadata: &FileMetadata) -> Result<Option<Vec<u8>>> {
58        match self {
59            Self::Plain(_) => Ok(None),
60            Self::Encrypted(output) => Ok(Some(
61                output
62                    .key_metadata_with_saved_file_metadata(file_metadata)
63                    .encode()?
64                    .into_vec(),
65            )),
66        }
67    }
68}
69
70/// The builder used to create a [`ManifestWriter`].
71pub struct ManifestWriterBuilder {
72    output: ManifestOutput,
73    location: String,
74    snapshot_id: Option<i64>,
75    schema: SchemaRef,
76    partition_spec: PartitionSpec,
77}
78
79impl ManifestWriterBuilder {
80    /// Create a new builder for unencrypted manifests.
81    pub fn new(
82        output: OutputFile,
83        snapshot_id: Option<i64>,
84        schema: SchemaRef,
85        partition_spec: PartitionSpec,
86    ) -> Self {
87        let location = output.location().to_owned();
88        Self {
89            output: ManifestOutput::Plain(output),
90            location,
91            snapshot_id,
92            schema,
93            partition_spec,
94        }
95    }
96
97    /// Create a new builder from an [`EncryptedOutputFile`].
98    ///
99    /// Use this when writing manifests with transparent encryption.
100    pub fn new_from_encrypted(
101        encrypted_output: EncryptedOutputFile,
102        snapshot_id: Option<i64>,
103        schema: SchemaRef,
104        partition_spec: PartitionSpec,
105    ) -> Result<Self> {
106        let location = encrypted_output.location().to_owned();
107        Ok(Self {
108            output: ManifestOutput::Encrypted(encrypted_output),
109            location,
110            snapshot_id,
111            schema,
112            partition_spec,
113        })
114    }
115
116    /// Build a [`ManifestWriter`] for format version 1.
117    pub fn build_v1(self) -> ManifestWriter {
118        let metadata = ManifestMetadata::builder()
119            .schema_id(self.schema.schema_id())
120            .schema(self.schema)
121            .partition_spec(self.partition_spec)
122            .format_version(FormatVersion::V1)
123            .content(ManifestContentType::Data)
124            .build();
125        ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
126    }
127
128    /// Build a [`ManifestWriter`] for format version 2, data content.
129    pub fn build_v2_data(self) -> ManifestWriter {
130        let metadata = ManifestMetadata::builder()
131            .schema_id(self.schema.schema_id())
132            .schema(self.schema)
133            .partition_spec(self.partition_spec)
134            .format_version(FormatVersion::V2)
135            .content(ManifestContentType::Data)
136            .build();
137
138        ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
139    }
140
141    /// Build a [`ManifestWriter`] for format version 2, deletes content.
142    pub fn build_v2_deletes(self) -> ManifestWriter {
143        let metadata = ManifestMetadata::builder()
144            .schema_id(self.schema.schema_id())
145            .schema(self.schema)
146            .partition_spec(self.partition_spec)
147            .format_version(FormatVersion::V2)
148            .content(ManifestContentType::Deletes)
149            .build();
150        ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
151    }
152
153    /// Build a [`ManifestWriter`] for format version 2, data content.
154    pub fn build_v3_data(self) -> ManifestWriter {
155        let metadata = ManifestMetadata::builder()
156            .schema_id(self.schema.schema_id())
157            .schema(self.schema)
158            .partition_spec(self.partition_spec)
159            .format_version(FormatVersion::V3)
160            .content(ManifestContentType::Data)
161            .build();
162        ManifestWriter::new(
163            self.output,
164            self.location,
165            self.snapshot_id,
166            metadata,
167            // First row id is assigned by the [`ManifestListWriter`] when the manifest
168            // is added to the list.
169            None,
170        )
171    }
172
173    /// Build a [`ManifestWriter`] for format version 3, deletes content.
174    pub fn build_v3_deletes(self) -> ManifestWriter {
175        let metadata = ManifestMetadata::builder()
176            .schema_id(self.schema.schema_id())
177            .schema(self.schema)
178            .partition_spec(self.partition_spec)
179            .format_version(FormatVersion::V3)
180            .content(ManifestContentType::Deletes)
181            .build();
182        ManifestWriter::new(self.output, self.location, self.snapshot_id, metadata, None)
183    }
184}
185
186/// A manifest writer.
187pub struct ManifestWriter {
188    output: ManifestOutput,
189    location: String,
190
191    snapshot_id: Option<i64>,
192
193    added_files: u32,
194    added_rows: u64,
195    existing_files: u32,
196    existing_rows: u64,
197    deleted_files: u32,
198    deleted_rows: u64,
199    first_row_id: Option<u64>,
200
201    min_seq_num: Option<i64>,
202
203    manifest_entries: Vec<ManifestEntry>,
204
205    metadata: ManifestMetadata,
206}
207
208impl ManifestWriter {
209    /// Create a new manifest writer.
210    pub(crate) fn new(
211        output: ManifestOutput,
212        location: String,
213        snapshot_id: Option<i64>,
214        metadata: ManifestMetadata,
215        first_row_id: Option<u64>,
216    ) -> Self {
217        Self {
218            output,
219            location,
220            snapshot_id,
221            added_files: 0,
222            added_rows: 0,
223            existing_files: 0,
224            existing_rows: 0,
225            deleted_files: 0,
226            deleted_rows: 0,
227            first_row_id,
228            min_seq_num: None,
229            manifest_entries: Vec::new(),
230            metadata,
231        }
232    }
233
234    fn construct_partition_summaries(
235        &mut self,
236        partition_type: &StructType,
237    ) -> Result<Vec<FieldSummary>> {
238        let mut field_stats: Vec<_> = partition_type
239            .fields()
240            .iter()
241            .map(|f| PartitionFieldStats::new(f.field_type.as_primitive_type().unwrap().clone()))
242            .collect();
243        for partition in self.manifest_entries.iter().map(|e| &e.data_file.partition) {
244            for (literal, stat) in partition.iter().zip_eq(field_stats.iter_mut()) {
245                let primitive_literal = literal.map(|v| v.as_primitive_literal().unwrap());
246                stat.update(primitive_literal)?;
247            }
248        }
249        Ok(field_stats.into_iter().map(|stat| stat.finish()).collect())
250    }
251
252    fn check_data_file(&self, data_file: &DataFile) -> Result<()> {
253        match self.metadata.content {
254            ManifestContentType::Data => {
255                if data_file.content != DataContentType::Data {
256                    return Err(invalid_data!(
257                        "Date file at path {} with manifest content type `data`, should have DataContentType `Data`, but has `{:?}`",
258                        data_file.file_path(),
259                        data_file.content
260                    ));
261                }
262            }
263            ManifestContentType::Deletes => {
264                if data_file.content != DataContentType::EqualityDeletes
265                    && data_file.content != DataContentType::PositionDeletes
266                {
267                    return Err(invalid_data!(
268                        "Date file at path {} with manifest content type `deletes`, should have DataContentType `Data`, but has `{:?}`",
269                        data_file.file_path(),
270                        data_file.content
271                    ));
272                }
273            }
274        }
275        Ok(())
276    }
277
278    /// Add a new manifest entry. This method will update following status of the entry:
279    /// - Update the entry status to `Added`
280    /// - Set the snapshot id to the current snapshot id
281    /// - Set the sequence number to `None` if it is invalid(smaller than 0)
282    /// - Set the file sequence number to `None`
283    pub(crate) fn add_entry(&mut self, mut entry: ManifestEntry) -> Result<()> {
284        self.check_data_file(&entry.data_file)?;
285        if entry.sequence_number().is_some_and(|n| n >= 0) {
286            entry.status = ManifestStatus::Added;
287            entry.snapshot_id = self.snapshot_id;
288            entry.file_sequence_number = None;
289        } else {
290            entry.status = ManifestStatus::Added;
291            entry.snapshot_id = self.snapshot_id;
292            entry.sequence_number = None;
293            entry.file_sequence_number = None;
294        };
295        self.add_entry_inner(entry)?;
296        Ok(())
297    }
298
299    /// Add file as an added entry with a specific sequence number. The entry's snapshot ID will be this manifest's snapshot ID. The entry's data sequence
300    /// number will be the provided data sequence number. The entry's file sequence number will be
301    /// assigned at commit.
302    pub fn add_file(&mut self, data_file: DataFile, sequence_number: i64) -> Result<()> {
303        self.check_data_file(&data_file)?;
304        let entry = ManifestEntry {
305            status: ManifestStatus::Added,
306            snapshot_id: self.snapshot_id,
307            sequence_number: (sequence_number >= 0).then_some(sequence_number),
308            file_sequence_number: None,
309            data_file,
310        };
311        self.add_entry_inner(entry)?;
312        Ok(())
313    }
314
315    /// Add a delete manifest entry. This method will update following status of the entry:
316    /// - Update the entry status to `Deleted`
317    /// - Set the snapshot id to the current snapshot id
318    ///
319    /// # TODO
320    /// Remove this allow later
321    #[allow(dead_code)]
322    pub(crate) fn add_delete_entry(&mut self, mut entry: ManifestEntry) -> Result<()> {
323        self.check_data_file(&entry.data_file)?;
324        entry.status = ManifestStatus::Deleted;
325        entry.snapshot_id = self.snapshot_id;
326        self.add_entry_inner(entry)?;
327        Ok(())
328    }
329
330    /// Add a file as delete manifest entry. The entry's snapshot ID will be this manifest's snapshot ID.
331    /// However, the original data and file sequence numbers of the file must be preserved when
332    /// the file is marked as deleted.
333    pub fn add_delete_file(
334        &mut self,
335        data_file: DataFile,
336        sequence_number: i64,
337        file_sequence_number: Option<i64>,
338    ) -> Result<()> {
339        self.check_data_file(&data_file)?;
340        let entry = ManifestEntry {
341            status: ManifestStatus::Deleted,
342            snapshot_id: self.snapshot_id,
343            sequence_number: Some(sequence_number),
344            file_sequence_number,
345            data_file,
346        };
347        self.add_entry_inner(entry)?;
348        Ok(())
349    }
350
351    /// Add an existing manifest entry. This method will update following status of the entry:
352    /// - Update the entry status to `Existing`
353    ///
354    /// # TODO
355    /// Remove this allow later
356    #[allow(dead_code)]
357    pub(crate) fn add_existing_entry(&mut self, mut entry: ManifestEntry) -> Result<()> {
358        self.check_data_file(&entry.data_file)?;
359        entry.status = ManifestStatus::Existing;
360        self.add_entry_inner(entry)?;
361        Ok(())
362    }
363
364    /// Add an file as existing manifest entry. The original data and file sequence numbers, snapshot ID,
365    /// which were assigned at commit, must be preserved when adding an existing entry.
366    pub fn add_existing_file(
367        &mut self,
368        data_file: DataFile,
369        snapshot_id: i64,
370        sequence_number: i64,
371        file_sequence_number: Option<i64>,
372    ) -> Result<()> {
373        self.check_data_file(&data_file)?;
374        let entry = ManifestEntry {
375            status: ManifestStatus::Existing,
376            snapshot_id: Some(snapshot_id),
377            sequence_number: Some(sequence_number),
378            file_sequence_number,
379            data_file,
380        };
381        self.add_entry_inner(entry)?;
382        Ok(())
383    }
384
385    fn add_entry_inner(&mut self, entry: ManifestEntry) -> Result<()> {
386        // Check if the entry has sequence number
387        if (entry.status == ManifestStatus::Deleted || entry.status == ManifestStatus::Existing)
388            && (entry.sequence_number.is_none() || entry.file_sequence_number.is_none())
389        {
390            return Err(invalid_data!(
391                "Manifest entry with status Existing or Deleted should have sequence number"
392            ));
393        }
394
395        // Update the statistics
396        match entry.status {
397            ManifestStatus::Added => {
398                self.added_files += 1;
399                self.added_rows += entry.data_file.record_count;
400            }
401            ManifestStatus::Deleted => {
402                self.deleted_files += 1;
403                self.deleted_rows += entry.data_file.record_count;
404            }
405            ManifestStatus::Existing => {
406                self.existing_files += 1;
407                self.existing_rows += entry.data_file.record_count;
408            }
409        }
410        if entry.is_alive()
411            && let Some(seq_num) = entry.sequence_number
412        {
413            self.min_seq_num = Some(self.min_seq_num.map_or(seq_num, |v| min(v, seq_num)));
414        }
415        self.manifest_entries.push(entry);
416        Ok(())
417    }
418
419    /// Write manifest file and return it.
420    pub async fn write_manifest_file(mut self) -> Result<ManifestFile> {
421        // Create the avro writer
422        let partition_type = self
423            .metadata
424            .partition_spec
425            .partition_type(&self.metadata.schema)?;
426        // Wrap once and reuse for every entry so the per-entry conversion does
427        // not rebuild the partition type's field-name lookup on each call.
428        let partition_struct_type = Type::Struct(partition_type.clone());
429        let table_schema = &self.metadata.schema;
430        let avro_schema = match self.metadata.format_version {
431            FormatVersion::V1 => manifest_schema_v1(&partition_type)?,
432            // Manifest schema did not change between V2 and V3
433            FormatVersion::V2 | FormatVersion::V3 => manifest_schema_v2(&partition_type)?,
434        };
435        let mut avro_writer = AvroWriter::new(&avro_schema, Vec::new())?;
436        avro_writer.add_user_metadata(
437            "schema".to_string(),
438            to_vec(table_schema)
439                .map_err(|err| invalid_data!("Fail to serialize table schema").with_source(err))?,
440        )?;
441        avro_writer.add_user_metadata(
442            "schema-id".to_string(),
443            table_schema.schema_id().to_string(),
444        )?;
445        avro_writer.add_user_metadata(
446            "partition-spec".to_string(),
447            to_vec(&self.metadata.partition_spec.fields()).map_err(|err| {
448                invalid_data!("Fail to serialize partition spec").with_source(err)
449            })?,
450        )?;
451        avro_writer.add_user_metadata(
452            "partition-spec-id".to_string(),
453            self.metadata.partition_spec.spec_id().to_string(),
454        )?;
455        avro_writer.add_user_metadata(
456            "format-version".to_string(),
457            (self.metadata.format_version as u8).to_string(),
458        )?;
459        match self.metadata.format_version {
460            FormatVersion::V1 => {}
461            FormatVersion::V2 | FormatVersion::V3 => {
462                avro_writer
463                    .add_user_metadata("content".to_string(), self.metadata.content.to_string())?;
464            }
465        }
466
467        let partition_summary = self.construct_partition_summaries(&partition_type)?;
468        // Write manifest entries
469        for entry in std::mem::take(&mut self.manifest_entries) {
470            let value = match self.metadata.format_version {
471                FormatVersion::V1 => {
472                    to_value(ManifestEntryV1::try_from(entry, &partition_struct_type)?)?
473                        .resolve(&avro_schema)?
474                }
475                // Manifest entry format did not change between V2 and V3
476                FormatVersion::V2 | FormatVersion::V3 => {
477                    to_value(ManifestEntryV2::try_from(entry, &partition_struct_type)?)?
478                        .resolve(&avro_schema)?
479                }
480            };
481
482            avro_writer.append_value(value)?;
483        }
484
485        let content = avro_writer.into_inner()?;
486        let mut writer = self.output.writer().await?;
487        writer.write(Bytes::from(content)).await?;
488        let file_metadata = writer.close().await?;
489        let key_metadata = self.output.encoded_key_metadata(&file_metadata)?;
490
491        Ok(ManifestFile {
492            manifest_path: self.location,
493            // Manifest lengths are on-disk sizes, including encryption overhead.
494            manifest_length: file_metadata.size.try_into()?,
495            partition_spec_id: self.metadata.partition_spec.spec_id(),
496            content: self.metadata.content,
497            // sequence_number and min_sequence_number with UNASSIGNED_SEQUENCE_NUMBER will be replace with
498            // real sequence number in `ManifestListWriter`.
499            sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
500            min_sequence_number: self.min_seq_num.unwrap_or(UNASSIGNED_SEQUENCE_NUMBER),
501            added_snapshot_id: self.snapshot_id.unwrap_or(UNASSIGNED_SNAPSHOT_ID),
502            added_files_count: Some(self.added_files),
503            existing_files_count: Some(self.existing_files),
504            deleted_files_count: Some(self.deleted_files),
505            added_rows_count: Some(self.added_rows),
506            existing_rows_count: Some(self.existing_rows),
507            deleted_rows_count: Some(self.deleted_rows),
508            partitions: Some(partition_summary),
509            key_metadata,
510            first_row_id: self.first_row_id,
511        })
512    }
513}
514
515struct PartitionFieldStats {
516    partition_type: PrimitiveType,
517
518    contains_null: bool,
519    contains_nan: Option<bool>,
520    lower_bound: Option<Datum>,
521    upper_bound: Option<Datum>,
522}
523
524impl PartitionFieldStats {
525    pub(crate) fn new(partition_type: PrimitiveType) -> Self {
526        Self {
527            partition_type,
528            contains_null: false,
529            contains_nan: Some(false),
530            upper_bound: None,
531            lower_bound: None,
532        }
533    }
534
535    pub(crate) fn update(&mut self, value: Option<PrimitiveLiteral>) -> Result<()> {
536        let Some(value) = value else {
537            self.contains_null = true;
538            return Ok(());
539        };
540        if !self.partition_type.compatible(&value) {
541            return Err(invalid_data!("value is not compatible with type"));
542        }
543        let value = Datum::new(self.partition_type.clone(), value);
544
545        if value.is_nan() {
546            self.contains_nan = Some(true);
547            return Ok(());
548        }
549
550        self.lower_bound = Some(self.lower_bound.take().map_or(value.clone(), |original| {
551            if value < original {
552                value.clone()
553            } else {
554                original
555            }
556        }));
557        self.upper_bound = Some(self.upper_bound.take().map_or(value.clone(), |original| {
558            if value > original { value } else { original }
559        }));
560
561        Ok(())
562    }
563
564    pub(crate) fn finish(self) -> FieldSummary {
565        FieldSummary {
566            contains_null: self.contains_null,
567            contains_nan: self.contains_nan,
568            upper_bound: self.upper_bound.map(|v| v.to_bytes().unwrap()),
569            lower_bound: self.lower_bound.map(|v| v.to_bytes().unwrap()),
570        }
571    }
572}
573
574#[cfg(test)]
575mod tests {
576    use std::collections::HashMap;
577    use std::fs;
578    use std::sync::Arc;
579
580    use tempfile::TempDir;
581
582    use super::*;
583    use crate::avro::define_named_types_repeatedly;
584    use crate::io::FileIO;
585    use crate::spec::{
586        DataContentType, DataFileBuilder, DataFileFormat, Literal, Manifest, NestedField,
587        PrimitiveType, Schema, Struct, Transform, Type,
588    };
589
590    #[tokio::test]
591    async fn test_add_delete_existing() {
592        let schema = Arc::new(
593            Schema::builder()
594                .with_fields(vec![
595                    Arc::new(NestedField::optional(
596                        1,
597                        "id",
598                        Type::Primitive(PrimitiveType::Int),
599                    )),
600                    Arc::new(NestedField::optional(
601                        2,
602                        "name",
603                        Type::Primitive(PrimitiveType::String),
604                    )),
605                ])
606                .build()
607                .unwrap(),
608        );
609        let metadata = ManifestMetadata {
610            schema_id: 0,
611            schema: schema.clone(),
612            partition_spec: PartitionSpec::builder(schema)
613                .with_spec_id(0)
614                .build()
615                .unwrap(),
616            content: ManifestContentType::Data,
617            format_version: FormatVersion::V2,
618        };
619        let mut entries = vec![
620                ManifestEntry {
621                    status: ManifestStatus::Added,
622                    snapshot_id: None,
623                    sequence_number: Some(1),
624                    file_sequence_number: Some(1),
625                    data_file: DataFile {
626                        content: DataContentType::Data,
627                        file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
628                        file_format: DataFileFormat::Parquet,
629                        partition: Struct::empty(),
630                        record_count: 1,
631                        file_size_in_bytes: 5442,
632                        column_sizes: HashMap::from([(1, 61), (2, 73)]),
633                        value_counts: HashMap::from([(1, 1), (2, 1)]),
634                        null_value_counts: HashMap::from([(1, 0), (2, 0)]),
635                        nan_value_counts: HashMap::new(),
636                        lower_bounds: HashMap::new(),
637                        upper_bounds: HashMap::new(),
638                        key_metadata: Some(Vec::new()),
639                        split_offsets: Some(vec![4]),
640                        equality_ids: None,
641                        sort_order_id: None,
642                        partition_spec_id: 0,
643                        first_row_id: None,
644                        referenced_data_file: None,
645                        content_offset: None,
646                        content_size_in_bytes: None,
647                    },
648                },
649                ManifestEntry {
650                    status: ManifestStatus::Deleted,
651                    snapshot_id: Some(1),
652                    sequence_number: Some(1),
653                    file_sequence_number: Some(1),
654                    data_file: DataFile {
655                        content: DataContentType::Data,
656                        file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
657                        file_format: DataFileFormat::Parquet,
658                        partition: Struct::empty(),
659                        record_count: 1,
660                        file_size_in_bytes: 5442,
661                        column_sizes: HashMap::from([(1, 61), (2, 73)]),
662                        value_counts: HashMap::from([(1, 1), (2, 1)]),
663                        null_value_counts: HashMap::from([(1, 0), (2, 0)]),
664                        nan_value_counts: HashMap::new(),
665                        lower_bounds: HashMap::new(),
666                        upper_bounds: HashMap::new(),
667                        key_metadata: Some(Vec::new()),
668                        split_offsets: Some(vec![4]),
669                        equality_ids: None,
670                        sort_order_id: None,
671                        partition_spec_id: 0,
672                        first_row_id: None,
673                        referenced_data_file: None,
674                        content_offset: None,
675                        content_size_in_bytes: None,
676                    },
677                },
678                ManifestEntry {
679                    status: ManifestStatus::Existing,
680                    snapshot_id: Some(1),
681                    sequence_number: Some(1),
682                    file_sequence_number: Some(1),
683                    data_file: DataFile {
684                        content: DataContentType::Data,
685                        file_path: "s3a://icebergdata/demo/s1/t1/data/00000-0-ba56fbfa-f2ff-40c9-bb27-565ad6dc2be8-00000.parquet".to_string(),
686                        file_format: DataFileFormat::Parquet,
687                        partition: Struct::empty(),
688                        record_count: 1,
689                        file_size_in_bytes: 5442,
690                        column_sizes: HashMap::from([(1, 61), (2, 73)]),
691                        value_counts: HashMap::from([(1, 1), (2, 1)]),
692                        null_value_counts: HashMap::from([(1, 0), (2, 0)]),
693                        nan_value_counts: HashMap::new(),
694                        lower_bounds: HashMap::new(),
695                        upper_bounds: HashMap::new(),
696                        key_metadata: Some(Vec::new()),
697                        split_offsets: Some(vec![4]),
698                        equality_ids: None,
699                        sort_order_id: None,
700                        partition_spec_id: 0,
701                        first_row_id: None,
702                        referenced_data_file: None,
703                        content_offset: None,
704                        content_size_in_bytes: None,
705                    },
706                },
707            ];
708
709        // write manifest to file
710        let tmp_dir = TempDir::new().unwrap();
711        let path = tmp_dir.path().join("test_manifest.avro");
712        let io = FileIO::new_with_fs();
713        let output_file = io.new_output(path.to_str().unwrap()).unwrap();
714        let mut writer = ManifestWriterBuilder::new(
715            output_file,
716            Some(3),
717            metadata.schema.clone(),
718            metadata.partition_spec.clone(),
719        )
720        .build_v2_data();
721        writer.add_entry(entries[0].clone()).unwrap();
722        writer.add_delete_entry(entries[1].clone()).unwrap();
723        writer.add_existing_entry(entries[2].clone()).unwrap();
724        writer.write_manifest_file().await.unwrap();
725
726        // read back the manifest file and check the content
727        let actual_manifest =
728            Manifest::parse_avro(fs::read(path).expect("read_file must succeed").as_slice())
729                .unwrap();
730
731        // The snapshot id is assigned when the entry is added and delete to the manifest. Existing entries are keep original.
732        entries[0].snapshot_id = Some(3);
733        entries[1].snapshot_id = Some(3);
734        // file sequence number is assigned to None when the entry is added and delete to the manifest.
735        entries[0].file_sequence_number = None;
736        assert_eq!(actual_manifest, Manifest::new(metadata, entries));
737    }
738
739    #[tokio::test]
740    async fn test_v3_delete_manifest_delete_file_roundtrip() {
741        let schema = Arc::new(
742            Schema::builder()
743                .with_fields(vec![
744                    Arc::new(NestedField::optional(
745                        1,
746                        "id",
747                        Type::Primitive(PrimitiveType::Long),
748                    )),
749                    Arc::new(NestedField::optional(
750                        2,
751                        "data",
752                        Type::Primitive(PrimitiveType::String),
753                    )),
754                ])
755                .build()
756                .unwrap(),
757        );
758
759        let partition_spec = PartitionSpec::builder(schema.clone())
760            .with_spec_id(0)
761            .build()
762            .unwrap();
763
764        // Create a position delete file entry
765        let delete_entry = ManifestEntry {
766            status: ManifestStatus::Added,
767            snapshot_id: None,
768            sequence_number: None,
769            file_sequence_number: None,
770            data_file: DataFile {
771                content: DataContentType::PositionDeletes,
772                file_path: "s3://bucket/table/data/delete-00000.parquet".to_string(),
773                file_format: DataFileFormat::Parquet,
774                partition: Struct::empty(),
775                record_count: 10,
776                file_size_in_bytes: 1024,
777                column_sizes: HashMap::new(),
778                value_counts: HashMap::new(),
779                null_value_counts: HashMap::new(),
780                nan_value_counts: HashMap::new(),
781                lower_bounds: HashMap::new(),
782                upper_bounds: HashMap::new(),
783                key_metadata: None,
784                split_offsets: None,
785                equality_ids: None,
786                sort_order_id: None,
787                partition_spec_id: 0,
788                first_row_id: None,
789                referenced_data_file: None,
790                content_offset: None,
791                content_size_in_bytes: None,
792            },
793        };
794
795        // Write a V3 delete manifest
796        let tmp_dir = TempDir::new().unwrap();
797        let path = tmp_dir.path().join("v3_delete_manifest.avro");
798        let io = FileIO::new_with_fs();
799        let output_file = io.new_output(path.to_str().unwrap()).unwrap();
800
801        let mut writer = ManifestWriterBuilder::new(
802            output_file,
803            Some(1),
804            schema.clone(),
805            partition_spec.clone(),
806        )
807        .build_v3_deletes();
808
809        writer.add_entry(delete_entry).unwrap();
810        let manifest_file = writer.write_manifest_file().await.unwrap();
811
812        // The returned ManifestFile correctly reports Deletes content
813        assert_eq!(manifest_file.content, ManifestContentType::Deletes);
814
815        // Read back the manifest file
816        let bytes = fs::read(&path).expect("read_file must succeed");
817        assert_eq!(manifest_file.manifest_length, bytes.len() as i64);
818        let actual_manifest = Manifest::parse_avro(&bytes).unwrap();
819
820        // Verify the content type is correctly preserved as Deletes
821        assert_eq!(
822            actual_manifest.metadata().content,
823            ManifestContentType::Deletes,
824        );
825    }
826
827    /// Writes a manifest partitioned by identity on two columns of `field_type`,
828    /// reads it back, and returns the manifest file.
829    async fn roundtrip_two_partition_fields_of_type(
830        field_type: Type,
831        values: [Literal; 2],
832    ) -> Vec<u8> {
833        let schema = Arc::new(
834            Schema::builder()
835                .with_fields(vec![
836                    NestedField::optional(1, "a", field_type.clone()).into(),
837                    NestedField::optional(2, "b", field_type).into(),
838                ])
839                .build()
840                .unwrap(),
841        );
842        let partition_spec = PartitionSpec::builder(schema.clone())
843            .with_spec_id(0)
844            .add_partition_field("a", "a", Transform::Identity)
845            .unwrap()
846            .add_partition_field("b", "b", Transform::Identity)
847            .unwrap()
848            .build()
849            .unwrap();
850        let partition = Struct::from_iter(values.map(Some));
851        let tmp_dir = TempDir::new().unwrap();
852        let path = tmp_dir.path().join("manifest.avro");
853        let output_file = FileIO::new_with_fs()
854            .new_output(path.to_str().unwrap())
855            .unwrap();
856        let mut writer = ManifestWriterBuilder::new(output_file, Some(1), schema, partition_spec)
857            .build_v2_data();
858        writer
859            .add_file(
860                DataFileBuilder::default()
861                    .content(DataContentType::Data)
862                    .file_path("s3://bucket/table/data/a.parquet".to_string())
863                    .file_format(DataFileFormat::Parquet)
864                    .partition(partition.clone())
865                    .record_count(1)
866                    .file_size_in_bytes(1)
867                    .partition_spec_id(0)
868                    .build()
869                    .unwrap(),
870                1,
871            )
872            .unwrap();
873        writer.write_manifest_file().await.unwrap();
874        let bs = fs::read(&path).unwrap();
875
876        // Also read the manifest as iceberg-rust wrote it before defining each
877        // named type once.
878        for bs in [bs.clone(), define_named_types_repeatedly(&bs)] {
879            let manifest = Manifest::parse_avro(&bs).unwrap();
880
881            assert_eq!(*manifest.entries()[0].data_file().partition(), partition);
882        }
883        bs
884    }
885
886    #[tokio::test]
887    async fn test_write_manifest_header_marks_int_keyed_maps() {
888        // Java and PyIceberg read an array of key-value records as a map only if
889        // it has `logicalType: map`. apache-avro 0.22 drops the attribute when it
890        // parses a schema (https://github.com/apache/avro-rs/issues/654), so this
891        // checks the schema JSON in the written header.
892        let bs = roundtrip_two_partition_fields_of_type(Type::Primitive(PrimitiveType::Int), [
893            Literal::int(1),
894            Literal::int(2),
895        ])
896        .await;
897
898        let file = String::from_utf8_lossy(&bs);
899        // `column_sizes`, `value_counts`, `null_value_counts`,
900        // `nan_value_counts`, `lower_bounds`, and `upper_bounds`.
901        assert_eq!(file.matches(r#""logicalType":"map""#).count(), 6);
902    }
903
904    #[tokio::test]
905    async fn test_write_manifest_with_repeated_decimal_partition_type() {
906        roundtrip_two_partition_fields_of_type(
907            Type::Primitive(PrimitiveType::Decimal {
908                precision: 10,
909                scale: 2,
910            }),
911            [Literal::decimal(12345), Literal::decimal(-678)],
912        )
913        .await;
914    }
915
916    #[tokio::test]
917    async fn test_write_manifest_with_repeated_fixed_partition_type() {
918        roundtrip_two_partition_fields_of_type(Type::Primitive(PrimitiveType::Fixed(4)), [
919            Literal::fixed(vec![1, 2, 3, 4]),
920            Literal::fixed(vec![5, 6, 7, 8]),
921        ])
922        .await;
923    }
924}