Skip to main content

iceberg/spec/manifest_list/
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::collections::HashMap;
19
20use apache_avro::Writer;
21use bytes::Bytes;
22
23use super::_const_schema::{
24    MANIFEST_LIST_AVRO_SCHEMA_V1, MANIFEST_LIST_AVRO_SCHEMA_V2, MANIFEST_LIST_AVRO_SCHEMA_V3,
25};
26use super::_serde::{ManifestFileV1, ManifestFileV2, ManifestFileV3};
27use super::{FormatVersion, ManifestContentType, ManifestFile, UNASSIGNED_SEQUENCE_NUMBER};
28use crate::error::{Result, invalid_data};
29use crate::io::{FileMetadata, FileWrite};
30use crate::{Error, ErrorKind};
31
32/// A manifest list writer.
33pub struct ManifestListWriter {
34    format_version: FormatVersion,
35    writer: Box<dyn FileWrite>,
36    avro_writer: Writer<'static, Vec<u8>>,
37    sequence_number: i64,
38    snapshot_id: i64,
39    next_row_id: Option<u64>,
40}
41
42impl std::fmt::Debug for ManifestListWriter {
43    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
44        f.debug_struct("ManifestListWriter")
45            .field("format_version", &self.format_version)
46            .field("avro_writer", &self.avro_writer.schema())
47            .finish_non_exhaustive()
48    }
49}
50
51impl ManifestListWriter {
52    /// Get the next row ID that will be assigned to the next data manifest added.
53    pub fn next_row_id(&self) -> Option<u64> {
54        self.next_row_id
55    }
56
57    /// Construct a v1 [`ManifestListWriter`] that writes to a provided [`FileWrite`].
58    pub fn v1(
59        writer: Box<dyn FileWrite>,
60        snapshot_id: i64,
61        parent_snapshot_id: Option<i64>,
62    ) -> Self {
63        let mut metadata = HashMap::from_iter([
64            ("snapshot-id".to_string(), snapshot_id.to_string()),
65            ("format-version".to_string(), "1".to_string()),
66        ]);
67        if let Some(parent_snapshot_id) = parent_snapshot_id {
68            metadata.insert(
69                "parent-snapshot-id".to_string(),
70                parent_snapshot_id.to_string(),
71            );
72        }
73        Self::new(FormatVersion::V1, writer, metadata, 0, snapshot_id, None)
74    }
75
76    /// Construct a v2 [`ManifestListWriter`] that writes to a provided [`FileWrite`].
77    pub fn v2(
78        writer: Box<dyn FileWrite>,
79        snapshot_id: i64,
80        parent_snapshot_id: Option<i64>,
81        sequence_number: i64,
82    ) -> Self {
83        let mut metadata = HashMap::from_iter([
84            ("snapshot-id".to_string(), snapshot_id.to_string()),
85            ("sequence-number".to_string(), sequence_number.to_string()),
86            ("format-version".to_string(), "2".to_string()),
87        ]);
88        metadata.insert(
89            "parent-snapshot-id".to_string(),
90            parent_snapshot_id
91                .map(|v| v.to_string())
92                .unwrap_or("null".to_string()),
93        );
94        Self::new(
95            FormatVersion::V2,
96            writer,
97            metadata,
98            sequence_number,
99            snapshot_id,
100            None,
101        )
102    }
103
104    /// Construct a v3 [`ManifestListWriter`] that writes to a provided [`FileWrite`].
105    pub fn v3(
106        writer: Box<dyn FileWrite>,
107        snapshot_id: i64,
108        parent_snapshot_id: Option<i64>,
109        sequence_number: i64,
110        first_row_id: Option<u64>, // Always None for delete manifests
111    ) -> Self {
112        let mut metadata = HashMap::from_iter([
113            ("snapshot-id".to_string(), snapshot_id.to_string()),
114            ("sequence-number".to_string(), sequence_number.to_string()),
115            ("format-version".to_string(), "3".to_string()),
116        ]);
117        metadata.insert(
118            "parent-snapshot-id".to_string(),
119            parent_snapshot_id
120                .map(|v| v.to_string())
121                .unwrap_or("null".to_string()),
122        );
123        metadata.insert(
124            "first-row-id".to_string(),
125            first_row_id
126                .map(|v| v.to_string())
127                .unwrap_or("null".to_string()),
128        );
129        Self::new(
130            FormatVersion::V3,
131            writer,
132            metadata,
133            sequence_number,
134            snapshot_id,
135            first_row_id,
136        )
137    }
138
139    fn new(
140        format_version: FormatVersion,
141        writer: Box<dyn FileWrite>,
142        metadata: HashMap<String, String>,
143        sequence_number: i64,
144        snapshot_id: i64,
145        first_row_id: Option<u64>,
146    ) -> Self {
147        let avro_schema = match format_version {
148            FormatVersion::V1 => &MANIFEST_LIST_AVRO_SCHEMA_V1,
149            FormatVersion::V2 => &MANIFEST_LIST_AVRO_SCHEMA_V2,
150            FormatVersion::V3 => &MANIFEST_LIST_AVRO_SCHEMA_V3,
151        };
152        let mut avro_writer = Writer::new(avro_schema, Vec::new()).expect(
153            "Manifest list Avro schemas should resolve because they refer to no named types.",
154        );
155        for (key, value) in metadata {
156            avro_writer
157                .add_user_metadata(key, value)
158                .expect("Avro metadata should be added to the writer before the first record.");
159        }
160        Self {
161            format_version,
162            writer,
163            avro_writer,
164            sequence_number,
165            snapshot_id,
166            next_row_id: first_row_id,
167        }
168    }
169
170    /// Append manifests to be written.
171    ///
172    /// If V3 Manifests are added and the `first_row_id` of any data manifest is unassigned,
173    /// it will be assigned based on the `next_row_id` of the writer, and the `next_row_id` of the writer will be updated accordingly.
174    /// If `first_row_id` is already assigned, it will be validated against the `next_row_id` of the writer.
175    pub fn add_manifests(&mut self, manifests: impl Iterator<Item = ManifestFile>) -> Result<()> {
176        match self.format_version {
177            FormatVersion::V1 => {
178                for manifest in manifests {
179                    let manifests: ManifestFileV1 = manifest.try_into()?;
180                    self.avro_writer.append_ser(manifests)?;
181                }
182            }
183            FormatVersion::V2 | FormatVersion::V3 => {
184                for mut manifest in manifests {
185                    self.assign_sequence_numbers(&mut manifest)?;
186
187                    if self.format_version == FormatVersion::V2 {
188                        let manifest_entry: ManifestFileV2 = manifest.try_into()?;
189                        self.avro_writer.append_ser(manifest_entry)?;
190                    } else if self.format_version == FormatVersion::V3 {
191                        self.assign_first_row_id(&mut manifest)?;
192                        let manifest_entry: ManifestFileV3 = manifest.try_into()?;
193                        self.avro_writer.append_ser(manifest_entry)?;
194                    }
195                }
196            }
197        }
198        Ok(())
199    }
200
201    /// Write the manifest list and return its stored size.
202    pub async fn close(mut self) -> Result<FileMetadata> {
203        let data = self.avro_writer.into_inner()?;
204        self.writer.write(Bytes::from(data)).await?;
205        self.writer.close().await
206    }
207
208    /// Assign sequence numbers to manifest if they are unassigned
209    fn assign_sequence_numbers(&self, manifest: &mut ManifestFile) -> Result<()> {
210        if manifest.sequence_number == UNASSIGNED_SEQUENCE_NUMBER {
211            if manifest.added_snapshot_id != self.snapshot_id {
212                return Err(invalid_data!(
213                    "Found unassigned sequence number for a manifest from snapshot {}.",
214                    manifest.added_snapshot_id
215                ));
216            }
217            manifest.sequence_number = self.sequence_number;
218        }
219
220        if manifest.min_sequence_number == UNASSIGNED_SEQUENCE_NUMBER {
221            if manifest.added_snapshot_id != self.snapshot_id {
222                return Err(invalid_data!(
223                    "Found unassigned sequence number for a manifest from snapshot {}.",
224                    manifest.added_snapshot_id
225                ));
226            }
227            manifest.min_sequence_number = self.sequence_number;
228        }
229
230        Ok(())
231    }
232
233    /// Returns number of newly assigned first-row-ids, if any.
234    fn assign_first_row_id(&mut self, manifest: &mut ManifestFile) -> Result<()> {
235        match manifest.content {
236            ManifestContentType::Data => {
237                match (self.next_row_id, manifest.first_row_id) {
238                    (Some(_), Some(_)) => {
239                        // Case: Manifest with already assigned first row ID.
240                        // No need to increase next_row_id, as this manifest is already assigned.
241                    }
242                    (None, Some(manifest_first_row_id)) => {
243                        // Case: Assigned first row ID for data manifest, but the writer does not have a next-row-id assigned.
244                        return Err(Error::new(
245                            ErrorKind::Unexpected,
246                            format!(
247                                "Found invalid first-row-id assignment for Manifest {}. Writer does not have a next-row-id assigned, but the manifest has first-row-id assigned to {}.",
248                                manifest.manifest_path, manifest_first_row_id,
249                            ),
250                        ));
251                    }
252                    (Some(writer_next_row_id), None) => {
253                        // Case: Unassigned first row ID for data manifest. This is either a new
254                        // manifest, or a manifest from a pre-v3 snapshot. We need to assign one.
255                        let (existing_rows_count, added_rows_count) =
256                            require_row_counts_in_manifest(manifest)?;
257                        manifest.first_row_id = Some(writer_next_row_id);
258
259                        self.next_row_id = writer_next_row_id
260                        .checked_add(existing_rows_count)
261                        .and_then(|sum| sum.checked_add(added_rows_count))
262                        .ok_or_else(|| {
263                            invalid_data!(
264                                    "Row ID overflow when computing next row ID for Manifest {}. Next Row ID: {writer_next_row_id}, Existing Rows Count: {existing_rows_count}, Added Rows Count: {added_rows_count}",
265                                    manifest.manifest_path
266                                )
267                        }).map(Some)?;
268                    }
269                    (None, None) => {
270                        // Case: Table without row lineage. No action needed.
271                    }
272                }
273            }
274            ManifestContentType::Deletes => {
275                // Deletes never have a first-row-id assigned.
276                manifest.first_row_id = None;
277            }
278        };
279
280        Ok(())
281    }
282}
283
284fn require_row_counts_in_manifest(manifest: &ManifestFile) -> Result<(u64, u64)> {
285    let existing_rows_count = manifest.existing_rows_count.ok_or_else(|| {
286        invalid_data!(
287                "Cannot include a Manifest without existing-rows-count to a table with row lineage enabled. Manifest path: {}",
288                manifest.manifest_path,
289            )
290    })?;
291    let added_rows_count = manifest.added_rows_count.ok_or_else(|| {
292        invalid_data!(
293                "Cannot include a Manifest without added-rows-count to a table with row lineage enabled. Manifest path: {}",
294                manifest.manifest_path,
295            )
296    })?;
297    Ok((existing_rows_count, added_rows_count))
298}
299
300#[cfg(test)]
301mod test {
302    use std::fs;
303    use std::path::Path;
304    use std::sync::Arc;
305
306    use tempfile::TempDir;
307
308    use super::ManifestListWriter;
309    use crate::encryption::kms::{KeyManagementClient, MemoryKeyManagementClient};
310    use crate::encryption::{EncryptedInputFile, EncryptionManager};
311    use crate::io::{FileIO, FileWrite};
312    use crate::spec::{
313        Datum, FieldSummary, ManifestContentType, ManifestFile, ManifestList,
314        UNASSIGNED_SEQUENCE_NUMBER,
315    };
316
317    #[tokio::test]
318    async fn test_manifest_list_writer_v1() {
319        let expected_manifest_list = ManifestList {
320            entries: vec![ManifestFile {
321                manifest_path: "/opt/bitnami/spark/warehouse/db/table/metadata/10d28031-9739-484c-92db-cdf2975cead4-m0.avro".to_string(),
322                manifest_length: 5806,
323                partition_spec_id: 1,
324                content: ManifestContentType::Data,
325                sequence_number: 0,
326                min_sequence_number: 0,
327                added_snapshot_id: 1646658105718557341,
328                added_files_count: Some(3),
329                existing_files_count: Some(0),
330                deleted_files_count: Some(0),
331                added_rows_count: Some(3),
332                existing_rows_count: Some(0),
333                deleted_rows_count: Some(0),
334                partitions: Some(
335                    vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}],
336                ),
337                key_metadata: None,
338                first_row_id: None,
339            }]
340        };
341
342        let temp_dir = TempDir::new().unwrap();
343        let path = temp_dir.path().join("manifest_list_v1.avro");
344        let io = FileIO::new_with_fs();
345        let file_writer = file_writer(&path, io).await;
346
347        let mut writer = ManifestListWriter::v1(file_writer, 1646658105718557341, Some(0));
348        writer
349            .add_manifests(expected_manifest_list.entries.clone().into_iter())
350            .unwrap();
351        writer.close().await.unwrap();
352
353        let bs = fs::read(path).unwrap();
354
355        let manifest_list =
356            ManifestList::parse_with_version(&bs, crate::spec::FormatVersion::V1).unwrap();
357        assert_eq!(manifest_list, expected_manifest_list);
358
359        temp_dir.close().unwrap();
360    }
361
362    #[tokio::test]
363    async fn test_manifest_list_writer_v2() {
364        let snapshot_id = 377075049360453639;
365        let seq_num = 1;
366        let mut expected_manifest_list = ManifestList {
367            entries: vec![ManifestFile {
368                manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m0.avro".to_string(),
369                manifest_length: 6926,
370                partition_spec_id: 1,
371                content: ManifestContentType::Data,
372                sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
373                min_sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
374                added_snapshot_id: snapshot_id,
375                added_files_count: Some(1),
376                existing_files_count: Some(0),
377                deleted_files_count: Some(0),
378                added_rows_count: Some(3),
379                existing_rows_count: Some(0),
380                deleted_rows_count: Some(0),
381                partitions: Some(
382                    vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
383                ),
384                key_metadata: None,
385                first_row_id: None,
386            }]
387        };
388
389        let temp_dir = TempDir::new().unwrap();
390        let path = temp_dir.path().join("manifest_list_v2.avro");
391        let io = FileIO::new_with_fs();
392        let file_writer = file_writer(&path, io).await;
393
394        let mut writer = ManifestListWriter::v2(file_writer, snapshot_id, Some(0), seq_num);
395        writer
396            .add_manifests(expected_manifest_list.entries.clone().into_iter())
397            .unwrap();
398        writer.close().await.unwrap();
399
400        let bs = fs::read(path).unwrap();
401        let manifest_list =
402            ManifestList::parse_with_version(&bs, crate::spec::FormatVersion::V2).unwrap();
403        expected_manifest_list.entries[0].sequence_number = seq_num;
404        expected_manifest_list.entries[0].min_sequence_number = seq_num;
405        assert_eq!(manifest_list, expected_manifest_list);
406
407        temp_dir.close().unwrap();
408    }
409
410    #[tokio::test]
411    async fn test_manifest_list_writer_v3() {
412        let snapshot_id = 377075049360453639;
413        let seq_num = 1;
414        let mut expected_manifest_list = ManifestList {
415            entries: vec![ManifestFile {
416                manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m0.avro".to_string(),
417                manifest_length: 6926,
418                partition_spec_id: 1,
419                content: ManifestContentType::Data,
420                sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
421                min_sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
422                added_snapshot_id: snapshot_id,
423                added_files_count: Some(1),
424                existing_files_count: Some(0),
425                deleted_files_count: Some(0),
426                added_rows_count: Some(3),
427                existing_rows_count: Some(0),
428                deleted_rows_count: Some(0),
429                partitions: Some(
430                    vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
431                ),
432                key_metadata: None,
433                first_row_id: Some(10),
434            }]
435        };
436
437        let temp_dir = TempDir::new().unwrap();
438        let path = temp_dir.path().join("manifest_list_v2.avro");
439        let io = FileIO::new_with_fs();
440        let file_writer = file_writer(&path, io).await;
441
442        let mut writer =
443            ManifestListWriter::v3(file_writer, snapshot_id, Some(0), seq_num, Some(10));
444        writer
445            .add_manifests(expected_manifest_list.entries.clone().into_iter())
446            .unwrap();
447        writer.close().await.unwrap();
448
449        let bs = fs::read(path).unwrap();
450        let manifest_list =
451            ManifestList::parse_with_version(&bs, crate::spec::FormatVersion::V3).unwrap();
452        expected_manifest_list.entries[0].sequence_number = seq_num;
453        expected_manifest_list.entries[0].min_sequence_number = seq_num;
454        expected_manifest_list.entries[0].first_row_id = Some(10);
455        assert_eq!(manifest_list, expected_manifest_list);
456
457        temp_dir.close().unwrap();
458    }
459
460    #[tokio::test]
461    async fn test_manifest_list_writer_v1_as_v2() {
462        let expected_manifest_list = ManifestList {
463            entries: vec![ManifestFile {
464                manifest_path: "/opt/bitnami/spark/warehouse/db/table/metadata/10d28031-9739-484c-92db-cdf2975cead4-m0.avro".to_string(),
465                manifest_length: 5806,
466                partition_spec_id: 1,
467                content: ManifestContentType::Data,
468                sequence_number: 0,
469                min_sequence_number: 0,
470                added_snapshot_id: 1646658105718557341,
471                added_files_count: Some(3),
472                existing_files_count: Some(0),
473                deleted_files_count: Some(0),
474                added_rows_count: Some(3),
475                existing_rows_count: Some(0),
476                deleted_rows_count: Some(0),
477                partitions: Some(
478                    vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
479                ),
480                key_metadata: None,
481                first_row_id: None,
482            }]
483        };
484
485        let temp_dir = TempDir::new().unwrap();
486        let path = temp_dir.path().join("manifest_list_v1.avro");
487        let io = FileIO::new_with_fs();
488        let file_writer = file_writer(&path, io).await;
489
490        let mut writer = ManifestListWriter::v1(file_writer, 1646658105718557341, Some(0));
491        writer
492            .add_manifests(expected_manifest_list.entries.clone().into_iter())
493            .unwrap();
494        writer.close().await.unwrap();
495
496        let bs = fs::read(path).unwrap();
497
498        let manifest_list =
499            ManifestList::parse_with_version(&bs, crate::spec::FormatVersion::V2).unwrap();
500        assert_eq!(manifest_list, expected_manifest_list);
501
502        temp_dir.close().unwrap();
503    }
504
505    #[tokio::test]
506    async fn test_manifest_list_writer_v1_as_v3() {
507        let expected_manifest_list = ManifestList {
508            entries: vec![ManifestFile {
509                manifest_path: "/opt/bitnami/spark/warehouse/db/table/metadata/10d28031-9739-484c-92db-cdf2975cead4-m0.avro".to_string(),
510                manifest_length: 5806,
511                partition_spec_id: 1,
512                content: ManifestContentType::Data,
513                sequence_number: 0,
514                min_sequence_number: 0,
515                added_snapshot_id: 1646658105718557341,
516                added_files_count: Some(3),
517                existing_files_count: Some(0),
518                deleted_files_count: Some(0),
519                added_rows_count: Some(3),
520                existing_rows_count: Some(0),
521                deleted_rows_count: Some(0),
522                partitions: Some(
523                    vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
524                ),
525                key_metadata: None,
526                first_row_id: None,
527            }]
528        };
529
530        let temp_dir = TempDir::new().unwrap();
531        let path = temp_dir.path().join("manifest_list_v1.avro");
532        let io = FileIO::new_with_fs();
533        let file_writer = file_writer(&path, io).await;
534
535        let mut writer = ManifestListWriter::v1(file_writer, 1646658105718557341, Some(0));
536        writer
537            .add_manifests(expected_manifest_list.entries.clone().into_iter())
538            .unwrap();
539        writer.close().await.unwrap();
540
541        let bs = fs::read(path).unwrap();
542
543        let manifest_list =
544            ManifestList::parse_with_version(&bs, crate::spec::FormatVersion::V3).unwrap();
545        assert_eq!(manifest_list, expected_manifest_list);
546
547        temp_dir.close().unwrap();
548    }
549
550    #[tokio::test]
551    async fn test_manifest_list_writer_v2_as_v3() {
552        let snapshot_id = 377075049360453639;
553        let seq_num = 1;
554        let mut expected_manifest_list = ManifestList {
555            entries: vec![ManifestFile {
556                manifest_path: "s3a://icebergdata/demo/s1/t1/metadata/05ffe08b-810f-49b3-a8f4-e88fc99b254a-m0.avro".to_string(),
557                manifest_length: 6926,
558                partition_spec_id: 1,
559                content: ManifestContentType::Data,
560                sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
561                min_sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
562                added_snapshot_id: snapshot_id,
563                added_files_count: Some(1),
564                existing_files_count: Some(0),
565                deleted_files_count: Some(0),
566                added_rows_count: Some(3),
567                existing_rows_count: Some(0),
568                deleted_rows_count: Some(0),
569                partitions: Some(
570                    vec![FieldSummary { contains_null: false, contains_nan: Some(false), lower_bound: Some(Datum::long(1).to_bytes().unwrap()), upper_bound: Some(Datum::long(1).to_bytes().unwrap())}]
571                ),
572                key_metadata: None,
573                first_row_id: None,
574            }]
575        };
576
577        let temp_dir = TempDir::new().unwrap();
578        let path = temp_dir.path().join("manifest_list_v2.avro");
579        let io = FileIO::new_with_fs();
580        let file_writer = file_writer(&path, io).await;
581
582        let mut writer = ManifestListWriter::v2(file_writer, snapshot_id, Some(0), seq_num);
583        writer
584            .add_manifests(expected_manifest_list.entries.clone().into_iter())
585            .unwrap();
586        writer.close().await.unwrap();
587
588        let bs = fs::read(path).unwrap();
589
590        let manifest_list =
591            ManifestList::parse_with_version(&bs, crate::spec::FormatVersion::V3).unwrap();
592        expected_manifest_list.entries[0].sequence_number = seq_num;
593        expected_manifest_list.entries[0].min_sequence_number = seq_num;
594        assert_eq!(manifest_list, expected_manifest_list);
595
596        temp_dir.close().unwrap();
597    }
598
599    #[tokio::test]
600    async fn test_manifest_list_writer_v3_encrypted_round_trip() {
601        let (mgr, file_io) = fresh_encryption_manager_and_io();
602        let path = "memory:///manifest_list_v3_encrypted.avro";
603
604        let encrypted_output = mgr.encrypt(file_io.new_output(path).unwrap());
605
606        let snapshot_id = 9_000_000_000_000_001i64;
607        let seq_num = 7i64;
608        let mut expected = ManifestList {
609            entries: vec![ManifestFile {
610                manifest_path: "memory:///encrypted/v3_m0.avro".to_string(),
611                manifest_length: 1234,
612                partition_spec_id: 0,
613                content: ManifestContentType::Data,
614                sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
615                min_sequence_number: UNASSIGNED_SEQUENCE_NUMBER,
616                added_snapshot_id: snapshot_id,
617                added_files_count: Some(2),
618                existing_files_count: Some(0),
619                deleted_files_count: Some(0),
620                added_rows_count: Some(10),
621                existing_rows_count: Some(0),
622                deleted_rows_count: Some(0),
623                partitions: Some(vec![FieldSummary {
624                    contains_null: false,
625                    contains_nan: Some(false),
626                    lower_bound: Some(Datum::long(1).to_bytes().unwrap()),
627                    upper_bound: Some(Datum::long(1).to_bytes().unwrap()),
628                }]),
629                key_metadata: None,
630                first_row_id: None,
631            }],
632        };
633
634        let file_writer = encrypted_output.writer().await.unwrap();
635        let mut writer = ManifestListWriter::v3(file_writer, snapshot_id, Some(0), seq_num, None);
636        writer
637            .add_manifests(expected.entries.clone().into_iter())
638            .unwrap();
639        let file_metadata = writer.close().await.unwrap();
640
641        let raw_bytes = file_io.new_input(path).unwrap().read().await.unwrap();
642        assert!(
643            ManifestList::parse_with_version(&raw_bytes, crate::spec::FormatVersion::V3).is_err(),
644            "raw bytes should be ciphertext, not parseable as Avro"
645        );
646
647        let key_metadata = encrypted_output.key_metadata_with_saved_file_metadata(&file_metadata);
648        assert_eq!(key_metadata.file_length(), Some(raw_bytes.len() as u64));
649        let plaintext = EncryptedInputFile::new(file_io.new_input(path).unwrap(), key_metadata)
650            .read()
651            .await
652            .unwrap();
653        let manifest_list =
654            ManifestList::parse_with_version(&plaintext, crate::spec::FormatVersion::V3).unwrap();
655
656        expected.entries[0].sequence_number = seq_num;
657        expected.entries[0].min_sequence_number = seq_num;
658        assert_eq!(manifest_list, expected);
659    }
660
661    fn fresh_encryption_manager_and_io() -> (EncryptionManager, FileIO) {
662        let kms = MemoryKeyManagementClient::new();
663        kms.add_master_key("master-1").unwrap();
664        let mgr = EncryptionManager::builder()
665            .kms_client(Arc::new(kms) as Arc<dyn KeyManagementClient>)
666            .table_key_id("master-1")
667            .build();
668        (mgr, FileIO::new_with_memory())
669    }
670
671    async fn file_writer(path: &Path, io: FileIO) -> Box<dyn FileWrite> {
672        io.new_output(path.to_str().unwrap())
673            .unwrap()
674            .writer()
675            .await
676            .unwrap()
677    }
678}