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