fix(rust): remove diagnostic provider detection

This commit is contained in:
Gatefixer
2026-08-06 21:43:22 +00:00
parent 7db2628c87
commit 79b65247bd
4 changed files with 85 additions and 62 deletions
+4 -13
View File
@@ -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<LanceNamespaceDatabase> {
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<String, String>,
session: Option<Arc<lance::session::Session>>,
) -> Result<Self> {
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,
+2 -5
View File
@@ -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);
+73 -22
View File
@@ -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<dyn ObjectStoreProvider>,
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<Mutex<Vec<Weak<ObjectStoreRegistry>>>> =
/// 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(&registry));
}
/// 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<lance::session::Session>>,
) -> Arc<lance::session::Session> {
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<dyn ObjectStore>,
}
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(), &params)
.await
.unwrap();
let session = atomic_aws_session(None);
session
.store_registry()
.get_provider("s3")
.unwrap()
.new_store(url::Url::parse("s3://bucket/table").unwrap(), &params)
.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(), &params)
.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(), &params)
.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<dyn ObjectStore> = 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 {
+6 -22
View File
@@ -2257,11 +2257,7 @@ impl NativeTable {
managed_versioning: Option<bool>,
) -> Result<Self> {
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