From 206b8fda51c107a2d16d6013b5e92ed5d1550e3b Mon Sep 17 00:00:00 2001 From: Xander Date: Mon, 13 Jul 2026 12:16:42 +0100 Subject: [PATCH 1/4] Write encrypted manifest files --- crates/iceberg/src/transaction/snapshot.rs | 205 ++++++++++++++++++++- 1 file changed, 195 insertions(+), 10 deletions(-) diff --git a/crates/iceberg/src/transaction/snapshot.rs b/crates/iceberg/src/transaction/snapshot.rs index d200c5ba9c..ecbc3f1fae 100644 --- a/crates/iceberg/src/transaction/snapshot.rs +++ b/crates/iceberg/src/transaction/snapshot.rs @@ -249,16 +249,25 @@ impl<'a> SnapshotProducer<'a> { DataFileFormat::Avro ); let output_file = self.table.file_io().new_output(new_manifest_path)?; - let builder = ManifestWriterBuilder::new( - output_file, - Some(self.snapshot_id), - self.table.metadata().current_schema().clone(), - self.table - .metadata() - .default_partition_spec() - .as_ref() - .clone(), - ); + let partition_spec = self + .table + .metadata() + .default_partition_spec() + .as_ref() + .clone(); + let schema = self.table.metadata().current_schema().clone(); + + let builder = if let Some(em) = self.table.encryption_manager() { + ManifestWriterBuilder::new_from_encrypted( + em.encrypt(output_file), + Some(self.snapshot_id), + schema, + partition_spec, + )? + } else { + ManifestWriterBuilder::new(output_file, Some(self.snapshot_id), schema, partition_spec) + }; + match self.table.metadata().format_version() { FormatVersion::V1 => Ok(builder.build_v1()), FormatVersion::V2 => match content { @@ -530,3 +539,179 @@ impl<'a> SnapshotProducer<'a> { Ok(ActionCommit::new(updates, requirements)) } } + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + use std::fs::File; + use std::io::BufReader; + use std::sync::Arc; + + use uuid::Uuid; + + use super::SnapshotProducer; + use crate::TableIdent; + use crate::encryption::kms::{KeyManagementClient, MemoryKeyManagementClient}; + use crate::encryption::{SensitiveBytes, StandardKeyMetadata}; + use crate::io::FileIO; + use crate::spec::{ + DataContentType, DataFileBuilder, DataFileFormat, Manifest, ManifestContentType, + ManifestStatus, Struct, TableMetadata, + }; + use crate::table::Table; + use crate::test_utils::test_runtime; + + /// Build a table backed by the V3 encryption fixture and an in-memory KMS, + /// so it has an [`EncryptionManager`](crate::encryption::EncryptionManager). + fn make_encrypted_table() -> Table { + let file = File::open(format!( + "{}/testdata/table_metadata/{}", + env!("CARGO_MANIFEST_DIR"), + "TableMetadataV3ValidEncryption.json" + )) + .unwrap(); + let reader = BufReader::new(file); + let metadata = serde_json::from_reader::<_, TableMetadata>(reader).unwrap(); + + let kms: Arc = { + let k = MemoryKeyManagementClient::new(); + k.add_master_key_bytes( + "master-1", + SensitiveBytes::new([ + 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, + 0x0d, 0x0e, 0x0f, + ]), + ) + .unwrap(); + Arc::new(k) + }; + + Table::builder() + .metadata(metadata) + .metadata_location("memory:///table/metadata/v1.json") + .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap()) + .file_io(FileIO::new_with_memory()) + .kms_client(kms) + .runtime(test_runtime()) + .build() + .unwrap() + } + + fn data_file() -> crate::spec::DataFile { + DataFileBuilder::default() + .content(DataContentType::Data) + .file_path("memory:///table/data/00000.parquet".to_string()) + .file_format(DataFileFormat::Parquet) + .partition(Struct::empty()) + .record_count(100) + .file_size_in_bytes(4096) + .partition_spec_id(0) + .build() + .unwrap() + } + + #[tokio::test] + async fn write_encrypted_manifest_roundtrips() { + let table = make_encrypted_table(); + assert!( + table.encryption_manager().is_some(), + "fixture table should have an EncryptionManager" + ); + + let mut producer = + SnapshotProducer::new(&table, Uuid::new_v4(), HashMap::new(), vec![data_file()]); + + let mut writer = producer + .new_manifest_writer(ManifestContentType::Data) + .unwrap(); + writer + .add_entry( + crate::spec::ManifestEntry::builder() + .status(ManifestStatus::Added) + .data_file(data_file()) + .build(), + ) + .unwrap(); + let manifest_file = writer.write_manifest_file().await.unwrap(); + + // The manifest list entry must carry decodable key metadata. + let key_metadata_bytes = manifest_file + .key_metadata + .as_ref() + .expect("encrypted manifest must record key metadata"); + assert!(!key_metadata_bytes.is_empty()); + StandardKeyMetadata::decode(key_metadata_bytes) + .expect("recorded key metadata must decode as StandardKeyMetadata"); + + // The bytes on disk must actually be encrypted: a plaintext Avro parse + // must fail. (Guards against silently taking the unencrypted branch.) + let raw = table + .file_io() + .new_input(&manifest_file.manifest_path) + .unwrap() + .read() + .await + .unwrap(); + assert!( + Manifest::parse_avro(&raw).is_err(), + "manifest bytes should not be readable as plaintext Avro" + ); + + // load_manifest self-decrypts using the recorded key metadata and must + // recover the entry we wrote. + let manifest = manifest_file.load_manifest(table.file_io()).await.unwrap(); + assert_eq!(manifest.entries().len(), 1); + assert_eq!( + manifest.entries()[0].data_file().file_path(), + "memory:///table/data/00000.parquet" + ); + } + + #[tokio::test] + async fn write_plaintext_manifest_when_no_encryption() { + let table = crate::transaction::tests::make_v2_table(); + assert!(table.encryption_manager().is_none()); + + let data_file = DataFileBuilder::default() + .content(DataContentType::Data) + .file_path("memory:///table/data/00000.parquet".to_string()) + .file_format(DataFileFormat::Parquet) + .partition(Struct::from_iter([Some(crate::spec::Literal::long(0))])) + .record_count(100) + .file_size_in_bytes(4096) + .partition_spec_id(0) + .build() + .unwrap(); + + let mut producer = SnapshotProducer::new(&table, Uuid::new_v4(), HashMap::new(), vec![ + data_file.clone(), + ]); + + let mut writer = producer + .new_manifest_writer(ManifestContentType::Data) + .unwrap(); + writer + .add_entry( + crate::spec::ManifestEntry::builder() + .status(ManifestStatus::Added) + .data_file(data_file) + .build(), + ) + .unwrap(); + let manifest_file = writer.write_manifest_file().await.unwrap(); + + assert!( + manifest_file.key_metadata.is_none(), + "plaintext manifest must not record key metadata" + ); + + let raw = table + .file_io() + .new_input(&manifest_file.manifest_path) + .unwrap() + .read() + .await + .unwrap(); + Manifest::parse_avro(&raw).expect("plaintext manifest bytes must parse as Avro"); + } +} From 9b1edf97617b7ff1333722e032c8c39d80efd778 Mon Sep 17 00:00:00 2001 From: Xander Date: Wed, 22 Jul 2026 10:30:34 +0100 Subject: [PATCH 2/4] comments --- crates/iceberg/src/transaction/append.rs | 77 +++++++++ crates/iceberg/src/transaction/mod.rs | 88 +++++++---- crates/iceberg/src/transaction/snapshot.rs | 176 --------------------- 3 files changed, 134 insertions(+), 207 deletions(-) diff --git a/crates/iceberg/src/transaction/append.rs b/crates/iceberg/src/transaction/append.rs index 9d94e18ede..d602f454f7 100644 --- a/crates/iceberg/src/transaction/append.rs +++ b/crates/iceberg/src/transaction/append.rs @@ -352,6 +352,83 @@ mod tests { ); } + #[tokio::test] + async fn test_fast_append_writes_encrypted_manifest() { + use crate::encryption::StandardKeyMetadata; + use crate::spec::Manifest; + use crate::transaction::tests::make_encrypted_table; + + let table = make_encrypted_table().await; + assert!( + table.encryption_manager().is_some(), + "fixture table should have an EncryptionManager" + ); + + let new_file = DataFileBuilder::default() + .content(DataContentType::Data) + .file_path("memory:///table/data/00000.parquet".to_string()) + .file_format(DataFileFormat::Parquet) + .partition(Struct::empty()) + .record_count(100) + .file_size_in_bytes(4096) + .partition_spec_id(table.metadata().default_partition_spec_id()) + .build() + .unwrap(); + + let tx = Transaction::new(&table); + let action = tx.fast_append().add_data_files(vec![new_file]); + let mut action_commit = Arc::new(action).commit(&table).await.unwrap(); + let updates = action_commit.take_updates(); + + let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] { + SnapshotRef::new(snapshot.clone()) + } else { + unreachable!("first update of a fast append should be AddSnapshot") + }; + + let manifest_list = table + .manifest_list_reader(&new_snapshot) + .load() + .await + .unwrap(); + let manifest_file = manifest_list + .entries() + .iter() + .find(|m| m.added_files_count.unwrap_or(0) > 0) + .expect("new snapshot should carry the appended data manifest"); + + // The manifest list entry must carry decodable key metadata. + let key_metadata_bytes = manifest_file + .key_metadata + .as_ref() + .expect("encrypted manifest must record key metadata"); + StandardKeyMetadata::decode(key_metadata_bytes) + .expect("recorded key metadata must decode as StandardKeyMetadata"); + + // The bytes on disk must actually be encrypted: a plaintext Avro parse + // must fail. (Guards against silently taking the unencrypted branch.) + let raw = table + .file_io() + .new_input(&manifest_file.manifest_path) + .unwrap() + .read() + .await + .unwrap(); + assert!( + Manifest::parse_avro(&raw).is_err(), + "manifest bytes should not be readable as plaintext Avro" + ); + + // load_manifest self-decrypts using the recorded key metadata and must + // recover the entry we appended. + let manifest = manifest_file.load_manifest(table.file_io()).await.unwrap(); + assert_eq!(manifest.entries().len(), 1); + assert_eq!( + manifest.entries()[0].data_file().file_path(), + "memory:///table/data/00000.parquet" + ); + } + #[tokio::test] async fn test_empty_data_append_action() { let table = make_v2_minimal_table(); diff --git a/crates/iceberg/src/transaction/mod.rs b/crates/iceberg/src/transaction/mod.rs index d78f41cd42..152e6dc3c7 100644 --- a/crates/iceberg/src/transaction/mod.rs +++ b/crates/iceberg/src/transaction/mod.rs @@ -331,6 +331,62 @@ mod tests { .unwrap() } + /// Build a table backed by the V3 encryption fixture and an in-memory KMS, + /// so it has an [`EncryptionManager`](crate::encryption::EncryptionManager). + /// + /// The fixture's snapshot references an encrypted manifest list; its bytes + /// (the `manifest-list-v3-encrypted.avro` testdata, an encrypted empty list) + /// are seeded into the in-memory `FileIO` at that path so callers can read + /// the current snapshot's manifest list. + pub(crate) async fn make_encrypted_table() -> Table { + let file = File::open(format!( + "{}/testdata/table_metadata/{}", + env!("CARGO_MANIFEST_DIR"), + "TableMetadataV3ValidEncryption.json" + )) + .unwrap(); + let reader = BufReader::new(file); + let metadata = serde_json::from_reader::<_, TableMetadata>(reader).unwrap(); + + let kms: Arc = { + let k = MemoryKeyManagementClient::new(); + k.add_master_key_bytes( + "master-1", + SensitiveBytes::new([ + 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, + 0x0d, 0x0e, 0x0f, + ]), + ) + .unwrap(); + Arc::new(k) + }; + + let file_io = FileIO::new_with_memory(); + + // Seed the encrypted (empty) manifest list at the path the snapshot references. + let manifest_list_bytes = std::fs::read(format!( + "{}/testdata/manifests_lists/manifest-list-v3-encrypted.avro", + env!("CARGO_MANIFEST_DIR"), + )) + .unwrap(); + file_io + .new_output(metadata.current_snapshot().unwrap().manifest_list()) + .unwrap() + .write(manifest_list_bytes.into()) + .await + .unwrap(); + + Table::builder() + .metadata(metadata) + .metadata_location("memory:///table/metadata/v1.json") + .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap()) + .file_io(file_io) + .kms_client(kms) + .runtime(test_runtime()) + .build() + .unwrap() + } + pub(crate) async fn make_v3_minimal_table_in_catalog(catalog: &impl Catalog) -> Table { let table_ident = TableIdent::from_strs([format!("ns1-{}", uuid::Uuid::new_v4()), "test1".to_string()]) @@ -579,37 +635,7 @@ mod tests { #[tokio::test] async fn test_commit_rejects_encrypted_table() { - let file = File::open(format!( - "{}/testdata/table_metadata/{}", - env!("CARGO_MANIFEST_DIR"), - "TableMetadataV3ValidEncryption.json" - )) - .unwrap(); - let reader = BufReader::new(file); - let resp = serde_json::from_reader::<_, TableMetadata>(reader).unwrap(); - - let kms: Arc = { - let k = MemoryKeyManagementClient::new(); - k.add_master_key_bytes( - "master-1", - SensitiveBytes::new([ - 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, - 0x0d, 0x0e, 0x0f, - ]), - ) - .unwrap(); - Arc::new(k) - }; - - let table = Table::builder() - .metadata(resp) - .metadata_location("s3://bucket/test/location/metadata/v1.json") - .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap()) - .file_io(FileIO::new_with_memory()) - .kms_client(kms) - .runtime(crate::test_utils::test_runtime()) - .build() - .unwrap(); + let table = make_encrypted_table().await; let tx = Transaction::new(&table); let tx = tx diff --git a/crates/iceberg/src/transaction/snapshot.rs b/crates/iceberg/src/transaction/snapshot.rs index ecbc3f1fae..7cd0cdcb68 100644 --- a/crates/iceberg/src/transaction/snapshot.rs +++ b/crates/iceberg/src/transaction/snapshot.rs @@ -539,179 +539,3 @@ impl<'a> SnapshotProducer<'a> { Ok(ActionCommit::new(updates, requirements)) } } - -#[cfg(test)] -mod tests { - use std::collections::HashMap; - use std::fs::File; - use std::io::BufReader; - use std::sync::Arc; - - use uuid::Uuid; - - use super::SnapshotProducer; - use crate::TableIdent; - use crate::encryption::kms::{KeyManagementClient, MemoryKeyManagementClient}; - use crate::encryption::{SensitiveBytes, StandardKeyMetadata}; - use crate::io::FileIO; - use crate::spec::{ - DataContentType, DataFileBuilder, DataFileFormat, Manifest, ManifestContentType, - ManifestStatus, Struct, TableMetadata, - }; - use crate::table::Table; - use crate::test_utils::test_runtime; - - /// Build a table backed by the V3 encryption fixture and an in-memory KMS, - /// so it has an [`EncryptionManager`](crate::encryption::EncryptionManager). - fn make_encrypted_table() -> Table { - let file = File::open(format!( - "{}/testdata/table_metadata/{}", - env!("CARGO_MANIFEST_DIR"), - "TableMetadataV3ValidEncryption.json" - )) - .unwrap(); - let reader = BufReader::new(file); - let metadata = serde_json::from_reader::<_, TableMetadata>(reader).unwrap(); - - let kms: Arc = { - let k = MemoryKeyManagementClient::new(); - k.add_master_key_bytes( - "master-1", - SensitiveBytes::new([ - 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, - 0x0d, 0x0e, 0x0f, - ]), - ) - .unwrap(); - Arc::new(k) - }; - - Table::builder() - .metadata(metadata) - .metadata_location("memory:///table/metadata/v1.json") - .identifier(TableIdent::from_strs(["ns1", "test1"]).unwrap()) - .file_io(FileIO::new_with_memory()) - .kms_client(kms) - .runtime(test_runtime()) - .build() - .unwrap() - } - - fn data_file() -> crate::spec::DataFile { - DataFileBuilder::default() - .content(DataContentType::Data) - .file_path("memory:///table/data/00000.parquet".to_string()) - .file_format(DataFileFormat::Parquet) - .partition(Struct::empty()) - .record_count(100) - .file_size_in_bytes(4096) - .partition_spec_id(0) - .build() - .unwrap() - } - - #[tokio::test] - async fn write_encrypted_manifest_roundtrips() { - let table = make_encrypted_table(); - assert!( - table.encryption_manager().is_some(), - "fixture table should have an EncryptionManager" - ); - - let mut producer = - SnapshotProducer::new(&table, Uuid::new_v4(), HashMap::new(), vec![data_file()]); - - let mut writer = producer - .new_manifest_writer(ManifestContentType::Data) - .unwrap(); - writer - .add_entry( - crate::spec::ManifestEntry::builder() - .status(ManifestStatus::Added) - .data_file(data_file()) - .build(), - ) - .unwrap(); - let manifest_file = writer.write_manifest_file().await.unwrap(); - - // The manifest list entry must carry decodable key metadata. - let key_metadata_bytes = manifest_file - .key_metadata - .as_ref() - .expect("encrypted manifest must record key metadata"); - assert!(!key_metadata_bytes.is_empty()); - StandardKeyMetadata::decode(key_metadata_bytes) - .expect("recorded key metadata must decode as StandardKeyMetadata"); - - // The bytes on disk must actually be encrypted: a plaintext Avro parse - // must fail. (Guards against silently taking the unencrypted branch.) - let raw = table - .file_io() - .new_input(&manifest_file.manifest_path) - .unwrap() - .read() - .await - .unwrap(); - assert!( - Manifest::parse_avro(&raw).is_err(), - "manifest bytes should not be readable as plaintext Avro" - ); - - // load_manifest self-decrypts using the recorded key metadata and must - // recover the entry we wrote. - let manifest = manifest_file.load_manifest(table.file_io()).await.unwrap(); - assert_eq!(manifest.entries().len(), 1); - assert_eq!( - manifest.entries()[0].data_file().file_path(), - "memory:///table/data/00000.parquet" - ); - } - - #[tokio::test] - async fn write_plaintext_manifest_when_no_encryption() { - let table = crate::transaction::tests::make_v2_table(); - assert!(table.encryption_manager().is_none()); - - let data_file = DataFileBuilder::default() - .content(DataContentType::Data) - .file_path("memory:///table/data/00000.parquet".to_string()) - .file_format(DataFileFormat::Parquet) - .partition(Struct::from_iter([Some(crate::spec::Literal::long(0))])) - .record_count(100) - .file_size_in_bytes(4096) - .partition_spec_id(0) - .build() - .unwrap(); - - let mut producer = SnapshotProducer::new(&table, Uuid::new_v4(), HashMap::new(), vec![ - data_file.clone(), - ]); - - let mut writer = producer - .new_manifest_writer(ManifestContentType::Data) - .unwrap(); - writer - .add_entry( - crate::spec::ManifestEntry::builder() - .status(ManifestStatus::Added) - .data_file(data_file) - .build(), - ) - .unwrap(); - let manifest_file = writer.write_manifest_file().await.unwrap(); - - assert!( - manifest_file.key_metadata.is_none(), - "plaintext manifest must not record key metadata" - ); - - let raw = table - .file_io() - .new_input(&manifest_file.manifest_path) - .unwrap() - .read() - .await - .unwrap(); - Manifest::parse_avro(&raw).expect("plaintext manifest bytes must parse as Avro"); - } -} From 10d04000e7f4dbdcad804e4d28882b35f979c2f8 Mon Sep 17 00:00:00 2001 From: Xander Date: Wed, 22 Jul 2026 10:34:38 +0100 Subject: [PATCH 3/4] format tetst imports --- crates/iceberg/src/transaction/append.rs | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/crates/iceberg/src/transaction/append.rs b/crates/iceberg/src/transaction/append.rs index 317f18da91..04e7545070 100644 --- a/crates/iceberg/src/transaction/append.rs +++ b/crates/iceberg/src/transaction/append.rs @@ -161,17 +161,17 @@ mod tests { use tempfile::TempDir; use uuid::Uuid; - use crate::encryption::SensitiveBytes; use crate::encryption::kms::MemoryKeyManagementClient; + use crate::encryption::{SensitiveBytes, StandardKeyMetadata}; use crate::io::FileIO; use crate::spec::{ - DataContentType, DataFile, DataFileBuilder, DataFileFormat, Literal, MAIN_BRANCH, + DataContentType, DataFile, DataFileBuilder, DataFileFormat, Literal, MAIN_BRANCH, Manifest, ManifestEntry, ManifestListWriter, ManifestStatus, ManifestWriterBuilder, SnapshotRef, Struct, TableMetadata, }; use crate::table::Table; use crate::test_utils::test_runtime; - use crate::transaction::tests::make_v2_minimal_table; + use crate::transaction::tests::{make_encrypted_table, make_v2_minimal_table}; use crate::transaction::{Transaction, TransactionAction}; use crate::{TableIdent, TableRequirement, TableUpdate}; @@ -389,10 +389,6 @@ mod tests { #[tokio::test] async fn test_fast_append_writes_encrypted_manifest() { - use crate::encryption::StandardKeyMetadata; - use crate::spec::Manifest; - use crate::transaction::tests::make_encrypted_table; - let table = make_encrypted_table().await; assert!( table.encryption_manager().is_some(), From 149009826234712b80c5cf472f92681e57d5aec4 Mon Sep 17 00:00:00 2001 From: Xander Date: Wed, 22 Jul 2026 13:18:00 +0100 Subject: [PATCH 4/4] fix test --- crates/iceberg/src/transaction/append.rs | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/crates/iceberg/src/transaction/append.rs b/crates/iceberg/src/transaction/append.rs index 04e7545070..fb03e2bbf0 100644 --- a/crates/iceberg/src/transaction/append.rs +++ b/crates/iceberg/src/transaction/append.rs @@ -411,11 +411,13 @@ mod tests { let mut action_commit = Arc::new(action).commit(&table).await.unwrap(); let updates = action_commit.take_updates(); - let new_snapshot: SnapshotRef = if let TableUpdate::AddSnapshot { snapshot } = &updates[0] { - SnapshotRef::new(snapshot.clone()) - } else { - unreachable!("first update of a fast append should be AddSnapshot") - }; + let new_snapshot: SnapshotRef = updates + .iter() + .find_map(|u| match u { + TableUpdate::AddSnapshot { snapshot } => Some(SnapshotRef::new(snapshot.clone())), + _ => None, + }) + .expect("a fast append should emit an AddSnapshot update"); let manifest_list = table .manifest_list_reader(&new_snapshot)