diff --git a/rust/lancedb/src/database/listing.rs b/rust/lancedb/src/database/listing.rs index be360b1da..4d41daef4 100644 --- a/rust/lancedb/src/database/listing.rs +++ b/rust/lancedb/src/database/listing.rs @@ -24,7 +24,7 @@ use crate::database::ReadConsistency; use crate::database::namespace::LanceNamespaceDatabase; use crate::error::{CreateDirSnafu, Error, Result}; use crate::io::object_store::{ - MirroringObjectStoreWrapper, install_atomic_aws_provider, is_aws_credential_option, + MirroringObjectStoreWrapper, atomic_aws_session, is_aws_credential_option, object_store_params_from_storage_options, set_storage_options, }; use crate::table::NativeTable; @@ -436,11 +436,7 @@ impl ListingDatabase { request: &ConnectRequest, ) -> Result { let options = ListingDatabaseOptions::parse_from_map(&request.options)?; - let session = request - .session - .clone() - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - install_atomic_aws_provider(&session); + let session = atomic_aws_session(request.session.clone()); let namespace_root = Self::prepare_namespace_root(&request.uri, &options.storage_options, session.clone()) .await?; @@ -542,11 +538,7 @@ impl ListingDatabase { url.to_string() }; - let session = request - .session - .clone() - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - install_atomic_aws_provider(&session); + let session = atomic_aws_session(request.session.clone()); let os_params = object_store_params_from_storage_options(options.storage_options.clone()); let (object_store, base_path) = ObjectStore::from_uri_and_params( @@ -613,8 +605,7 @@ impl ListingDatabase { namespace_client_properties: HashMap, session: Option>, ) -> Result { - let session = session.unwrap_or_else(|| Arc::new(lance::session::Session::default())); - install_atomic_aws_provider(&session); + let session = atomic_aws_session(session); let (object_store, base_path) = ObjectStore::from_uri_and_params( session.store_registry(), path, diff --git a/rust/lancedb/src/database/namespace.rs b/rust/lancedb/src/database/namespace.rs index 66ba2432f..f46ad715d 100644 --- a/rust/lancedb/src/database/namespace.rs +++ b/rust/lancedb/src/database/namespace.rs @@ -34,7 +34,7 @@ use crate::database::read_freshness::{ FreshnessBaselines, ReadFreshnessContextProvider, TableFreshness, }; use crate::error::{Error, Result}; -use crate::io::object_store::install_atomic_aws_provider; +use crate::io::object_store::{atomic_aws_session, install_atomic_aws_provider}; use crate::table::{NativeTable, map_namespace_lance_error}; use lance::dataset::WriteMode; @@ -160,10 +160,7 @@ impl LanceNamespaceDatabase { // Namespace construction needs a protected session even when the connection did not // supply one. Keep the original option separately so per-operation sessions retain // precedence when tables are opened or created later. - let builder_session = session - .clone() - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - install_atomic_aws_provider(&builder_session); + let builder_session = atomic_aws_session(session.clone()); let mut builder = ConnectBuilder::new(ns_impl); for (key, value) in ns_properties.clone() { builder = builder.property(key, value); diff --git a/rust/lancedb/src/io/object_store.rs b/rust/lancedb/src/io/object_store.rs index 48f8759e8..c7d0466c0 100644 --- a/rust/lancedb/src/io/object_store.rs +++ b/rust/lancedb/src/io/object_store.rs @@ -20,7 +20,7 @@ use lance::io::{ObjectStoreParams, WrappingObjectStore}; #[cfg(feature = "aws")] use lance_io::object_store::{ ObjectStore as LanceObjectStore, ObjectStoreProvider, ObjectStoreRegistry, StorageOptions, - providers::aws::{AwsStoreProvider, build_aws_credential}, + providers::aws::build_aws_credential, throttle::{AimdThrottleConfig, AimdThrottledStore}, }; use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider}; @@ -418,19 +418,13 @@ impl ObjectStore for AtomicOpenDalStore { #[derive(Debug)] struct AtomicAwsStoreProvider { inner: Arc, + adapt_builtin_opendal: bool, } #[cfg(feature = "aws")] impl AtomicAwsStoreProvider { const CACHE_GENERATION: &'static str = "lancedb-atomic-aws-v1"; - /// Lance does not expose provider downcasting. Its built-in provider is a unit struct with a - /// stable derived Debug representation, which is the only available way to distinguish the - /// one provider whose OpenDAL implementation needs adaptation from an unknown custom one. - fn is_builtin_aws_provider(&self) -> bool { - format!("{:?}", self.inner.as_ref()) == format!("{:?}", AwsStoreProvider) - } - fn generated_prefix( &self, url: &url::Url, @@ -456,7 +450,7 @@ impl AtomicAwsStoreProvider { // its returned store may add encryption, authorization, or wrapping behavior that // cannot be reconstructed from ObjectStoreParams. Only adapt Lance's known built-in // provider, whose OpenDAL path ignores both supported credential-provider fields. - if !self.is_builtin_aws_provider() { + if !self.adapt_builtin_opendal { return self.inner.new_store(base_path, params).await; } if storage_options @@ -605,7 +599,10 @@ static ATOMIC_AWS_REGISTRIES: LazyLock>>> = /// Install the credential-safe S3 provider once on a session's shared object-store registry. #[cfg(feature = "aws")] -pub(crate) fn install_atomic_aws_provider(session: &lance::session::Session) { +fn install_atomic_aws_provider_inner( + session: &lance::session::Session, + adapt_builtin_opendal: bool, +) { let registry = session.store_registry(); let mut installed = ATOMIC_AWS_REGISTRIES .lock() @@ -621,15 +618,49 @@ pub(crate) fn install_atomic_aws_provider(session: &lance::session::Session) { for scheme in ["s3", "s3+ddb"] { if let Some(inner) = registry.get_provider(scheme) { - registry.insert(scheme, Arc::new(AtomicAwsStoreProvider { inner })); + registry.insert( + scheme, + Arc::new(AtomicAwsStoreProvider { + inner, + adapt_builtin_opendal, + }), + ); } } installed.push(Arc::downgrade(®istry)); } +/// Protect a caller-supplied session while treating every registered provider as unknown. +/// +/// Unknown providers retain unconditional control of their OpenDAL store composition. Only a +/// fresh default session created by [`atomic_aws_session`] has the explicit capability to adapt +/// Lance's built-in AWS provider. +#[cfg(feature = "aws")] +pub(crate) fn install_atomic_aws_provider(session: &lance::session::Session) { + install_atomic_aws_provider_inner(session, false); +} + #[cfg(not(feature = "aws"))] pub(crate) fn install_atomic_aws_provider(_session: &lance::session::Session) {} +/// Select a supplied session or create a fresh session with a known built-in AWS provider. +pub(crate) fn atomic_aws_session( + session: Option>, +) -> Arc { + match session { + Some(session) => { + install_atomic_aws_provider(&session); + session + } + None => { + let session = Arc::new(lance::session::Session::default()); + #[cfg(feature = "aws")] + install_atomic_aws_provider_inner(&session, true); + session + } + } +} + /// Apply storage options to object store parameters. /// /// Credential providers are deliberately installed by [`AtomicAwsStoreProvider`] only after a @@ -920,11 +951,17 @@ mod credential_tests { } } - #[derive(Debug)] struct CustomStoreProvider { marker: Arc, } + impl std::fmt::Debug for CustomStoreProvider { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + // Diagnostic output must never grant authority to replace a custom provider's store. + formatter.write_str("AwsStoreProvider") + } + } + #[async_trait] impl ObjectStoreProvider for CustomStoreProvider { async fn new_store( @@ -999,6 +1036,7 @@ mod credential_tests { inner: Arc::new(ResolvingProvider { resolved_credential: resolved_credential.clone(), }), + adapt_builtin_opendal: false, }; provider @@ -1218,12 +1256,14 @@ mod credential_tests { ..Default::default() }; - AtomicAwsStoreProvider { - inner: Arc::new(AwsStoreProvider), - } - .new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms) - .await - .unwrap(); + let session = atomic_aws_session(None); + session + .store_registry() + .get_provider("s3") + .unwrap() + .new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms) + .await + .unwrap(); assert_eq!(fetches.load(Ordering::SeqCst), 1); } @@ -1245,6 +1285,7 @@ mod credential_tests { let error = AtomicAwsStoreProvider { inner: Arc::new(AwsStoreProvider), + adapt_builtin_opendal: true, } .new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms) .await @@ -1276,6 +1317,7 @@ mod credential_tests { inner: Arc::new(ResolvingProvider { resolved_credential: resolved_credential.clone(), }), + adapt_builtin_opendal: false, } .new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms) .await @@ -1295,6 +1337,7 @@ mod credential_tests { fn wrapper_delegates_custom_path_extraction() { let provider = AtomicAwsStoreProvider { inner: Arc::new(CustomPathProvider), + adapt_builtin_opendal: false, }; assert_eq!( @@ -1338,15 +1381,21 @@ mod credential_tests { #[tokio::test] async fn opendal_preserves_custom_provider_store_behavior() { let marker: Arc = Arc::new(object_store::memory::InMemory::new()); - let provider = AtomicAwsStoreProvider { - inner: Arc::new(CustomStoreProvider { + let registry = Arc::new(ObjectStoreRegistry::default()); + registry.insert( + "s3", + Arc::new(CustomStoreProvider { marker: marker.clone(), }), - }; + ); + let session = lance::session::Session::new(16, 16, registry.clone()); + install_atomic_aws_provider(&session); let mut options = local_s3_options(); options.insert("use_opendal".to_string(), "true".to_string()); - let store = provider + let store = registry + .get_provider("s3") + .unwrap() .new_store( url::Url::parse("s3://bucket/table").unwrap(), &object_store_params_from_storage_options(options), @@ -1406,6 +1455,7 @@ mod credential_tests { inner: Arc::new(ResolvingProvider { resolved_credential: resolved_credential.clone(), }), + adapt_builtin_opendal: false, }; let params = ObjectStoreParams { storage_options_accessor: Some(Arc::new( @@ -1448,6 +1498,7 @@ mod credential_tests { inner: Arc::new(ResolvingProvider { resolved_credential: resolved_credential.clone(), }), + adapt_builtin_opendal: false, }; let params = ObjectStoreParams { aws_credentials: Some(Arc::new(StaticCredentialProvider::new(AwsCredential { diff --git a/rust/lancedb/src/table.rs b/rust/lancedb/src/table.rs index 5f1b36962..cb7b05d87 100644 --- a/rust/lancedb/src/table.rs +++ b/rust/lancedb/src/table.rs @@ -2257,11 +2257,7 @@ impl NativeTable { managed_versioning: Option, ) -> Result { let mut params = params.unwrap_or_default(); - let effective_session = params - .session - .clone() - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - crate::io::object_store::install_atomic_aws_provider(&effective_session); + let effective_session = crate::io::object_store::atomic_aws_session(params.session.clone()); params.session(effective_session); // patch the params if we have a write store wrapper let params = match write_store_wrapper.clone() { @@ -2421,12 +2417,8 @@ impl NativeTable { // Advanced operation parameters take precedence over the connection session. A fresh // session is needed only when neither source supplied one. - let effective_session = params - .session - .clone() - .or(session) - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - crate::io::object_store::install_atomic_aws_provider(&effective_session); + let effective_session = + crate::io::object_store::atomic_aws_session(params.session.clone().or(session)); params.session(effective_session); // patch the params if we have a write store wrapper @@ -2539,11 +2531,7 @@ impl NativeTable { let mut params = params.unwrap_or(WriteParams { ..Default::default() }); - let effective_session = params - .session - .clone() - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - crate::io::object_store::install_atomic_aws_provider(&effective_session); + let effective_session = crate::io::object_store::atomic_aws_session(params.session.clone()); params.session = Some(effective_session); // patch the params if we have a write store wrapper let params = match write_store_wrapper.clone() { @@ -2655,12 +2643,8 @@ impl NativeTable { // Advanced operation parameters take precedence over the connection session. A fresh // session is needed only when neither source supplied one. - let effective_session = params - .session - .clone() - .or(session) - .unwrap_or_else(|| Arc::new(lance::session::Session::default())); - crate::io::object_store::install_atomic_aws_provider(&effective_session); + let effective_session = + crate::io::object_store::atomic_aws_session(params.session.clone().or(session)); params.session = Some(effective_session); // Ensure store_params exists and set the storage options provider