From e3f755fcd0adc5f8e0d34852cb66826f5b2e676c Mon Sep 17 00:00:00 2001 From: Xander Date: Mon, 13 Jul 2026 18:27:56 +0100 Subject: [PATCH 1/3] Encrypt / Decrypt puffin files --- crates/iceberg/public-api.txt | 4 +- crates/iceberg/src/puffin/metadata.rs | 114 ++++++++++++++------------ crates/iceberg/src/puffin/reader.rs | 41 ++++++--- crates/iceberg/src/puffin/writer.rs | 104 +++++++++++++++++++++-- 4 files changed, 190 insertions(+), 73 deletions(-) diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 610675c83d..3fc4cb824a 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -1248,12 +1248,14 @@ pub struct iceberg::puffin::PuffinReader impl iceberg::puffin::PuffinReader pub async fn iceberg::puffin::PuffinReader::blob(&self, blob_metadata: &iceberg::puffin::BlobMetadata) -> iceberg::Result pub async fn iceberg::puffin::PuffinReader::file_metadata(&self) -> iceberg::Result<&iceberg::puffin::FileMetadata> -pub fn iceberg::puffin::PuffinReader::new(input_file: iceberg::io::InputFile) -> Self +pub async fn iceberg::puffin::PuffinReader::new(input_file: iceberg::io::InputFile) -> iceberg::Result +pub async fn iceberg::puffin::PuffinReader::new_from_encrypted(encrypted_input: iceberg::encryption::EncryptedInputFile) -> iceberg::Result pub struct iceberg::puffin::PuffinWriter impl iceberg::puffin::PuffinWriter pub async fn iceberg::puffin::PuffinWriter::add(&mut self, blob: iceberg::puffin::Blob, compression_codec: iceberg::compression::CompressionCodec) -> iceberg::Result<()> pub async fn iceberg::puffin::PuffinWriter::close(self) -> iceberg::Result<()> pub async fn iceberg::puffin::PuffinWriter::new(output_file: &iceberg::io::OutputFile, properties: std::collections::hash::map::HashMap, compress_footer: bool) -> iceberg::Result +pub async fn iceberg::puffin::PuffinWriter::new_from_encrypted(encrypted_output: &iceberg::encryption::EncryptedOutputFile, properties: std::collections::hash::map::HashMap, compress_footer: bool) -> iceberg::Result pub const iceberg::puffin::APACHE_DATASKETCHES_THETA_V1: &str pub const iceberg::puffin::CREATED_BY_PROPERTY: &str pub const iceberg::puffin::DELETION_VECTOR_V1: &str diff --git a/crates/iceberg/src/puffin/metadata.rs b/crates/iceberg/src/puffin/metadata.rs index 35984a5ef1..51b317161f 100644 --- a/crates/iceberg/src/puffin/metadata.rs +++ b/crates/iceberg/src/puffin/metadata.rs @@ -21,7 +21,7 @@ use bytes::Bytes; use serde::{Deserialize, Serialize}; use crate::compression::CompressionCodec; -use crate::io::{FileRead, InputFile}; +use crate::io::FileRead; use crate::{Error, ErrorKind, Result}; /// Human-readable identification of the application writing the file, along with its version. @@ -291,16 +291,13 @@ impl FileMetadata { } /// Returns the file metadata about a Puffin file - pub(crate) async fn read(input_file: &InputFile) -> Result { - let file_read = input_file.reader().await?; - - let input_file_length = input_file.metadata().await?.size; - if input_file_length < FileMetadata::MIN_FILE_LENGTH { + pub(crate) async fn read(file_read: &dyn FileRead, file_length: u64) -> Result { + if file_length < FileMetadata::MIN_FILE_LENGTH { return Err(Error::new( ErrorKind::DataInvalid, format!( "File length {} is too short to be a Puffin file, expected at least {} bytes", - input_file_length, + file_length, FileMetadata::MIN_FILE_LENGTH ), )); @@ -310,13 +307,9 @@ impl FileMetadata { FileMetadata::check_magic(&first_four_bytes)?; let footer_payload_length = - FileMetadata::read_footer_payload_length(file_read.as_ref(), input_file_length).await?; - let footer_bytes = FileMetadata::read_footer_bytes( - file_read.as_ref(), - input_file_length, - footer_payload_length, - ) - .await?; + FileMetadata::read_footer_payload_length(file_read, file_length).await?; + let footer_bytes = + FileMetadata::read_footer_bytes(file_read, file_length, footer_payload_length).await?; let magic_length = FileMetadata::MAGIC_LENGTH as usize; // check first four bytes of footer @@ -337,16 +330,14 @@ impl FileMetadata { /// read option. #[allow(dead_code)] pub(crate) async fn read_with_prefetch( - input_file: &InputFile, + file_read: &dyn FileRead, + file_length: u64, prefetch_hint: u8, ) -> Result { if prefetch_hint > 16 { - let input_file_length = input_file.metadata().await?.size; - let file_read = input_file.reader().await?; - // Hint cannot be larger than input file - if prefetch_hint as u64 > input_file_length { - return FileMetadata::read(input_file).await; + if prefetch_hint as u64 > file_length { + return FileMetadata::read(file_read, file_length).await; } // Validate file header magic @@ -354,8 +345,8 @@ impl FileMetadata { FileMetadata::check_magic(&first_four_bytes)?; // Read footer based on prefetch hint - let start = input_file_length - prefetch_hint as u64; - let end = input_file_length; + let start = file_length - prefetch_hint as u64; + let end = file_length; let footer_bytes = file_read.read(start..end).await?; let payload_length_start = @@ -375,7 +366,7 @@ impl FileMetadata { + FileMetadata::FOOTER_STRUCT_LENGTH as usize + FileMetadata::MAGIC_LENGTH as usize; if footer_length > prefetch_hint as usize { - return FileMetadata::read(input_file).await; + return FileMetadata::read(file_read, file_length).await; } // Read footer bytes @@ -394,7 +385,7 @@ impl FileMetadata { return FileMetadata::from_json_str(&footer_payload_str); } - FileMetadata::read(input_file).await + FileMetadata::read(file_read, file_length).await } #[inline] @@ -417,7 +408,6 @@ mod tests { use bytes::Bytes; use tempfile::TempDir; - use crate::ErrorKind; use crate::io::{FileIO, InputFile}; use crate::puffin::metadata::{BlobMetadata, CompressionCodec, FileMetadata}; use crate::puffin::test_utils::{ @@ -426,9 +416,27 @@ mod tests { java_zstd_compressed_metric_input_file, uncompressed_metric_file_metadata, zstd_compressed_metric_file_metadata, }; + use crate::{ErrorKind, Result}; const INVALID_MAGIC_VALUE: [u8; 4] = [80, 70, 65, 0]; + /// Reads file metadata from an [`InputFile`], resolving its reader and length. + async fn read_file_metadata(input_file: &InputFile) -> Result { + let file_read = input_file.reader().await?; + let file_length = input_file.metadata().await?.size; + FileMetadata::read(file_read.as_ref(), file_length).await + } + + /// Reads file metadata with a prefetch hint from an [`InputFile`]. + async fn read_file_metadata_with_prefetch( + input_file: &InputFile, + prefetch_hint: u8, + ) -> Result { + let file_read = input_file.reader().await?; + let file_length = input_file.metadata().await?.size; + FileMetadata::read_with_prefetch(file_read.as_ref(), file_length, prefetch_hint).await + } + async fn input_file_with_bytes(temp_dir: &TempDir, slice: &[u8]) -> InputFile { let file_io = FileIO::new_with_fs(); @@ -473,7 +481,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -496,7 +504,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -519,7 +527,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -544,7 +552,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -569,7 +577,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -592,7 +600,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -615,7 +623,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -641,7 +649,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -656,7 +664,7 @@ mod tests { // Only the file header magic, nothing else. let input_file = input_file_with_bytes(&temp_dir, &FileMetadata::MAGIC).await; - let err = FileMetadata::read(&input_file).await.unwrap_err(); + let err = read_file_metadata(&input_file).await.unwrap_err(); assert_eq!(err.kind(), ErrorKind::DataInvalid); assert!( err.to_string().contains("too short to be a Puffin file"), @@ -679,7 +687,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; - let err = FileMetadata::read(&input_file).await.unwrap_err(); + let err = read_file_metadata(&input_file).await.unwrap_err(); assert_eq!(err.kind(), ErrorKind::DataInvalid); assert!( err.to_string().contains("exceeds file length"), @@ -702,7 +710,7 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await.unwrap(), FileMetadata { blobs: vec![], properties: HashMap::new(), @@ -726,7 +734,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await.unwrap(), FileMetadata { blobs: vec![], properties: { @@ -755,7 +763,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await.unwrap(), FileMetadata { blobs: vec![], properties: { @@ -781,7 +789,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -804,7 +812,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -844,7 +852,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await.unwrap(), FileMetadata { blobs: vec![ BlobMetadata { @@ -898,7 +906,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await.unwrap(), FileMetadata { blobs: vec![BlobMetadata { r#type: "type-a".to_string(), @@ -945,7 +953,7 @@ mod tests { .await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -962,7 +970,7 @@ mod tests { let input_file = input_file_with_payload(&temp_dir, r#""blobs" = []"#).await; assert_eq!( - FileMetadata::read(&input_file) + read_file_metadata(&input_file) .await .unwrap_err() .to_string(), @@ -974,7 +982,7 @@ mod tests { async fn test_read_file_metadata_of_uncompressed_empty_file() { let input_file = java_empty_uncompressed_input_file(); - let file_metadata = FileMetadata::read(&input_file).await.unwrap(); + let file_metadata = read_file_metadata(&input_file).await.unwrap(); assert_eq!(file_metadata, empty_footer_payload()) } @@ -982,7 +990,7 @@ mod tests { async fn test_read_file_metadata_of_uncompressed_metric_data() { let input_file = java_uncompressed_metric_input_file(); - let file_metadata = FileMetadata::read(&input_file).await.unwrap(); + let file_metadata = read_file_metadata(&input_file).await.unwrap(); assert_eq!(file_metadata, uncompressed_metric_file_metadata()) } @@ -990,7 +998,7 @@ mod tests { async fn test_read_file_metadata_of_zstd_compressed_metric_data() { let input_file = java_zstd_compressed_metric_input_file(); - let file_metadata = FileMetadata::read_with_prefetch(&input_file, 64) + let file_metadata = read_file_metadata_with_prefetch(&input_file, 64) .await .unwrap(); assert_eq!(file_metadata, zstd_compressed_metric_file_metadata()) @@ -999,7 +1007,7 @@ mod tests { #[tokio::test] async fn test_read_file_metadata_of_empty_file_with_prefetching() { let input_file = java_empty_uncompressed_input_file(); - let file_metadata = FileMetadata::read_with_prefetch(&input_file, 64) + let file_metadata = read_file_metadata_with_prefetch(&input_file, 64) .await .unwrap(); @@ -1009,7 +1017,7 @@ mod tests { #[tokio::test] async fn test_read_file_metadata_of_uncompressed_metric_data_with_prefetching() { let input_file = java_uncompressed_metric_input_file(); - let file_metadata = FileMetadata::read_with_prefetch(&input_file, 64) + let file_metadata = read_file_metadata_with_prefetch(&input_file, 64) .await .unwrap(); @@ -1019,7 +1027,7 @@ mod tests { #[tokio::test] async fn test_read_file_metadata_of_zstd_compressed_metric_data_with_prefetching() { let input_file = java_zstd_compressed_metric_input_file(); - let file_metadata = FileMetadata::read_with_prefetch(&input_file, 64) + let file_metadata = read_file_metadata_with_prefetch(&input_file, 64) .await .unwrap(); @@ -1046,11 +1054,11 @@ mod tests { let input_file = input_file_with_bytes(&temp_dir, &bytes).await; assert_eq!( - FileMetadata::read(&input_file).await.unwrap_err().kind(), + read_file_metadata(&input_file).await.unwrap_err().kind(), ErrorKind::DataInvalid, ); assert_eq!( - FileMetadata::read_with_prefetch(&input_file, prefetch_hint) + read_file_metadata_with_prefetch(&input_file, prefetch_hint) .await .unwrap_err() .kind(), @@ -1081,7 +1089,7 @@ mod tests { let input_file = input_file_with_payload(&temp_dir, payload).await; // Reading metadata should succeed (lazy validation) - let result = FileMetadata::read(&input_file).await; + let result = read_file_metadata(&input_file).await; assert!(result.is_ok()); let metadata = result.unwrap(); assert_eq!(metadata.blobs.len(), 1); diff --git a/crates/iceberg/src/puffin/reader.rs b/crates/iceberg/src/puffin/reader.rs index 0aced4186f..01601f48e0 100644 --- a/crates/iceberg/src/puffin/reader.rs +++ b/crates/iceberg/src/puffin/reader.rs @@ -19,21 +19,41 @@ use tokio::sync::OnceCell; use super::validate_puffin_compression; use crate::Result; -use crate::io::InputFile; +use crate::encryption::EncryptedInputFile; +use crate::io::{FileRead, InputFile}; use crate::puffin::blob::Blob; use crate::puffin::metadata::{BlobMetadata, FileMetadata}; /// Puffin reader pub struct PuffinReader { - input_file: InputFile, + file_read: Box, + file_length: u64, file_metadata: OnceCell, } impl PuffinReader { - /// Returns a new Puffin reader - pub fn new(input_file: InputFile) -> Self { + /// Returns a new Puffin reader for an unencrypted file. + pub async fn new(input_file: InputFile) -> Result { + let file_length = input_file.metadata().await?.size; + let file_read = input_file.reader().await?; + Ok(Self::from_parts(file_read, file_length)) + } + + /// Returns a new Puffin reader from an [`EncryptedInputFile`]. + /// + /// Use this when reading Puffin files with transparent decryption. The + /// reader operates over plaintext offsets and length, so all blob and + /// footer positions match those written to the unencrypted file. + pub async fn new_from_encrypted(encrypted_input: EncryptedInputFile) -> Result { + let file_length = encrypted_input.metadata().await?.size; + let file_read = encrypted_input.reader().await?; + Ok(Self::from_parts(file_read, file_length)) + } + + fn from_parts(file_read: Box, file_length: u64) -> Self { Self { - input_file, + file_read, + file_length, file_metadata: OnceCell::new(), } } @@ -41,7 +61,7 @@ impl PuffinReader { /// Returns file metadata pub async fn file_metadata(&self) -> Result<&FileMetadata> { self.file_metadata - .get_or_try_init(|| FileMetadata::read(&self.input_file)) + .get_or_try_init(|| FileMetadata::read(self.file_read.as_ref(), self.file_length)) .await } @@ -49,10 +69,9 @@ impl PuffinReader { pub async fn blob(&self, blob_metadata: &BlobMetadata) -> Result { validate_puffin_compression(blob_metadata.compression_codec)?; - let file_read = self.input_file.reader().await?; let start = blob_metadata.offset; let end = start + blob_metadata.length; - let bytes = file_read.read(start..end).await?; + let bytes = self.file_read.read(start..end).await?; let data = blob_metadata.compression_codec.decompress(bytes.to_vec())?; Ok(Blob { @@ -83,7 +102,7 @@ mod tests { #[tokio::test] async fn test_puffin_reader_uncompressed_metric_data() { let input_file = java_uncompressed_metric_input_file(); - let puffin_reader = PuffinReader::new(input_file); + let puffin_reader = PuffinReader::new(input_file).await.unwrap(); let file_metadata = puffin_reader.file_metadata().await.unwrap().clone(); assert_eq!(file_metadata, uncompressed_metric_file_metadata()); @@ -108,7 +127,7 @@ mod tests { #[tokio::test] async fn test_puffin_reader_zstd_compressed_metric_data() { let input_file = java_zstd_compressed_metric_input_file(); - let puffin_reader = PuffinReader::new(input_file); + let puffin_reader = PuffinReader::new(input_file).await.unwrap(); let file_metadata = puffin_reader.file_metadata().await.unwrap().clone(); assert_eq!(file_metadata, zstd_compressed_metric_file_metadata()); @@ -134,7 +153,7 @@ mod tests { async fn test_gzip_compression_rejected_on_blob_access() { // Use a real puffin file let input_file = java_uncompressed_metric_input_file(); - let reader = PuffinReader::new(input_file); + let reader = PuffinReader::new(input_file).await.unwrap(); // Create a BlobMetadata with Gzip compression let gzip_blob_metadata = BlobMetadata { diff --git a/crates/iceberg/src/puffin/writer.rs b/crates/iceberg/src/puffin/writer.rs index 4af4970b04..e1a84489eb 100644 --- a/crates/iceberg/src/puffin/writer.rs +++ b/crates/iceberg/src/puffin/writer.rs @@ -22,6 +22,7 @@ use bytes::Bytes; use super::validate_puffin_compression; use crate::Result; use crate::compression::CompressionCodec; +use crate::encryption::EncryptedOutputFile; use crate::io::{FileWrite, OutputFile}; use crate::puffin::blob::Blob; use crate::puffin::metadata::{BlobMetadata, FileMetadata, Flag}; @@ -38,12 +39,41 @@ pub struct PuffinWriter { } impl PuffinWriter { - /// Returns a new Puffin writer + /// Returns a new Puffin writer for an unencrypted file. pub async fn new( output_file: &OutputFile, properties: HashMap, compress_footer: bool, ) -> Result { + Ok(Self::from_writer( + output_file.writer().await?, + properties, + compress_footer, + )) + } + + /// Returns a new Puffin writer from an [`EncryptedOutputFile`]. + /// + /// Use this when writing Puffin files with transparent encryption. Blob + /// and footer offsets are recorded as plaintext positions, matching what + /// an unencrypted writer would produce. + pub async fn new_from_encrypted( + encrypted_output: &EncryptedOutputFile, + properties: HashMap, + compress_footer: bool, + ) -> Result { + Ok(Self::from_writer( + encrypted_output.writer().await?, + properties, + compress_footer, + )) + } + + fn from_writer( + writer: Box, + properties: HashMap, + compress_footer: bool, + ) -> Self { let mut flags = HashSet::::new(); let footer_compression_codec = if compress_footer { flags.insert(Flag::FooterPayloadCompressed); @@ -52,15 +82,15 @@ impl PuffinWriter { CompressionCodec::None }; - Ok(Self { - writer: output_file.writer().await?, + Self { + writer, is_header_written: false, num_bytes_written: 0, written_blobs_metadata: Vec::new(), properties, footer_compression_codec, flags, - }) + } } /// Adds blob to Puffin file @@ -184,8 +214,16 @@ mod tests { Ok(output_file) } + async fn read_file_metadata(input_file: &InputFile) -> FileMetadata { + let file_read = input_file.reader().await.unwrap(); + let file_length = input_file.metadata().await.unwrap().size; + FileMetadata::read(file_read.as_ref(), file_length) + .await + .unwrap() + } + async fn read_all_blobs_from_puffin_file(input_file: InputFile) -> Vec { - let puffin_reader = PuffinReader::new(input_file); + let puffin_reader = PuffinReader::new(input_file).await.unwrap(); let mut blobs = Vec::new(); let blobs_metadata = puffin_reader.file_metadata().await.unwrap().clone().blobs; for blob_metadata in blobs_metadata { @@ -204,7 +242,7 @@ mod tests { .to_input_file(); assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await, empty_footer_payload() ); @@ -240,7 +278,7 @@ mod tests { .to_input_file(); assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await, uncompressed_metric_file_metadata() ); @@ -260,7 +298,7 @@ mod tests { .to_input_file(); assert_eq!( - FileMetadata::read(&input_file).await.unwrap(), + read_file_metadata(&input_file).await, zstd_compressed_metric_file_metadata() ); @@ -354,4 +392,54 @@ mod tests { .contains("is not supported for Puffin files") ); } + + #[tokio::test] + async fn test_encrypted_write_read_roundtrip() { + use crate::encryption::{EncryptedInputFile, EncryptedOutputFile, StandardKeyMetadata}; + + let key_metadata = || { + StandardKeyMetadata::try_new(b"0123456789abcdef") + .unwrap() + .with_aad_prefix(b"test-aad-prefix!") + }; + + let file_io = FileIO::new_with_memory(); + let path = "memory:///test/encrypted.puffin"; + let blobs = vec![blob_0(), blob_1()]; + + // Write through the encrypting writer. + let encrypted_output = + EncryptedOutputFile::new(file_io.new_output(path).unwrap(), key_metadata()); + let mut writer = + PuffinWriter::new_from_encrypted(&encrypted_output, file_properties(), false) + .await + .unwrap(); + for blob in blobs.clone() { + writer.add(blob, CompressionCodec::None).await.unwrap(); + } + writer.close().await.unwrap(); + + // The ciphertext on disk must not equal a plaintext puffin file. + let raw = file_io.new_input(path).unwrap().read().await.unwrap(); + assert_ne!( + &raw[..FileMetadata::MAGIC_LENGTH as usize], + FileMetadata::MAGIC + ); + + // Read back through the decrypting reader over plaintext offsets. + let encrypted_input = + EncryptedInputFile::new(file_io.new_input(path).unwrap(), key_metadata()); + let reader = PuffinReader::new_from_encrypted(encrypted_input) + .await + .unwrap(); + + let file_metadata = reader.file_metadata().await.unwrap().clone(); + assert_eq!(file_metadata, uncompressed_metric_file_metadata()); + + let mut read_blobs = Vec::new(); + for blob_metadata in &file_metadata.blobs { + read_blobs.push(reader.blob(blob_metadata).await.unwrap()); + } + assert_eq!(read_blobs, blobs); + } } From e7d211e395cea6e4b3a250985e52caeb42067b48 Mon Sep 17 00:00:00 2001 From: Xander Date: Tue, 28 Jul 2026 23:18:16 +0100 Subject: [PATCH 2/3] refactor tests --- crates/iceberg/src/puffin/metadata.rs | 22 +++------------------- crates/iceberg/src/puffin/test_utils.rs | 18 ++++++++++++++++++ crates/iceberg/src/puffin/writer.rs | 18 +++++------------- 3 files changed, 26 insertions(+), 32 deletions(-) diff --git a/crates/iceberg/src/puffin/metadata.rs b/crates/iceberg/src/puffin/metadata.rs index 51b317161f..78d904f066 100644 --- a/crates/iceberg/src/puffin/metadata.rs +++ b/crates/iceberg/src/puffin/metadata.rs @@ -413,30 +413,14 @@ mod tests { use crate::puffin::test_utils::{ empty_footer_payload, empty_footer_payload_bytes, empty_footer_payload_bytes_length_bytes, java_empty_uncompressed_input_file, java_uncompressed_metric_input_file, - java_zstd_compressed_metric_input_file, uncompressed_metric_file_metadata, + java_zstd_compressed_metric_input_file, read_file_metadata, + read_file_metadata_with_prefetch, uncompressed_metric_file_metadata, zstd_compressed_metric_file_metadata, }; - use crate::{ErrorKind, Result}; + use crate::ErrorKind; const INVALID_MAGIC_VALUE: [u8; 4] = [80, 70, 65, 0]; - /// Reads file metadata from an [`InputFile`], resolving its reader and length. - async fn read_file_metadata(input_file: &InputFile) -> Result { - let file_read = input_file.reader().await?; - let file_length = input_file.metadata().await?.size; - FileMetadata::read(file_read.as_ref(), file_length).await - } - - /// Reads file metadata with a prefetch hint from an [`InputFile`]. - async fn read_file_metadata_with_prefetch( - input_file: &InputFile, - prefetch_hint: u8, - ) -> Result { - let file_read = input_file.reader().await?; - let file_length = input_file.metadata().await?.size; - FileMetadata::read_with_prefetch(file_read.as_ref(), file_length, prefetch_hint).await - } - async fn input_file_with_bytes(temp_dir: &TempDir, slice: &[u8]) -> InputFile { let file_io = FileIO::new_with_fs(); diff --git a/crates/iceberg/src/puffin/test_utils.rs b/crates/iceberg/src/puffin/test_utils.rs index e0844e2002..26a6834a65 100644 --- a/crates/iceberg/src/puffin/test_utils.rs +++ b/crates/iceberg/src/puffin/test_utils.rs @@ -18,6 +18,7 @@ use std::collections::HashMap; use super::blob::Blob; +use crate::Result; use crate::compression::CompressionCodec; use crate::io::{FileIO, InputFile}; use crate::puffin::metadata::{BlobMetadata, CREATED_BY_PROPERTY, FileMetadata}; @@ -33,6 +34,23 @@ fn input_file_for_test_data(path: &str) -> InputFile { .unwrap() } +/// Reads file metadata from an [`InputFile`], resolving its reader and length. +pub(crate) async fn read_file_metadata(input_file: &InputFile) -> Result { + let file_read = input_file.reader().await?; + let file_length = input_file.metadata().await?.size; + FileMetadata::read(file_read.as_ref(), file_length).await +} + +/// Reads file metadata with a prefetch hint from an [`InputFile`]. +pub(crate) async fn read_file_metadata_with_prefetch( + input_file: &InputFile, + prefetch_hint: u8, +) -> Result { + let file_read = input_file.reader().await?; + let file_length = input_file.metadata().await?.size; + FileMetadata::read_with_prefetch(file_read.as_ref(), file_length, prefetch_hint).await +} + pub(crate) fn java_empty_uncompressed_input_file() -> InputFile { input_file_for_test_data(&[JAVA_TESTDATA, EMPTY_UNCOMPRESSED].join("/")) } diff --git a/crates/iceberg/src/puffin/writer.rs b/crates/iceberg/src/puffin/writer.rs index e1a84489eb..0437bd5bf0 100644 --- a/crates/iceberg/src/puffin/writer.rs +++ b/crates/iceberg/src/puffin/writer.rs @@ -188,8 +188,8 @@ mod tests { use crate::puffin::test_utils::{ blob_0, blob_1, empty_footer_payload, empty_footer_payload_bytes, file_properties, java_empty_uncompressed_input_file, java_uncompressed_metric_input_file, - java_zstd_compressed_metric_input_file, uncompressed_metric_file_metadata, - zstd_compressed_metric_file_metadata, + java_zstd_compressed_metric_input_file, read_file_metadata, + uncompressed_metric_file_metadata, zstd_compressed_metric_file_metadata, }; use crate::puffin::writer::PuffinWriter; use crate::{ErrorKind, Result}; @@ -214,14 +214,6 @@ mod tests { Ok(output_file) } - async fn read_file_metadata(input_file: &InputFile) -> FileMetadata { - let file_read = input_file.reader().await.unwrap(); - let file_length = input_file.metadata().await.unwrap().size; - FileMetadata::read(file_read.as_ref(), file_length) - .await - .unwrap() - } - async fn read_all_blobs_from_puffin_file(input_file: InputFile) -> Vec { let puffin_reader = PuffinReader::new(input_file).await.unwrap(); let mut blobs = Vec::new(); @@ -242,7 +234,7 @@ mod tests { .to_input_file(); assert_eq!( - read_file_metadata(&input_file).await, + read_file_metadata(&input_file).await.unwrap(), empty_footer_payload() ); @@ -278,7 +270,7 @@ mod tests { .to_input_file(); assert_eq!( - read_file_metadata(&input_file).await, + read_file_metadata(&input_file).await.unwrap(), uncompressed_metric_file_metadata() ); @@ -298,7 +290,7 @@ mod tests { .to_input_file(); assert_eq!( - read_file_metadata(&input_file).await, + read_file_metadata(&input_file).await.unwrap(), zstd_compressed_metric_file_metadata() ); From 3af2cc8a123cdc34fed7e105e01f90ea43765845 Mon Sep 17 00:00:00 2001 From: Xander Date: Tue, 28 Jul 2026 23:18:39 +0100 Subject: [PATCH 3/3] fmt --- crates/iceberg/src/puffin/metadata.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/crates/iceberg/src/puffin/metadata.rs b/crates/iceberg/src/puffin/metadata.rs index 78d904f066..1ee954b873 100644 --- a/crates/iceberg/src/puffin/metadata.rs +++ b/crates/iceberg/src/puffin/metadata.rs @@ -408,6 +408,7 @@ mod tests { use bytes::Bytes; use tempfile::TempDir; + use crate::ErrorKind; use crate::io::{FileIO, InputFile}; use crate::puffin::metadata::{BlobMetadata, CompressionCodec, FileMetadata}; use crate::puffin::test_utils::{ @@ -417,7 +418,6 @@ mod tests { read_file_metadata_with_prefetch, uncompressed_metric_file_metadata, zstd_compressed_metric_file_metadata, }; - use crate::ErrorKind; const INVALID_MAGIC_VALUE: [u8; 4] = [80, 70, 65, 0];