laskoviymishka commented on code in PR #3286:
URL: https://github.com/apache/iceberg-rust/pull/3286#discussion_r4153153691
##########
crates/catalog/loader/src/lib.rs:
##########
@@ -57,6 +57,10 @@ pub trait BoxedCatalogBuilder: Send {
kms_client_factory: Arc<dyn KmsClientFactory>,
) -> Box<dyn BoxedCatalogBuilder>;
+ /// Sets the runtime the catalog, and the tables it creates, spawn their
Review Comment:
small thing while we're here — `with_runtime` is the only method on this
trait with a doc comment; `with_storage_factory` and `with_kms_client_factory`
have none. I'd add a matching one-liner to those two so the trait reads
consistently, or drop this one.
##########
crates/catalog/loader/src/lib.rs:
##########
@@ -287,6 +298,72 @@ 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().clone().unwrap();
+ let thread = runtime
+ .io()
+ .spawn(async {
std::thread::current().name().map(str::to_string) })
+ .await
+ .unwrap();
+ assert_eq!(thread.as_deref(), Some("loader-test-runtime"));
Review Comment:
This leans on tokio naming the worker thread exactly `loader-test-runtime`
with no suffix. Fine today with one worker, but a future tokio that appends a
counter would break it silently. A prefix check keeps the same signal without
pinning the exact string: `assert!(thread.as_deref().is_some_and(|n|
n.starts_with("loader-test-runtime")), "got: {thread:?}");`
##########
crates/catalog/loader/src/lib.rs:
##########
@@ -287,6 +298,72 @@ 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().clone().unwrap();
Review Comment:
`recorded.lock().unwrap().clone().unwrap()` clones through the `MutexGuard`
auto-deref, so it reads like it's cloning the guard rather than the inner
`Option<Runtime>`. `.take()` makes the intent explicit and skips the clone:
`let runtime = recorded.lock().unwrap().take().unwrap();`
##########
crates/catalog/loader/src/lib.rs:
##########
@@ -57,6 +57,10 @@ pub trait BoxedCatalogBuilder: Send {
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>;
Review Comment:
Adding this as a required method (no default) to the `pub
BoxedCatalogBuilder` trait breaks any downstream crate that implements the
trait directly instead of going through the blanket `impl<T: CatalogBuilder +
'static>` — a logging/decorator wrapper around `Box<dyn BoxedCatalogBuilder>`
is the natural case. It matches what `with_storage_factory` and
`with_kms_client_factory` already do, so this is consistent rather than novel,
but each required method compounds the break.
A default pass-through keeps direct implementors compiling:
```rust
fn with_runtime(self: Box<Self>, _runtime: Runtime) -> Box<dyn
BoxedCatalogBuilder> { self }
```
If we'd rather keep it required, I'd at least call out the `public-api.txt`
delta in the changelog so the break is deliberate.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]