From de9b217ee3d2a3a6e97bc32fffc2e64fccd0bca3 Mon Sep 17 00:00:00 2001 From: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com> Date: Thu, 6 Aug 2026 02:20:42 +0000 Subject: [PATCH] test(rust): verify concurrent table opens reuse object store --- rust/lancedb/src/database/listing.rs | 54 ++++++++++++++++++++++++++++ 1 file changed, 54 insertions(+) diff --git a/rust/lancedb/src/database/listing.rs b/rust/lancedb/src/database/listing.rs index 454498d54..d5cb6b0a4 100644 --- a/rust/lancedb/src/database/listing.rs +++ b/rust/lancedb/src/database/listing.rs @@ -1297,6 +1297,7 @@ mod tests { use crate::table::WriteOptions; use arrow_array::{Int32Array, RecordBatch, StringArray}; use arrow_schema::{DataType, Field, Schema}; + use futures::future::try_join_all; use std::path::PathBuf; use tempfile::tempdir; @@ -1322,6 +1323,59 @@ mod tests { (tempdir, db) } + #[tokio::test] + async fn test_concurrent_open_table_reuses_connection_object_store() { + let tempdir = tempdir().unwrap(); + let uri = tempdir.path().to_str().unwrap(); + let session = Arc::new(lance::session::Session::default()); + let request = ConnectRequest { + uri: uri.to_string(), + #[cfg(feature = "remote")] + client_config: Default::default(), + options: Default::default(), + namespace_client_properties: Default::default(), + manifest_enabled: false, + read_consistency_interval: None, + session: Some(session.clone()), + }; + let db = ListingDatabase::connect_with_options(&request) + .await + .unwrap(); + let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)])); + db.create_table(CreateTableRequest { + name: "test".to_string(), + namespace_path: vec![], + data: Box::new(RecordBatch::new_empty(schema)) as Box, + mode: CreateTableMode::Create, + write_options: Default::default(), + location: None, + namespace_client: None, + }) + .await + .unwrap(); + + let before = session.store_registry().stats(); + let opened_tables = try_join_all((0..32).map(|_| { + db.open_table(OpenTableRequest { + name: "test".to_string(), + namespace_path: vec![], + index_cache_size: None, + lance_read_params: None, + location: None, + namespace_client: None, + managed_versioning: None, + }) + })) + .await + .unwrap(); + let after = session.store_registry().stats(); + + assert_eq!(opened_tables.len(), 32); + assert_eq!(after.misses, before.misses); + assert_eq!(after.active_stores, before.active_stores); + assert!(after.hits >= before.hits + 32); + } + #[tokio::test] async fn test_listing_database_root_ops_do_not_create_manifest() { let tempdir = tempdir().unwrap();