diff --git a/crates/iceberg/src/transaction/append.rs b/crates/iceberg/src/transaction/append.rs index dabb46461e..fb03e2bbf0 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}; @@ -387,6 +387,81 @@ mod tests { .collect() } + #[tokio::test] + async fn test_fast_append_writes_encrypted_manifest() { + 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 = 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) + .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 c2fde69fca..7b8e02da36 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(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 5685786922..b4f1946a7c 100644 --- a/crates/iceberg/src/transaction/snapshot.rs +++ b/crates/iceberg/src/transaction/snapshot.rs @@ -246,16 +246,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 {