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 49a24b6b9 feat(encryption): [16/N] Enable encryption for purge table
(#2918)
49a24b6b9 is described below
commit 49a24b6b96eb96fe3f0c5758aa244dd8caaa09e7
Author: Xander <[email protected]>
AuthorDate: Thu Jul 30 09:39:05 2026 +0100
feat(encryption): [16/N] Enable encryption for purge table (#2918)
## 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?
We previously used `FileIO` to read manifest files rather than going
through the `load_manifest` method which handles decryption correctly.
<!--
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?
Tested added to show we can purge encrypted tables.
<!--
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 | 1 +
crates/iceberg/src/catalog/utils.rs | 24 +++++++------
crates/iceberg/src/test_utils.rs | 62 +++++++++++++++++++++++++++++++-
crates/iceberg/src/transaction/append.rs | 4 +--
crates/iceberg/src/transaction/mod.rs | 60 +------------------------------
5 files changed, 78 insertions(+), 73 deletions(-)
diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt
index 8a6227729..114cc6be7 100644
--- a/crates/iceberg/public-api.txt
+++ b/crates/iceberg/public-api.txt
@@ -3137,6 +3137,7 @@ pub fn iceberg::table::TableBuilder::readonly(self,
readonly: bool) -> Self
pub fn iceberg::table::TableBuilder::runtime(self, runtime: iceberg::Runtime)
-> Self
pub mod iceberg::test_utils
pub fn iceberg::test_utils::check_record_batches(record_batches:
alloc::vec::Vec<arrow_array::record_batch::RecordBatch>, expected_schema:
expect_test::Expect, expected_data: expect_test::Expect, ignore_check_columns:
&[&str], sort_column: core::option::Option<&str>)
+pub async fn iceberg::test_utils::make_encrypted_table() ->
iceberg::table::Table
pub fn iceberg::test_utils::test_runtime() -> iceberg::Runtime
pub mod iceberg::transaction
pub struct iceberg::transaction::ActionCommit
diff --git a/crates/iceberg/src/catalog/utils.rs
b/crates/iceberg/src/catalog/utils.rs
index 853fc8d09..8e743e7d7 100644
--- a/crates/iceberg/src/catalog/utils.rs
+++ b/crates/iceberg/src/catalog/utils.rs
@@ -17,12 +17,13 @@
//! Utility functions for catalog operations.
-use std::collections::HashSet;
+use std::collections::{HashMap, HashSet};
use futures::{TryStreamExt, stream};
use crate::Result;
use crate::io::FileIO;
+use crate::spec::ManifestFile;
use crate::table::Table;
const DELETE_CONCURRENCY: usize = 10;
@@ -38,7 +39,7 @@ const DELETE_CONCURRENCY: usize = 10;
/// may share the same data files.
pub async fn drop_table_data(table_info: &Table) -> Result<()> {
let mut manifest_lists_to_delete: HashSet<String> = HashSet::new();
- let mut manifests_to_delete: HashSet<String> = HashSet::new();
+ let mut manifests_to_delete: HashMap<String, ManifestFile> =
HashMap::new();
let metadata = table_info.metadata_ref();
let io = table_info.file_io();
@@ -55,7 +56,7 @@ pub async fn drop_table_data(table_info: &Table) ->
Result<()> {
manifest_lists_to_delete.insert(manifest_list_location);
}
for manifest_file in manifest_list.entries() {
- manifests_to_delete.insert(manifest_file.manifest_path.clone());
+ manifests_to_delete.insert(manifest_file.manifest_path.clone(),
manifest_file.clone());
}
}
@@ -65,7 +66,8 @@ pub async fn drop_table_data(table_info: &Table) ->
Result<()> {
}
// Delete manifest files
- io.delete_stream(stream::iter(manifests_to_delete)).await?;
+ let manifest_paths: Vec<String> =
manifests_to_delete.into_keys().collect();
+ io.delete_stream(stream::iter(manifest_paths)).await?;
// Delete manifest lists
io.delete_stream(stream::iter(manifest_lists_to_delete))
@@ -103,13 +105,13 @@ pub async fn drop_table_data(table_info: &Table) ->
Result<()> {
}
/// Reads manifests concurrently and deletes the data files referenced within.
-async fn delete_data_files(io: &FileIO, manifest_paths: &HashSet<String>) ->
Result<()> {
- stream::iter(manifest_paths.iter().map(Ok))
- .try_for_each_concurrent(DELETE_CONCURRENCY, |manifest_path| async
move {
- let input = io.new_input(manifest_path)?;
- let manifest_content = input.read().await?;
- let manifest =
crate::spec::Manifest::parse_avro(&manifest_content)?;
-
+async fn delete_data_files(
+ io: &FileIO,
+ manifest_files: &HashMap<String, ManifestFile>,
+) -> Result<()> {
+ stream::iter(manifest_files.values().map(Ok))
+ .try_for_each_concurrent(DELETE_CONCURRENCY, |manifest_file| async
move {
+ let manifest = manifest_file.load_manifest(io).await?;
let data_file_paths = manifest
.entries()
.iter()
diff --git a/crates/iceberg/src/test_utils.rs b/crates/iceberg/src/test_utils.rs
index d47c39950..abfd2af9f 100644
--- a/crates/iceberg/src/test_utils.rs
+++ b/crates/iceberg/src/test_utils.rs
@@ -19,13 +19,19 @@
//! This module is pub just for internal testing.
//! It is subject to change and is not intended to be used by external users.
-use std::sync::OnceLock;
+use std::sync::{Arc, OnceLock};
use arrow_array::RecordBatch;
use expect_test::Expect;
use itertools::Itertools;
+use crate::TableIdent;
+use crate::encryption::SensitiveBytes;
+use crate::encryption::kms::{KeyManagementClient, MemoryKeyManagementClient};
+use crate::io::FileIO;
use crate::runtime::Runtime;
+use crate::spec::TableMetadata;
+use crate::table::Table;
/// Returns a process-wide [`Runtime`] suitable for tests that need to
construct
/// a [`Table`](crate::table::Table) outside a tokio context.
@@ -98,3 +104,57 @@ pub fn check_record_batches(
.format(",\n")
));
}
+
+/// 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 async fn make_encrypted_table() -> Table {
+ let metadata_json = std::fs::read_to_string(format!(
+ "{}/testdata/table_metadata/TableMetadataV3ValidEncryption.json",
+ env!("CARGO_MANIFEST_DIR"),
+ ))
+ .unwrap();
+ let metadata: TableMetadata =
serde_json::from_str(&metadata_json).unwrap();
+
+ let kms: Arc<dyn KeyManagementClient> = {
+ 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()
+}
diff --git a/crates/iceberg/src/transaction/append.rs
b/crates/iceberg/src/transaction/append.rs
index 799abdeff..a19e297ac 100644
--- a/crates/iceberg/src/transaction/append.rs
+++ b/crates/iceberg/src/transaction/append.rs
@@ -170,8 +170,8 @@ mod tests {
Struct, TableMetadata,
};
use crate::table::Table;
- use crate::test_utils::test_runtime;
- use crate::transaction::tests::{make_encrypted_table,
make_v2_minimal_table};
+ use crate::test_utils::{make_encrypted_table, test_runtime};
+ use crate::transaction::tests::make_v2_minimal_table;
use crate::transaction::{Transaction, TransactionAction};
use crate::{TableIdent, TableRequirement, TableUpdate};
diff --git a/crates/iceberg/src/transaction/mod.rs
b/crates/iceberg/src/transaction/mod.rs
index 7b8e02da3..3e0a4e939 100644
--- a/crates/iceberg/src/transaction/mod.rs
+++ b/crates/iceberg/src/transaction/mod.rs
@@ -259,15 +259,13 @@ mod tests {
use std::sync::atomic::{AtomicU32, Ordering};
use crate::catalog::MockCatalog;
- use crate::encryption::SensitiveBytes;
- use crate::encryption::kms::{KeyManagementClient,
MemoryKeyManagementClient};
use crate::io::FileIO;
use crate::memory::tests::new_memory_catalog;
use crate::spec::{
DataContentType, DataFileBuilder, DataFileFormat, Literal, Struct,
TableMetadata,
};
use crate::table::Table;
- use crate::test_utils::test_runtime;
+ use crate::test_utils::{make_encrypted_table, test_runtime};
use crate::transaction::{ApplyTransactionAction, Transaction};
use crate::{Catalog, Error, ErrorKind, TableCreation, TableIdent};
@@ -331,62 +329,6 @@ 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<dyn KeyManagementClient> = {
- 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()])