This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] 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 3c91b3942 feat(catalog-loader): add with_runtime to 
BoxedCatalogBuilder (#3286)
3c91b3942 is described below

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") {

Reply via email to