This is an automated email from the ASF dual-hosted git repository.
blackmwk pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git
The following commit(s) were added to refs/heads/main by this push:
new c1d877da9 feat(encryption) [14/N] Read / Write Encrypted puffin files
(#2822)
c1d877da9 is described below
commit c1d877da953e57cd4916b17dde63ad36ff128028
Author: Xander <[email protected]>
AuthorDate: Thu Jul 30 10:42:55 2026 +0100
feat(encryption) [14/N] Read / Write Encrypted puffin files (#2822)
## Which issue does this PR close?
<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes #123` indicates that this PR will close issue #123.
-->
- Working towards https://github.com/apache/iceberg-rust/issues/2034
## What changes are included in this PR?
Puffin files aren't yet used in the repo but we should support
encryption nevertheless, similar to
https://github.com/apache/iceberg-rust/pull/2568, we add a
`new_from_encrypted` constructor for the reader and writer.
<!--
Provide a summary of the modifications in this PR. List the main changes
such as new features, bug fixes, refactoring, or any other updates.
-->
## Are these changes tested?
Yes - roundtrip test
<!--
Specify what test covers (unit test, integration test, etc.).
If tests are not included in your PR, please explain why (for example,
are they covered by existing tests)?
-->
---
crates/iceberg/public-api.txt | 4 +-
crates/iceberg/src/puffin/metadata.rs | 98 ++++++++++++++-----------------
crates/iceberg/src/puffin/reader.rs | 41 +++++++++----
crates/iceberg/src/puffin/test_utils.rs | 18 ++++++
crates/iceberg/src/puffin/writer.rs | 100 ++++++++++++++++++++++++++++----
5 files changed, 186 insertions(+), 75 deletions(-)
diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt
index 114cc6be7..12f844b36 100644
--- a/crates/iceberg/public-api.txt
+++ b/crates/iceberg/public-api.txt
@@ -1250,12 +1250,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<iceberg::puffin::Blob>
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<Self>
+pub async fn
iceberg::puffin::PuffinReader::new_from_encrypted(encrypted_input:
iceberg::encryption::EncryptedInputFile) -> iceberg::Result<Self>
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<alloc::string::String,
alloc::string::String>, compress_footer: bool) -> iceberg::Result<Self>
+pub async fn
iceberg::puffin::PuffinWriter::new_from_encrypted(encrypted_output:
&iceberg::encryption::EncryptedOutputFile, properties:
std::collections::hash::map::HashMap<alloc::string::String,
alloc::string::String>, compress_footer: bool) -> iceberg::Result<Self>
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 35984a5ef..1ee954b87 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<FileMetadata> {
- 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<FileMetadata> {
+ 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<FileMetadata> {
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]
@@ -423,7 +414,8 @@ 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,
};
@@ -473,7 +465,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 +488,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 +511,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 +536,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 +561,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 +584,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 +607,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 +633,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 +648,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 +671,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 +694,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 +718,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 +747,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 +773,7 @@ mod tests {
.await;
assert_eq!(
- FileMetadata::read(&input_file)
+ read_file_metadata(&input_file)
.await
.unwrap_err()
.to_string(),
@@ -804,7 +796,7 @@ mod tests {
.await;
assert_eq!(
- FileMetadata::read(&input_file)
+ read_file_metadata(&input_file)
.await
.unwrap_err()
.to_string(),
@@ -844,7 +836,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 +890,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 +937,7 @@ mod tests {
.await;
assert_eq!(
- FileMetadata::read(&input_file)
+ read_file_metadata(&input_file)
.await
.unwrap_err()
.to_string(),
@@ -962,7 +954,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 +966,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 +974,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 +982,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 +991,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 +1001,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 +1011,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 +1038,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 +1073,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 0aced4186..01601f48e 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<dyn FileRead>,
+ file_length: u64,
file_metadata: OnceCell<FileMetadata>,
}
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<Self> {
+ 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<Self> {
+ 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<dyn FileRead>, 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<Blob> {
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/test_utils.rs
b/crates/iceberg/src/puffin/test_utils.rs
index e0844e200..26a6834a6 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<FileMetadata> {
+ 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<FileMetadata> {
+ 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 4af4970b0..0437bd5bf 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<String, String>,
compress_footer: bool,
) -> Result<Self> {
+ 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<String, String>,
+ compress_footer: bool,
+ ) -> Result<Self> {
+ Ok(Self::from_writer(
+ encrypted_output.writer().await?,
+ properties,
+ compress_footer,
+ ))
+ }
+
+ fn from_writer(
+ writer: Box<dyn FileWrite>,
+ properties: HashMap<String, String>,
+ compress_footer: bool,
+ ) -> Self {
let mut flags = HashSet::<Flag>::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
@@ -158,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};
@@ -185,7 +215,7 @@ mod tests {
}
async fn read_all_blobs_from_puffin_file(input_file: InputFile) ->
Vec<Blob> {
- 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 +234,7 @@ mod tests {
.to_input_file();
assert_eq!(
- FileMetadata::read(&input_file).await.unwrap(),
+ read_file_metadata(&input_file).await.unwrap(),
empty_footer_payload()
);
@@ -240,7 +270,7 @@ mod tests {
.to_input_file();
assert_eq!(
- FileMetadata::read(&input_file).await.unwrap(),
+ read_file_metadata(&input_file).await.unwrap(),
uncompressed_metric_file_metadata()
);
@@ -260,7 +290,7 @@ mod tests {
.to_input_file();
assert_eq!(
- FileMetadata::read(&input_file).await.unwrap(),
+ read_file_metadata(&input_file).await.unwrap(),
zstd_compressed_metric_file_metadata()
);
@@ -354,4 +384,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);
+ }
}