Skip to main content

iceberg/spec/manifest/
reader.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 super::Manifest;
19use crate::encryption::{EncryptedInputFile, StandardKeyMetadata};
20use crate::error::{Result, invalid_data};
21use crate::io::FileIO;
22use crate::spec::{ManifestContentType, ManifestEntry, ManifestFile};
23
24/// Reads a manifest file referenced by a manifest list entry, transparently
25/// decrypting it when the entry records key metadata.
26pub struct ManifestReader {
27    file_io: FileIO,
28}
29
30impl ManifestReader {
31    /// Create a manifest reader.
32    pub(crate) fn new(file_io: FileIO) -> Self {
33        Self { file_io }
34    }
35
36    /// Read, decrypt, parse and return the manifest described by
37    /// `manifest_file`.
38    pub async fn read(self, manifest_file: &ManifestFile) -> Result<Manifest> {
39        let input_file = self.file_io.new_input(&manifest_file.manifest_path)?;
40        let key_metadata = manifest_file
41            .key_metadata
42            .as_deref()
43            .map(StandardKeyMetadata::decode)
44            .transpose()?;
45        let bytes = match key_metadata {
46            Some(key_metadata) => {
47                EncryptedInputFile::new(input_file, key_metadata)
48                    .read()
49                    .await?
50            }
51            None => input_file.read().await?,
52        };
53
54        let (metadata, mut entries) =
55            Manifest::try_from_avro_bytes(&bytes, Some(&manifest_file.manifest_path))?;
56
57        for entry in &mut entries {
58            entry.inherit_data(manifest_file);
59        }
60
61        self.assign_first_row_ids(manifest_file, &mut entries)?;
62
63        Ok(Manifest::new(metadata, entries))
64    }
65
66    /// Assigns `first_row_id` to the live data-file entries, following the
67    /// row-lineage inheritance rules in
68    /// <https://github.com/apache/iceberg/blob/main/format/spec.md#first-row-id-inheritance>.
69    ///
70    /// With a manifest-level `first_row_id`, each live entry lacking one is
71    /// assigned the running id, which then advances by that entry's record
72    /// count; entries that already carry a `first_row_id` keep it and do not
73    /// advance the counter. Without a manifest-level `first_row_id`, any
74    /// inherited per-entry value is cleared so callers never observe a stale id.
75    fn assign_first_row_ids(
76        &self,
77        manifest_file: &ManifestFile,
78        entries: &mut [ManifestEntry],
79    ) -> Result<()> {
80        // A `first_row_id` is only valid on data manifests. Delete files always
81        // have a null `first_row_id`, so there is nothing to assign or clear; a
82        // stray value on a delete manifest is a spec violation by the writer,
83        // which we surface without failing the read.
84        if manifest_file.content != ManifestContentType::Data {
85            if let Some(manifest_first_row_id) = manifest_file.first_row_id {
86                tracing::warn!(
87                    "Ignoring first_row_id {manifest_first_row_id} on delete manifest {}",
88                    manifest_file.manifest_path
89                );
90            }
91
92            return Ok(());
93        }
94
95        let Some(manifest_first_row_id) = manifest_file.first_row_id else {
96            // A data manifest with no manifest-level `first_row_id` predates row
97            // lineage; clear any per-entry value inherited from an earlier read.
98            for entry in entries {
99                entry.data_file.first_row_id = None;
100            }
101
102            return Ok(());
103        };
104
105        let mut next_row_id = i64::try_from(manifest_first_row_id).map_err(|_| {
106            invalid_data!("Invalid first_row_id: {manifest_first_row_id} (exceeds i64::MAX)")
107        })?;
108
109        for entry in entries {
110            if !entry.is_alive() {
111                continue;
112            }
113
114            if entry.data_file.first_row_id.is_none() {
115                let file_first_row_id = next_row_id;
116                entry.data_file.first_row_id = Some(file_first_row_id);
117                let record_count = entry.data_file.record_count;
118                next_row_id = file_first_row_id.checked_add_unsigned(record_count).ok_or_else(|| {
119                    invalid_data!("Row ID overflow assigning first_row_id in {}. File first_row_id: {file_first_row_id}, record count: {record_count}",
120                            manifest_file.manifest_path)
121                })?;
122            }
123        }
124
125        Ok(())
126    }
127}
128
129#[cfg(test)]
130mod tests {
131    use std::collections::HashMap;
132    use std::sync::Arc;
133
134    use super::*;
135    use crate::ErrorKind;
136    use crate::encryption::{EncryptedOutputFile, StandardKeyMetadata};
137    use crate::io::FileIO;
138    use crate::spec::{
139        DataContentType, DataFile, DataFileBuilder, DataFileFormat, ManifestEntry, ManifestStatus,
140        ManifestWriterBuilder, NestedField, PartitionSpec, PrimitiveType, Schema, Struct, Type,
141    };
142
143    #[tokio::test]
144    async fn test_read_plaintext_manifest_inherits_entries() {
145        let schema = test_schema();
146        let partition_spec = PartitionSpec::builder(schema.clone())
147            .with_spec_id(0)
148            .build()
149            .unwrap();
150
151        let io = FileIO::new_with_memory();
152        let path = "memory:///table/metadata/plain.avro";
153        let mut writer = ManifestWriterBuilder::new(
154            io.new_output(path).unwrap(),
155            Some(1),
156            schema.clone(),
157            partition_spec,
158        )
159        .build_v2_data();
160        writer.add_entry(test_entry()).unwrap();
161        // Writing the manifest yields the manifest list entry describing it.
162        let manifest_file = writer.write_manifest_file().await.unwrap();
163
164        let manifest = ManifestReader::new(io).read(&manifest_file).await.unwrap();
165        assert_eq!(manifest.entries().len(), 1);
166        assert_eq!(
167            manifest.entries()[0].data_file().file_path(),
168            "memory:///table/data/00000.parquet"
169        );
170        // Entries must inherit values from the manifest list entry.
171        assert_eq!(
172            manifest.entries()[0].sequence_number(),
173            Some(manifest_file.sequence_number)
174        );
175        assert_eq!(
176            manifest.entries()[0].snapshot_id(),
177            Some(manifest_file.added_snapshot_id)
178        );
179    }
180
181    #[tokio::test]
182    async fn test_read_encrypted_manifest_roundtrip() {
183        let schema = test_schema();
184        let partition_spec = PartitionSpec::builder(schema.clone())
185            .with_spec_id(0)
186            .build()
187            .unwrap();
188
189        let io = FileIO::new_with_memory();
190        let path = "memory:///table/metadata/encrypted.avro";
191        let encrypted_output =
192            EncryptedOutputFile::new(io.new_output(path).unwrap(), key_metadata());
193
194        let mut writer = ManifestWriterBuilder::new_from_encrypted(
195            encrypted_output,
196            Some(1),
197            schema.clone(),
198            partition_spec,
199        )
200        .unwrap()
201        .build_v3_data();
202        writer.add_entry(test_entry()).unwrap();
203        // The returned manifest list entry records the key metadata.
204        let manifest_file = writer.write_manifest_file().await.unwrap();
205        assert!(manifest_file.key_metadata.is_some());
206
207        // Reading with the recorded key metadata must recover the entry.
208        let manifest = ManifestReader::new(io.clone())
209            .read(&manifest_file)
210            .await
211            .unwrap();
212        assert_eq!(manifest.entries().len(), 1);
213        assert_eq!(
214            manifest.entries()[0].data_file().file_path(),
215            "memory:///table/data/00000.parquet"
216        );
217
218        // Without the key metadata the encrypted bytes must not read as plaintext.
219        let mut plaintext_entry = manifest_file.clone();
220        plaintext_entry.key_metadata = None;
221        assert!(
222            ManifestReader::new(io)
223                .read(&plaintext_entry)
224                .await
225                .is_err(),
226            "encrypted manifest must not parse as plaintext"
227        );
228    }
229
230    fn test_schema() -> Arc<Schema> {
231        Arc::new(
232            Schema::builder()
233                .with_fields(vec![Arc::new(NestedField::optional(
234                    1,
235                    "id",
236                    Type::Primitive(PrimitiveType::Long),
237                ))])
238                .build()
239                .unwrap(),
240        )
241    }
242
243    fn test_entry() -> ManifestEntry {
244        ManifestEntry {
245            status: ManifestStatus::Added,
246            snapshot_id: None,
247            sequence_number: None,
248            file_sequence_number: None,
249            data_file: DataFile {
250                content: DataContentType::Data,
251                file_path: "memory:///table/data/00000.parquet".to_string(),
252                file_format: DataFileFormat::Parquet,
253                partition: Struct::empty(),
254                record_count: 1,
255                file_size_in_bytes: 4096,
256                column_sizes: HashMap::new(),
257                value_counts: HashMap::new(),
258                null_value_counts: HashMap::new(),
259                nan_value_counts: HashMap::new(),
260                lower_bounds: HashMap::new(),
261                upper_bounds: HashMap::new(),
262                key_metadata: None,
263                split_offsets: None,
264                equality_ids: None,
265                sort_order_id: None,
266                partition_spec_id: 0,
267                first_row_id: None,
268                referenced_data_file: None,
269                content_offset: None,
270                content_size_in_bytes: None,
271            },
272        }
273    }
274
275    fn key_metadata() -> StandardKeyMetadata {
276        StandardKeyMetadata::try_new(b"0123456789abcdef").unwrap()
277    }
278
279    /// A reader with an unused in-memory `FileIO`, for exercising the pure
280    /// row-id assignment logic without touching storage.
281    fn test_reader() -> ManifestReader {
282        ManifestReader::new(FileIO::new_with_memory())
283    }
284
285    /// Builds a data-file manifest entry with the given status, record count,
286    /// and pre-existing `first_row_id`.
287    fn data_entry(
288        status: ManifestStatus,
289        record_count: u64,
290        first_row_id: Option<i64>,
291    ) -> ManifestEntry {
292        let data_file = DataFileBuilder::default()
293            .content(DataContentType::Data)
294            .file_path("s3://bucket/table/data/00000.parquet".to_string())
295            .file_format(DataFileFormat::Parquet)
296            .file_size_in_bytes(4096)
297            .record_count(record_count)
298            .first_row_id(first_row_id)
299            .build()
300            .unwrap();
301
302        ManifestEntry::builder()
303            .status(status)
304            .data_file(data_file)
305            .build()
306    }
307
308    /// Builds a manifest file with the given content type and manifest-level
309    /// `first_row_id`. Other fields are irrelevant to row-id assignment.
310    fn manifest_file(content: ManifestContentType, first_row_id: Option<u64>) -> ManifestFile {
311        ManifestFile {
312            manifest_path: "memory:///m.avro".to_string(),
313            manifest_length: 0,
314            partition_spec_id: 0,
315            content,
316            sequence_number: 0,
317            min_sequence_number: 0,
318            added_snapshot_id: 0,
319            added_files_count: None,
320            existing_files_count: None,
321            deleted_files_count: None,
322            added_rows_count: None,
323            existing_rows_count: None,
324            deleted_rows_count: None,
325            partitions: None,
326            key_metadata: None,
327            first_row_id,
328        }
329    }
330
331    #[test]
332    fn test_assign_first_row_ids_interleaved() {
333        let manifest = manifest_file(ManifestContentType::Data, Some(10));
334        let mut entries = vec![
335            data_entry(ManifestStatus::Added, 3, None),
336            // A pre-assigned entry between two assigned ones: it keeps its id and
337            // must not advance the running counter.
338            data_entry(ManifestStatus::Added, 5, Some(100)),
339            // A deleted entry with a pre-set id: it is skipped, so the id is
340            // preserved verbatim and does not advance the counter.
341            data_entry(ManifestStatus::Deleted, 7, Some(999)),
342            data_entry(ManifestStatus::Existing, 2, None),
343        ];
344
345        test_reader()
346            .assign_first_row_ids(&manifest, &mut entries)
347            .unwrap();
348
349        assert_eq!(entries[0].data_file.first_row_id, Some(10));
350        assert_eq!(entries[1].data_file.first_row_id, Some(100));
351        assert_eq!(entries[2].data_file.first_row_id, Some(999));
352        // 10 + 3 = 13; the preserved and deleted entries in between do not move it.
353        assert_eq!(entries[3].data_file.first_row_id, Some(13));
354    }
355
356    #[test]
357    fn test_assign_first_row_ids_clears_without_manifest_first_row_id() {
358        // A data manifest with no manifest-level first_row_id predates row
359        // lineage: any per-entry value inherited from an earlier read is cleared
360        // so callers never observe a stale id.
361        let manifest = manifest_file(ManifestContentType::Data, None);
362        let mut entries = vec![
363            data_entry(ManifestStatus::Added, 3, None),
364            data_entry(ManifestStatus::Existing, 5, Some(100)),
365        ];
366
367        test_reader()
368            .assign_first_row_ids(&manifest, &mut entries)
369            .unwrap();
370
371        assert_eq!(entries[0].data_file.first_row_id, None);
372        assert_eq!(entries[1].data_file.first_row_id, None);
373    }
374
375    #[test]
376    fn test_assign_first_row_ids_ignores_delete_manifest() {
377        // A stray first_row_id on a delete manifest is a writer-side spec
378        // violation; the read ignores it rather than failing, and does not
379        // assign ids to the entries.
380        let manifest = manifest_file(ManifestContentType::Deletes, Some(10));
381        let mut entries = vec![data_entry(ManifestStatus::Added, 3, None)];
382
383        test_reader()
384            .assign_first_row_ids(&manifest, &mut entries)
385            .unwrap();
386
387        assert_eq!(entries[0].data_file.first_row_id, None);
388    }
389
390    #[test]
391    fn test_assign_first_row_ids_rejects_oversized_manifest_first_row_id() {
392        // A manifest-level first_row_id above i64::MAX cannot be represented as the
393        // signed running counter and must be rejected.
394        let manifest = manifest_file(ManifestContentType::Data, Some(i64::MAX as u64 + 1));
395        let mut entries = vec![data_entry(ManifestStatus::Added, 3, None)];
396
397        let err = test_reader()
398            .assign_first_row_ids(&manifest, &mut entries)
399            .expect_err("an oversized manifest first_row_id must be rejected");
400        assert_eq!(err.kind(), ErrorKind::DataInvalid);
401        assert!(err.message().contains("Invalid first_row_id"));
402    }
403
404    #[test]
405    fn test_assign_first_row_ids_rejects_counter_overflow() {
406        // Advancing the running counter past i64::MAX must be rejected rather than
407        // wrapping to a negative value that would corrupt subsequent assignments.
408        let manifest = manifest_file(ManifestContentType::Data, Some(i64::MAX as u64));
409        let mut entries = vec![data_entry(ManifestStatus::Added, 1, None)];
410
411        let err = test_reader()
412            .assign_first_row_ids(&manifest, &mut entries)
413            .expect_err("counter overflow past i64::MAX must be rejected");
414        assert_eq!(err.kind(), ErrorKind::DataInvalid);
415        assert!(err.message().contains("Row ID overflow"));
416    }
417}