1use super::Manifest;
19use crate::encryption::{EncryptedInputFile, StandardKeyMetadata};
20use crate::error::{Result, invalid_data};
21use crate::io::FileIO;
22use crate::spec::{ManifestContentType, ManifestEntry, ManifestFile};
23
24pub struct ManifestReader {
27 file_io: FileIO,
28}
29
30impl ManifestReader {
31 pub(crate) fn new(file_io: FileIO) -> Self {
33 Self { file_io }
34 }
35
36 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 fn assign_first_row_ids(
76 &self,
77 manifest_file: &ManifestFile,
78 entries: &mut [ManifestEntry],
79 ) -> Result<()> {
80 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 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 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 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 let manifest_file = writer.write_manifest_file().await.unwrap();
205 assert!(manifest_file.key_metadata.is_some());
206
207 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 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 fn test_reader() -> ManifestReader {
282 ManifestReader::new(FileIO::new_with_memory())
283 }
284
285 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 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 data_entry(ManifestStatus::Added, 5, Some(100)),
339 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 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 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 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 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 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}