NoahKusaba commented on code in PR #3286:
URL: https://github.com/apache/iceberg-rust/pull/3286#discussion_r4154561828


##########
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:
   Added matching one-liners to `with_storage_factory` and 
`with_kms_client_factory`, so all three methods are documented.



##########
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:
   I'd like to keep it required. The only types that implement 
`BoxedCatalogBuilder` directly are wrappers around another `Box<dyn 
BoxedCatalogBuilder>`, and for those a `{ self }` default would compile but 
drop the runtime, falling back to `Runtime::current()`, which is the silent 
fallback this PR is fixing. Keeping it required makes the compiler flag 
wrappers that need to forward it. It also matches 
`CatalogBuilder::with_runtime` in core, which has no default. I've added a 
Breaking Changes entry to CHANGELOG.md so the break is called out.



##########
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:
   Done, it's `.take()` now.



##########
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:
   Done, it now checks the prefix and prints the name on failure.



-- 
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]

Reply via email to