This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-3286-ffa864bf2fd3ba27ff221746f993d4bf7810440b in repository https://gitbox.apache.org/repos/asf/iceberg-rust.git
commit 3c91b3942ab40ec005dfbb9fdbccea848392cccb Author: NoahKusaba <[email protected]> AuthorDate: Fri Oct 2 18:56:13 2026 +0000 feat(catalog-loader): add with_runtime to BoxedCatalogBuilder (#3286) * init * address review: doc sibling methods, tidy test, note break in changelog --- CHANGELOG.md | 1 + crates/catalog/loader/public-api.txt | 2 + crates/catalog/loader/src/lib.rs | 94 ++++++++++++++++++++++++++++++++++-- 3 files changed, 93 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 83eeeb31e..f357a03e5 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -32,6 +32,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/). * `EncryptedOutputFile::key_metadata()` is replaced by `key_metadata_with_saved_file_metadata(&FileMetadata)`. Pass the metadata returned by `write()` or the writer's `close()` to include the stored length before encoding key metadata. * `EncryptedInputFile::metadata()` is now synchronous and derives the plaintext size from key metadata without a storage stat. Remove `.await` from calls to this method. * AGS1 readers now require `StandardKeyMetadata::file_length` and reject missing or invalid lengths without falling back to a storage stat. AGS1-encrypted manifests, manifest lists, and Puffin files written by earlier development builds without this field must be rewritten using a build that can still read them before upgrading. This matches the Java client's read contract. +* `iceberg_catalog_loader::BoxedCatalogBuilder` has a new required method, `with_runtime`. Types implementing `CatalogBuilder` get it through the blanket impl; types implementing `BoxedCatalogBuilder` directly (for example, a wrapper around another `Box<dyn BoxedCatalogBuilder>`) must implement it, typically by forwarding the runtime to the builder they wrap. ## [v0.10.1] - 2026-07-28 diff --git a/crates/catalog/loader/public-api.txt b/crates/catalog/loader/public-api.txt index 1efaf831d..4e3c28861 100644 --- a/crates/catalog/loader/public-api.txt +++ b/crates/catalog/loader/public-api.txt @@ -7,10 +7,12 @@ pub fn iceberg_catalog_loader::CatalogLoader<'a>::from(s: &'a str) -> Self pub trait iceberg_catalog_loader::BoxedCatalogBuilder: core::marker::Send pub fn iceberg_catalog_loader::BoxedCatalogBuilder::load<'async_trait>(self: alloc::boxed::Box<Self>, name: alloc::string::String, props: std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg::catalog::Catalog>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait pub fn iceberg_catalog_loader::BoxedCatalogBuilder::with_kms_client_factory(self: alloc::boxed::Box<Self>, kms_client_factory: alloc::sync::Arc<dyn iceberg::encryption::kms::factory::KmsClientFactory>) -> alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder> +pub fn iceberg_catalog_loader::BoxedCatalogBuilder::with_runtime(self: alloc::boxed::Box<Self>, runtime: iceberg::runtime::Runtime) -> alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder> pub fn iceberg_catalog_loader::BoxedCatalogBuilder::with_storage_factory(self: alloc::boxed::Box<Self>, storage_factory: alloc::sync::Arc<dyn iceberg::io::storage::StorageFactory>) -> alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder> impl<T: iceberg::catalog::CatalogBuilder + 'static> iceberg_catalog_loader::BoxedCatalogBuilder for T pub fn T::load<'async_trait>(self: alloc::boxed::Box<Self>, name: alloc::string::String, props: std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg::catalog::Catalog>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait pub fn T::with_kms_client_factory(self: alloc::boxed::Box<Self>, kms_client_factory: alloc::sync::Arc<dyn iceberg::encryption::kms::factory::KmsClientFactory>) -> alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder> +pub fn T::with_runtime(self: alloc::boxed::Box<Self>, runtime: iceberg::runtime::Runtime) -> alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder> pub fn T::with_storage_factory(self: alloc::boxed::Box<Self>, storage_factory: alloc::sync::Arc<dyn iceberg::io::storage::StorageFactory>) -> alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder> pub fn iceberg_catalog_loader::load(type: &str) -> iceberg::error::Result<alloc::boxed::Box<dyn iceberg_catalog_loader::BoxedCatalogBuilder>> pub fn iceberg_catalog_loader::supported_types() -> alloc::vec::Vec<&'static str> diff --git a/crates/catalog/loader/src/lib.rs b/crates/catalog/loader/src/lib.rs index d45decc40..3993de676 100644 --- a/crates/catalog/loader/src/lib.rs +++ b/crates/catalog/loader/src/lib.rs @@ -21,7 +21,7 @@ use std::sync::Arc; use async_trait::async_trait; use iceberg::encryption::kms::KmsClientFactory; use iceberg::io::StorageFactory; -use iceberg::{Catalog, CatalogBuilder, Error, ErrorKind, Result}; +use iceberg::{Catalog, CatalogBuilder, Error, ErrorKind, Result, Runtime}; use iceberg_catalog_glue::GlueCatalogBuilder; use iceberg_catalog_hms::HmsCatalogBuilder; use iceberg_catalog_rest::RestCatalogBuilder; @@ -47,16 +47,24 @@ pub fn supported_types() -> Vec<&'static str> { #[async_trait] pub trait BoxedCatalogBuilder: Send { + /// Sets the storage factory used to build the catalog's `FileIO`; see + /// [`CatalogBuilder::with_storage_factory`]. fn with_storage_factory( self: Box<Self>, storage_factory: Arc<dyn StorageFactory>, ) -> Box<dyn BoxedCatalogBuilder>; + /// Sets the KMS client factory used to enable table encryption; see + /// [`CatalogBuilder::with_kms_client_factory`]. fn with_kms_client_factory( self: Box<Self>, kms_client_factory: Arc<dyn KmsClientFactory>, ) -> Box<dyn BoxedCatalogBuilder>; + /// Sets the runtime the catalog, and the tables it creates, spawn their + /// tasks on; see [`CatalogBuilder::with_runtime`]. + fn with_runtime(self: Box<Self>, runtime: Runtime) -> Box<dyn BoxedCatalogBuilder>; + async fn load( self: Box<Self>, name: String, @@ -83,6 +91,10 @@ impl<T: CatalogBuilder + 'static> BoxedCatalogBuilder for T { )) } + fn with_runtime(self: Box<Self>, runtime: Runtime) -> Box<dyn BoxedCatalogBuilder> { + Box::new(CatalogBuilder::with_runtime(*self, runtime)) + } + async fn load( self: Box<Self>, name: String, @@ -138,13 +150,16 @@ impl CatalogLoader<'_> { #[cfg(test)] mod tests { use std::collections::HashMap; - use std::sync::Arc; + use std::sync::{Arc, Mutex}; - use iceberg::io::LocalFsStorageFactory; + use iceberg::encryption::kms::KmsClientFactory; + use iceberg::io::{LocalFsStorageFactory, StorageFactory}; + use iceberg::memory::{MEMORY_CATALOG_WAREHOUSE, MemoryCatalog, MemoryCatalogBuilder}; + use iceberg::{CatalogBuilder, Result, Runtime}; use sqlx::migrate::MigrateDatabase; use tempfile::TempDir; - use crate::{CatalogLoader, load}; + use crate::{BoxedCatalogBuilder, CatalogLoader, load}; #[tokio::test] async fn test_load_unsupported_catalog() { @@ -287,6 +302,77 @@ mod tests { assert!(catalog.is_ok()); } + /// A memory catalog builder that records the runtime it is given. + #[derive(Debug, Default)] + struct RuntimeRecordingBuilder { + runtime: Arc<Mutex<Option<Runtime>>>, + } + + impl CatalogBuilder for RuntimeRecordingBuilder { + type C = MemoryCatalog; + + fn with_storage_factory(self, _storage_factory: Arc<dyn StorageFactory>) -> Self { + self + } + + fn with_kms_client_factory(self, _kms_client_factory: Arc<dyn KmsClientFactory>) -> Self { + self + } + + fn with_runtime(self, runtime: Runtime) -> Self { + *self.runtime.lock().unwrap() = Some(runtime); + self + } + + fn load( + self, + name: impl Into<String>, + props: HashMap<String, String>, + ) -> impl Future<Output = Result<MemoryCatalog>> + Send { + MemoryCatalogBuilder::default().load(name, props) + } + } + + #[test] + fn test_with_runtime_reaches_the_catalog_builder() { + let tokio_runtime = tokio::runtime::Builder::new_multi_thread() + .worker_threads(1) + .thread_name("loader-test-runtime") + .enable_all() + .build() + .unwrap(); + let recorded = Arc::new(Mutex::new(None)); + let builder: Box<dyn BoxedCatalogBuilder> = Box::new(RuntimeRecordingBuilder { + runtime: recorded.clone(), + }); + + tokio_runtime.block_on(async { + builder + .with_runtime(Runtime::new(&tokio_runtime)) + .load( + "memory".to_string(), + HashMap::from([(MEMORY_CATALOG_WAREHOUSE.to_string(), temp_path())]), + ) + .await + .unwrap(); + + // The builder got the runtime passed to the boxed builder: its + // tasks run on that runtime's threads. + let runtime = recorded.lock().unwrap().take().unwrap(); + let thread = runtime + .io() + .spawn(async { std::thread::current().name().map(str::to_string) }) + .await + .unwrap(); + assert!( + thread + .as_deref() + .is_some_and(|n| n.starts_with("loader-test-runtime")), + "got: {thread:?}" + ); + }); + } + #[tokio::test] async fn test_error_message_includes_supported_types() { let err = match load("does-not-exist") {
