mirror of
https://github.com/lancedb/lancedb.git
synced 2026-08-18 12:08:35 +00:00
fix(rust): preserve AWS object store contracts
This commit is contained in:
@@ -1336,9 +1336,11 @@ mod tests {
|
||||
|
||||
let write_params = db.prepare_write_params(&request, None, None, None);
|
||||
|
||||
let store_params = write_params.store_params.unwrap();
|
||||
assert!(store_params.storage_options_accessor.is_some());
|
||||
assert!(
|
||||
write_params.store_params.unwrap().aws_credentials.is_some(),
|
||||
"operation-only explicit credentials must be installed atomically"
|
||||
store_params.aws_credentials.is_none(),
|
||||
"credential allocation must happen after the registry cache lookup"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -1361,15 +1363,11 @@ mod tests {
|
||||
|
||||
db.prepare_open_table_request(&mut request);
|
||||
|
||||
let store_params = request.lance_read_params.unwrap().store_options.unwrap();
|
||||
assert!(store_params.storage_options_accessor.is_some());
|
||||
assert!(
|
||||
request
|
||||
.lance_read_params
|
||||
.unwrap()
|
||||
.store_options
|
||||
.unwrap()
|
||||
.aws_credentials
|
||||
.is_some(),
|
||||
"operation-only explicit credentials must be installed atomically"
|
||||
store_params.aws_credentials.is_none(),
|
||||
"credential allocation must happen after the registry cache lookup"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -102,8 +102,9 @@ impl LanceNamespaceDatabase {
|
||||
session: Option<Arc<lance::session::Session>>,
|
||||
namespace_client_pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
|
||||
) -> Self {
|
||||
let session = session.unwrap_or_else(|| Arc::new(lance::session::Session::default()));
|
||||
install_atomic_aws_provider(&session);
|
||||
if let Some(session) = &session {
|
||||
install_atomic_aws_provider(session);
|
||||
}
|
||||
// Client is pre-built, so we can't install the freshness provider here;
|
||||
// baselines are still tracked for a uniform bump path.
|
||||
let delimiter = resolve_delimiter(&namespace_client_properties);
|
||||
@@ -111,7 +112,7 @@ impl LanceNamespaceDatabase {
|
||||
namespace: namespace_client,
|
||||
storage_options,
|
||||
read_consistency_interval,
|
||||
session: Some(session),
|
||||
session,
|
||||
uri: format!("namespace://{}", namespace_client_impl),
|
||||
pushdown_operations: namespace_client_pushdown_operations,
|
||||
ns_impl: namespace_client_impl,
|
||||
@@ -156,13 +157,18 @@ impl LanceNamespaceDatabase {
|
||||
pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
|
||||
new_table_config: NewTableConfig,
|
||||
) -> Result<Self> {
|
||||
let session = session.unwrap_or_else(|| Arc::new(lance::session::Session::default()));
|
||||
install_atomic_aws_provider(&session);
|
||||
// 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 mut builder = ConnectBuilder::new(ns_impl);
|
||||
for (key, value) in ns_properties.clone() {
|
||||
builder = builder.property(key, value);
|
||||
}
|
||||
builder = builder.session(session.clone());
|
||||
builder = builder.session(builder_session);
|
||||
|
||||
// Install the read-freshness provider before building the client.
|
||||
let freshness_baselines: FreshnessBaselines = Arc::new(Mutex::new(HashMap::new()));
|
||||
@@ -180,7 +186,7 @@ impl LanceNamespaceDatabase {
|
||||
namespace,
|
||||
storage_options,
|
||||
read_consistency_interval,
|
||||
session: Some(session),
|
||||
session,
|
||||
uri: format!("namespace://{}", ns_impl),
|
||||
pushdown_operations,
|
||||
ns_impl: ns_impl.to_string(),
|
||||
@@ -631,9 +637,37 @@ mod tests {
|
||||
use crate::query::ExecutableQuery;
|
||||
use arrow_array::{Int32Array, RecordBatch, StringArray};
|
||||
use arrow_schema::{DataType, Field, Schema};
|
||||
use async_trait::async_trait;
|
||||
use futures::TryStreamExt;
|
||||
use lance_io::object_store::{
|
||||
ObjectStore as LanceObjectStore, ObjectStoreParams, ObjectStoreProvider,
|
||||
ObjectStoreRegistry, uri_to_url,
|
||||
};
|
||||
use tempfile::tempdir;
|
||||
|
||||
#[derive(Debug)]
|
||||
struct FailingFileProvider;
|
||||
|
||||
#[async_trait]
|
||||
impl ObjectStoreProvider for FailingFileProvider {
|
||||
async fn new_store(
|
||||
&self,
|
||||
_base_path: url::Url,
|
||||
_params: &ObjectStoreParams,
|
||||
) -> lance_core::Result<LanceObjectStore> {
|
||||
Err(lance_core::Error::invalid_input(
|
||||
"operation-supplied session was used",
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
fn file_object_store_uri(path: &str) -> String {
|
||||
let file_url = uri_to_url(path).unwrap();
|
||||
let mut url = url::Url::parse("file-object-store:///").unwrap();
|
||||
url.set_path(file_url.path());
|
||||
url.to_string()
|
||||
}
|
||||
|
||||
/// Helper function to create test data
|
||||
fn create_test_data() -> RecordBatch {
|
||||
let schema = Arc::new(Schema::new(vec![
|
||||
@@ -667,7 +701,81 @@ mod tests {
|
||||
.downcast_ref::<LanceNamespaceDatabase>()
|
||||
.unwrap();
|
||||
|
||||
assert!(database.session.is_some());
|
||||
assert!(
|
||||
database.session.is_none(),
|
||||
"a builder-only default session must not override operation parameters"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn operation_read_session_takes_precedence_when_connection_has_no_session() {
|
||||
let tmp_dir = tempdir().unwrap();
|
||||
let properties = HashMap::from([(
|
||||
"root".to_string(),
|
||||
file_object_store_uri(tmp_dir.path().to_str().unwrap()),
|
||||
)]);
|
||||
let connection = connect_namespace("dir", properties)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
connection
|
||||
.create_table("test", create_test_data())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let registry = Arc::new(ObjectStoreRegistry::default());
|
||||
registry.insert("file-object-store", Arc::new(FailingFileProvider));
|
||||
let read_params = lance::dataset::ReadParams {
|
||||
session: Some(Arc::new(lance::session::Session::new(16, 16, registry))),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let error = connection
|
||||
.open_table("test")
|
||||
.lance_read_params(read_params)
|
||||
.execute()
|
||||
.await
|
||||
.expect_err("operation-supplied session was silently replaced");
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("operation-supplied session was used")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn operation_write_session_takes_precedence_when_connection_has_no_session() {
|
||||
let tmp_dir = tempdir().unwrap();
|
||||
let properties = HashMap::from([(
|
||||
"root".to_string(),
|
||||
file_object_store_uri(tmp_dir.path().to_str().unwrap()),
|
||||
)]);
|
||||
let connection = connect_namespace("dir", properties)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let registry = Arc::new(ObjectStoreRegistry::default());
|
||||
registry.insert("file-object-store", Arc::new(FailingFileProvider));
|
||||
let write_params = lance::dataset::WriteParams {
|
||||
session: Some(Arc::new(lance::session::Session::new(16, 16, registry))),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let error = connection
|
||||
.create_table("test", create_test_data())
|
||||
.write_options(crate::table::WriteOptions {
|
||||
lance_write_params: Some(write_params),
|
||||
})
|
||||
.execute()
|
||||
.await
|
||||
.expect_err("operation-supplied session was silently replaced");
|
||||
assert!(
|
||||
error
|
||||
.to_string()
|
||||
.contains("operation-supplied session was used")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
+667
-143
@@ -6,13 +6,21 @@
|
||||
use std::{collections::HashMap, fmt::Formatter, sync::Arc};
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
use std::sync::{LazyLock, Mutex, Weak};
|
||||
use std::{
|
||||
ops::Range,
|
||||
sync::{LazyLock, Mutex, Weak},
|
||||
};
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
use bytes::Bytes;
|
||||
use futures::{StreamExt, TryFutureExt, stream::BoxStream};
|
||||
#[cfg(feature = "aws")]
|
||||
use futures::{TryStreamExt, stream};
|
||||
use lance::io::{ObjectStoreParams, WrappingObjectStore};
|
||||
#[cfg(feature = "aws")]
|
||||
use lance_io::object_store::{
|
||||
ObjectStore as LanceObjectStore, ObjectStoreProvider, ObjectStoreRegistry,
|
||||
ObjectStore as LanceObjectStore, ObjectStoreProvider, ObjectStoreRegistry, StorageOptions,
|
||||
providers::aws::build_aws_credential,
|
||||
throttle::{AimdThrottleConfig, AimdThrottledStore},
|
||||
};
|
||||
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
|
||||
@@ -23,7 +31,7 @@ use object_store::{
|
||||
};
|
||||
#[cfg(feature = "aws")]
|
||||
use object_store::{
|
||||
StaticCredentialProvider,
|
||||
CredentialProvider, RenameOptions, StaticCredentialProvider,
|
||||
aws::{AmazonS3ConfigKey, AwsCredential},
|
||||
};
|
||||
#[cfg(feature = "aws")]
|
||||
@@ -32,22 +40,14 @@ use object_store_opendal::OpendalStore;
|
||||
use opendal::{Operator, services::S3};
|
||||
#[cfg(feature = "aws")]
|
||||
use std::str::FromStr;
|
||||
#[cfg(feature = "aws")]
|
||||
use tokio::sync::RwLock as TokioRwLock;
|
||||
|
||||
use async_trait::async_trait;
|
||||
|
||||
#[cfg(test)]
|
||||
pub mod io_tracking;
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
fn explicit_aws_credentials(
|
||||
storage_options: &HashMap<String, String>,
|
||||
) -> Option<object_store::aws::AwsCredentialProvider> {
|
||||
explicit_aws_credential(storage_options)
|
||||
.ok()
|
||||
.flatten()
|
||||
.map(|credential| Arc::new(StaticCredentialProvider::new(credential)) as _)
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
fn explicit_aws_credential(
|
||||
storage_options: &HashMap<String, String>,
|
||||
@@ -80,6 +80,14 @@ fn explicit_aws_credential(
|
||||
}))
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
fn has_aws_credential_member(storage_options: &HashMap<String, String>) -> bool {
|
||||
storage_options.keys().any(|key| {
|
||||
AmazonS3ConfigKey::from_str(&key.to_ascii_lowercase())
|
||||
.is_ok_and(|key| is_aws_credential_key(&key))
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
fn is_aws_credential_key(key: &AmazonS3ConfigKey) -> bool {
|
||||
matches!(
|
||||
@@ -90,42 +98,24 @@ fn is_aws_credential_key(key: &AmazonS3ConfigKey) -> bool {
|
||||
)
|
||||
}
|
||||
|
||||
/// Resolve OpenDAL options while keeping one explicit AWS credential family atomic.
|
||||
///
|
||||
/// Lance's native S3 provider accepts an explicit credential provider, but OpenDAL does not.
|
||||
/// For that backend we therefore merge non-credential environment options here and disable
|
||||
/// OpenDAL's second credential lookup. Wholly ambient credentials never enter this path.
|
||||
#[cfg(feature = "aws")]
|
||||
fn atomic_opendal_options(
|
||||
fn canonical_noncredential_options(
|
||||
storage_options: &HashMap<String, String>,
|
||||
environment: impl IntoIterator<Item = (String, String)>,
|
||||
) -> lance_core::Result<Option<HashMap<String, String>>> {
|
||||
let Some(credential) = explicit_aws_credential(storage_options)? else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let mut options = HashMap::new();
|
||||
for (key, value) in storage_options {
|
||||
match AmazonS3ConfigKey::from_str(&key.to_ascii_lowercase()) {
|
||||
Ok(config_key) if is_aws_credential_key(&config_key) => {}
|
||||
Ok(config_key) => {
|
||||
options.insert(config_key.as_ref().to_string(), value.clone());
|
||||
}
|
||||
Err(_) => {
|
||||
options.insert(key.clone(), value.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
for (key, value) in environment {
|
||||
if let Ok(config_key) = AmazonS3ConfigKey::from_str(&key.to_ascii_lowercase())
|
||||
&& !is_aws_credential_key(&config_key)
|
||||
{
|
||||
options
|
||||
.entry(config_key.as_ref().to_string())
|
||||
.or_insert(value);
|
||||
}
|
||||
}
|
||||
) -> HashMap<String, String> {
|
||||
storage_options
|
||||
.iter()
|
||||
.filter_map(
|
||||
|(key, value)| match AmazonS3ConfigKey::from_str(&key.to_ascii_lowercase()) {
|
||||
Ok(config_key) if is_aws_credential_key(&config_key) => None,
|
||||
Ok(config_key) => Some((config_key.as_ref().to_string(), value.clone())),
|
||||
Err(_) => Some((key.clone(), value.clone())),
|
||||
},
|
||||
)
|
||||
.collect()
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
fn insert_aws_credential(options: &mut HashMap<String, String>, credential: AwsCredential) {
|
||||
options.insert(
|
||||
AmazonS3ConfigKey::AccessKeyId.as_ref().to_string(),
|
||||
credential.key_id,
|
||||
@@ -137,8 +127,280 @@ fn atomic_opendal_options(
|
||||
if let Some(token) = credential.token {
|
||||
options.insert(AmazonS3ConfigKey::Token.as_ref().to_string(), token);
|
||||
}
|
||||
options.insert("disable_config_load".to_string(), "true".to_string());
|
||||
Ok(Some(options))
|
||||
}
|
||||
|
||||
/// Merge an OpenDAL configuration without ever combining two AWS credential families.
|
||||
#[cfg(feature = "aws")]
|
||||
fn atomic_opendal_options(
|
||||
base_options: &HashMap<String, String>,
|
||||
dynamic_options: &HashMap<String, String>,
|
||||
credential: Option<AwsCredential>,
|
||||
environment: impl IntoIterator<Item = (String, String)>,
|
||||
) -> lance_core::Result<HashMap<String, String>> {
|
||||
let mut options = canonical_noncredential_options(base_options);
|
||||
for (key, value) in environment {
|
||||
match AmazonS3ConfigKey::from_str(&key.to_ascii_lowercase()) {
|
||||
Ok(config_key) if is_aws_credential_key(&config_key) => {}
|
||||
Ok(config_key) => {
|
||||
options
|
||||
.entry(config_key.as_ref().to_string())
|
||||
.or_insert(value);
|
||||
}
|
||||
Err(_) => {}
|
||||
}
|
||||
}
|
||||
options.extend(canonical_noncredential_options(dynamic_options));
|
||||
|
||||
let credential = match credential {
|
||||
Some(credential) => Some(credential),
|
||||
None if has_aws_credential_member(dynamic_options) => {
|
||||
explicit_aws_credential(dynamic_options)?
|
||||
}
|
||||
None => explicit_aws_credential(base_options)?,
|
||||
};
|
||||
if let Some(credential) = credential {
|
||||
insert_aws_credential(&mut options, credential);
|
||||
// OpenDAL must not run another credential lookup after an explicit family wins.
|
||||
options.insert("disable_config_load".to_string(), "true".to_string());
|
||||
}
|
||||
Ok(options)
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
#[derive(Debug)]
|
||||
struct AtomicAccessorAwsCredentialProvider {
|
||||
accessor: Arc<StorageOptionsAccessor>,
|
||||
fallback: Option<object_store::aws::AwsCredentialProvider>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
#[async_trait]
|
||||
impl CredentialProvider for AtomicAccessorAwsCredentialProvider {
|
||||
type Credential = AwsCredential;
|
||||
|
||||
async fn get_credential(&self) -> object_store::Result<Arc<Self::Credential>> {
|
||||
let options = self
|
||||
.accessor
|
||||
.get_storage_options()
|
||||
.await
|
||||
.map_err(|error| Error::Generic {
|
||||
store: "AtomicAwsCredentialProvider",
|
||||
source: Box::new(error),
|
||||
})?
|
||||
.0;
|
||||
match explicit_aws_credential(&options).map_err(|error| Error::Generic {
|
||||
store: "AtomicAwsCredentialProvider",
|
||||
source: Box::new(error),
|
||||
})? {
|
||||
Some(credential) => Ok(Arc::new(credential)),
|
||||
None => match &self.fallback {
|
||||
Some(fallback) => fallback.get_credential().await,
|
||||
None => Err(Error::Generic {
|
||||
store: "AtomicAwsCredentialProvider",
|
||||
source: "Explicit AWS credentials require both aws_access_key_id and aws_secret_access_key".into(),
|
||||
}),
|
||||
},
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
#[derive(Debug, Clone)]
|
||||
struct CachedOpenDalStore {
|
||||
config: HashMap<String, String>,
|
||||
store: Arc<OpendalStore>,
|
||||
}
|
||||
|
||||
/// OpenDAL store that refreshes namespace-vended credentials and caches by normalized config.
|
||||
#[cfg(feature = "aws")]
|
||||
#[derive(Clone)]
|
||||
struct AtomicOpenDalStore {
|
||||
base_options: Arc<HashMap<String, String>>,
|
||||
accessor: Option<Arc<StorageOptionsAccessor>>,
|
||||
aws_credentials: Option<object_store::aws::AwsCredentialProvider>,
|
||||
bucket: Arc<str>,
|
||||
has_root: bool,
|
||||
cache: Arc<TokioRwLock<Option<CachedOpenDalStore>>>,
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
impl std::fmt::Debug for AtomicOpenDalStore {
|
||||
fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
|
||||
formatter
|
||||
.debug_struct("AtomicOpenDalStore")
|
||||
.field("bucket", &self.bucket)
|
||||
.field("accessor", &self.accessor)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
impl std::fmt::Display for AtomicOpenDalStore {
|
||||
fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result {
|
||||
write!(formatter, "AtomicOpenDalStore({})", self.bucket)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
impl AtomicOpenDalStore {
|
||||
async fn current_store(&self) -> lance_core::Result<Arc<OpendalStore>> {
|
||||
let dynamic_options = match &self.accessor {
|
||||
Some(accessor) if accessor.has_provider() => accessor.get_storage_options().await?.0,
|
||||
_ => HashMap::new(),
|
||||
};
|
||||
let credential = match &self.aws_credentials {
|
||||
Some(provider) => Some({
|
||||
let credential = provider
|
||||
.get_credential()
|
||||
.await
|
||||
.map_err(|error| lance_core::Error::io_source(Box::new(error)))?;
|
||||
AwsCredential {
|
||||
key_id: credential.key_id.clone(),
|
||||
secret_key: credential.secret_key.clone(),
|
||||
token: credential.token.clone(),
|
||||
}
|
||||
}),
|
||||
None => None,
|
||||
};
|
||||
let mut config = atomic_opendal_options(
|
||||
&self.base_options,
|
||||
&dynamic_options,
|
||||
credential,
|
||||
std::env::vars_os().filter_map(|(key, value)| {
|
||||
Some((key.into_string().ok()?, value.into_string().ok()?))
|
||||
}),
|
||||
)?;
|
||||
config.insert("bucket".to_string(), self.bucket.to_string());
|
||||
if self.has_root {
|
||||
config.insert("root".to_string(), "/".to_string());
|
||||
} else {
|
||||
config.remove("root");
|
||||
}
|
||||
|
||||
{
|
||||
let cache = self.cache.read().await;
|
||||
if let Some(cached) = cache.as_ref()
|
||||
&& cached.config == config
|
||||
{
|
||||
return Ok(cached.store.clone());
|
||||
}
|
||||
}
|
||||
|
||||
let operator = Operator::from_iter::<S3>(config.clone()).map_err(|error| {
|
||||
lance_core::Error::invalid_input(format!("Failed to create S3 operator: {error:?}"))
|
||||
})?;
|
||||
let store = Arc::new(OpendalStore::new(operator));
|
||||
let mut cache = self.cache.write().await;
|
||||
if let Some(cached) = cache.as_ref()
|
||||
&& cached.config == config
|
||||
{
|
||||
return Ok(cached.store.clone());
|
||||
}
|
||||
*cache = Some(CachedOpenDalStore {
|
||||
config,
|
||||
store: store.clone(),
|
||||
});
|
||||
Ok(store)
|
||||
}
|
||||
|
||||
fn map_store_error(error: lance_core::Error) -> Error {
|
||||
Error::Generic {
|
||||
store: "AtomicOpenDalStore",
|
||||
source: Box::new(error),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
#[async_trait]
|
||||
impl ObjectStore for AtomicOpenDalStore {
|
||||
async fn put_opts(
|
||||
&self,
|
||||
location: &Path,
|
||||
payload: PutPayload,
|
||||
options: PutOptions,
|
||||
) -> Result<PutResult> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.put_opts(location, payload, options)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn put_multipart_opts(
|
||||
&self,
|
||||
location: &Path,
|
||||
options: PutMultipartOptions,
|
||||
) -> Result<Box<dyn MultipartUpload>> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.put_multipart_opts(location, options)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.get_opts(location, options)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn get_ranges(&self, location: &Path, ranges: &[Range<u64>]) -> Result<Vec<Bytes>> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.get_ranges(location, ranges)
|
||||
.await
|
||||
}
|
||||
|
||||
fn delete_stream(
|
||||
&self,
|
||||
locations: BoxStream<'static, Result<Path>>,
|
||||
) -> BoxStream<'static, Result<Path>> {
|
||||
let this = self.clone();
|
||||
stream::once(async move {
|
||||
let store = this.current_store().await.map_err(Self::map_store_error)?;
|
||||
Ok::<_, Error>((store, locations))
|
||||
})
|
||||
.map_ok(|(store, locations)| store.delete_stream(locations))
|
||||
.try_flatten()
|
||||
.boxed()
|
||||
}
|
||||
|
||||
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
|
||||
let prefix = prefix.cloned();
|
||||
let this = self.clone();
|
||||
stream::once(async move { this.current_store().await.map_err(Self::map_store_error) })
|
||||
.map_ok(move |store| store.list(prefix.as_ref()))
|
||||
.try_flatten()
|
||||
.boxed()
|
||||
}
|
||||
|
||||
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.list_with_delimiter(prefix)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.copy_opts(from, to, options)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
|
||||
self.current_store()
|
||||
.await
|
||||
.map_err(Self::map_store_error)?
|
||||
.rename_opts(from, to, options)
|
||||
.await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "aws")]
|
||||
@@ -155,41 +417,11 @@ impl ObjectStoreProvider for AtomicAwsStoreProvider {
|
||||
base_path: url::Url,
|
||||
params: &ObjectStoreParams,
|
||||
) -> lance_core::Result<LanceObjectStore> {
|
||||
// Caller-supplied credential providers and refreshable storage options are already
|
||||
// atomic credential authorities. Preserve Lance's precedence and refresh behavior.
|
||||
if params.aws_credentials.is_some()
|
||||
|| params
|
||||
.storage_options_accessor
|
||||
.as_ref()
|
||||
.is_some_and(|accessor| accessor.has_provider())
|
||||
{
|
||||
return self.inner.new_store(base_path, params).await;
|
||||
}
|
||||
|
||||
let storage_options = params.storage_options().cloned().unwrap_or_default();
|
||||
let Some(credential) = explicit_aws_credential(&storage_options)? else {
|
||||
return self.inner.new_store(base_path, params).await;
|
||||
};
|
||||
|
||||
let use_opendal = storage_options
|
||||
.get("use_opendal")
|
||||
.is_some_and(|value| value == "true");
|
||||
|
||||
// The native provider gives an explicit credential provider highest precedence, so its
|
||||
// later environment merge cannot splice in an unrelated session token. It also supplies
|
||||
// the correctly initialized Lance ObjectStore shell used below for OpenDAL.
|
||||
let mut native_options = storage_options.clone();
|
||||
native_options.insert("use_opendal".to_string(), "false".to_string());
|
||||
let mut native_params = params.clone();
|
||||
native_params.storage_options_accessor = Some(Arc::new(
|
||||
StorageOptionsAccessor::with_static_options(native_options),
|
||||
));
|
||||
native_params.aws_credentials = Some(Arc::new(StaticCredentialProvider::new(credential)));
|
||||
let mut store = self
|
||||
.inner
|
||||
.new_store(base_path.clone(), &native_params)
|
||||
.await?;
|
||||
|
||||
if use_opendal {
|
||||
if storage_options
|
||||
.get("aws_provider_scheme")
|
||||
@@ -199,33 +431,110 @@ impl ObjectStoreProvider for AtomicAwsStoreProvider {
|
||||
"OpendalStore does not support an explicit aws_provider_scheme".to_string(),
|
||||
));
|
||||
}
|
||||
let mut config = atomic_opendal_options(
|
||||
&storage_options,
|
||||
std::env::vars_os().filter_map(|(key, value)| {
|
||||
Some((key.into_string().ok()?, value.into_string().ok()?))
|
||||
}),
|
||||
)?
|
||||
.expect("explicit credentials were already validated");
|
||||
let has_dynamic_options = params
|
||||
.storage_options_accessor
|
||||
.as_ref()
|
||||
.is_some_and(|accessor| accessor.has_provider());
|
||||
if params.aws_credentials.is_none()
|
||||
&& !has_dynamic_options
|
||||
&& explicit_aws_credential(&storage_options)?.is_none()
|
||||
{
|
||||
return self.inner.new_store(base_path, params).await;
|
||||
}
|
||||
|
||||
let bucket = base_path.host_str().ok_or_else(|| {
|
||||
lance_core::Error::invalid_input("S3 URL must contain bucket name")
|
||||
})?;
|
||||
config.insert("bucket".to_string(), bucket.to_string());
|
||||
if !base_path.path().trim_start_matches('/').is_empty() {
|
||||
config.insert("root".to_string(), "/".to_string());
|
||||
}
|
||||
let operator = Operator::from_iter::<S3>(config).map_err(|error| {
|
||||
lance_core::Error::invalid_input(format!("Failed to create S3 operator: {error:?}"))
|
||||
})?;
|
||||
let opendal_store: Arc<dyn ObjectStore> = Arc::new(OpendalStore::new(operator));
|
||||
let dynamic_store = AtomicOpenDalStore {
|
||||
base_options: Arc::new(storage_options.clone()),
|
||||
accessor: params.storage_options_accessor.clone(),
|
||||
aws_credentials: params.aws_credentials.clone(),
|
||||
bucket: Arc::from(bucket),
|
||||
has_root: !base_path.path().trim_start_matches('/').is_empty(),
|
||||
cache: Arc::new(TokioRwLock::new(None)),
|
||||
};
|
||||
// Preflight the actual current credential family. This rejects incomplete dynamic
|
||||
// credentials before publishing the store and primes the normalized-config cache.
|
||||
dynamic_store.current_store().await?;
|
||||
|
||||
let mut store = self.inner.new_store(base_path, params).await?;
|
||||
let opendal_store: Arc<dyn ObjectStore> = Arc::new(dynamic_store);
|
||||
let throttle_config = AimdThrottleConfig::from_storage_options(Some(&storage_options))?;
|
||||
store.inner = if throttle_config.is_disabled() {
|
||||
opendal_store
|
||||
} else {
|
||||
Arc::new(AimdThrottledStore::new(opendal_store, throttle_config)?)
|
||||
};
|
||||
return Ok(store);
|
||||
}
|
||||
|
||||
Ok(store)
|
||||
if params.aws_credentials.is_some() {
|
||||
return self.inner.new_store(base_path, params).await;
|
||||
}
|
||||
|
||||
let Some(accessor) = params.storage_options_accessor.as_ref() else {
|
||||
return self.inner.new_store(base_path, params).await;
|
||||
};
|
||||
let credential_provider: object_store::aws::AwsCredentialProvider =
|
||||
if accessor.has_provider() {
|
||||
// Validate the currently vended family first. A complete dynamic family replaces
|
||||
// the whole static family, while a provider returning no AWS options must never
|
||||
// make a partial static family fall through to Lance's environment-merged map.
|
||||
let current_options = accessor.get_storage_options().await?.0;
|
||||
let current_credential = explicit_aws_credential(¤t_options)?;
|
||||
let static_credential = explicit_aws_credential(&storage_options);
|
||||
if current_credential.is_none() {
|
||||
static_credential
|
||||
.as_ref()
|
||||
.map_err(|error| lance_core::Error::invalid_input(error.to_string()))?;
|
||||
}
|
||||
|
||||
let s3_options = storage_options
|
||||
.iter()
|
||||
.filter_map(|(key, value)| {
|
||||
AmazonS3ConfigKey::from_str(&key.to_ascii_lowercase())
|
||||
.ok()
|
||||
.map(|key| (key, value.clone()))
|
||||
})
|
||||
.collect::<HashMap<_, _>>();
|
||||
let provider_scheme =
|
||||
StorageOptions::new(storage_options.clone()).aws_provider_scheme()?;
|
||||
let region = s3_options.get(&AmazonS3ConfigKey::Region).cloned();
|
||||
let fallback = if static_credential.is_ok() {
|
||||
Some(
|
||||
build_aws_credential(
|
||||
params.s3_credentials_refresh_offset,
|
||||
None,
|
||||
Some(&s3_options),
|
||||
region,
|
||||
None,
|
||||
provider_scheme,
|
||||
)
|
||||
.await?
|
||||
.0,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
Arc::new(AtomicAccessorAwsCredentialProvider {
|
||||
accessor: accessor.clone(),
|
||||
fallback,
|
||||
})
|
||||
} else if let Some(credential) = explicit_aws_credential(&storage_options)? {
|
||||
Arc::new(StaticCredentialProvider::new(credential))
|
||||
} else {
|
||||
return self.inner.new_store(base_path, params).await;
|
||||
};
|
||||
|
||||
// This allocation occurs only after the registry cache miss. Cache identity therefore
|
||||
// remains the semantic identity of the original storage-options accessor.
|
||||
let mut atomic_params = params.clone();
|
||||
atomic_params.aws_credentials = Some(credential_provider);
|
||||
self.inner.new_store(base_path, &atomic_params).await
|
||||
}
|
||||
|
||||
fn extract_path(&self, url: &url::Url) -> lance_core::Result<Path> {
|
||||
self.inner.extract_path(url)
|
||||
}
|
||||
|
||||
fn calculate_object_store_prefix(
|
||||
@@ -235,6 +544,7 @@ impl ObjectStoreProvider for AtomicAwsStoreProvider {
|
||||
) -> lance_core::Result<String> {
|
||||
self.inner
|
||||
.calculate_object_store_prefix(url, storage_options)
|
||||
.map(|prefix| format!("{prefix}$lancedb-atomic-aws-v1"))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -271,24 +581,13 @@ pub(crate) fn install_atomic_aws_provider(_session: &lance::session::Session) {}
|
||||
|
||||
/// Apply storage options to object store parameters.
|
||||
///
|
||||
/// Static credentials for Lance's native S3 backend are installed directly. OpenDAL options are
|
||||
/// left for the session's credential-safe provider because that backend ignores this field.
|
||||
/// Caller-supplied and refreshable credential providers retain their existing precedence.
|
||||
/// Credential providers are deliberately installed by [`AtomicAwsStoreProvider`] only after a
|
||||
/// registry cache miss, preserving semantic cache reuse for identical option maps.
|
||||
pub(crate) fn set_storage_options(
|
||||
params: &mut ObjectStoreParams,
|
||||
storage_options: HashMap<String, String>,
|
||||
provider: Option<Arc<dyn StorageOptionsProvider>>,
|
||||
) {
|
||||
#[cfg(feature = "aws")]
|
||||
if provider.is_none()
|
||||
&& params.aws_credentials.is_none()
|
||||
&& !storage_options
|
||||
.get("use_opendal")
|
||||
.is_some_and(|value| value == "true")
|
||||
{
|
||||
params.aws_credentials = explicit_aws_credentials(&storage_options);
|
||||
}
|
||||
|
||||
params.storage_options_accessor = match (storage_options.is_empty(), provider) {
|
||||
(true, None) => None,
|
||||
(true, Some(provider)) => Some(Arc::new(StorageOptionsAccessor::with_provider(provider))),
|
||||
@@ -483,7 +782,7 @@ impl WrappingObjectStore for MirroringObjectStoreWrapper {
|
||||
#[cfg(all(test, feature = "aws"))]
|
||||
mod credential_tests {
|
||||
use super::*;
|
||||
use lance_io::object_store::providers::aws::build_aws_credential;
|
||||
use lance_io::object_store::providers::aws::{AwsStoreProvider, build_aws_credential};
|
||||
use std::sync::{
|
||||
Mutex,
|
||||
atomic::{AtomicBool, AtomicUsize, Ordering},
|
||||
@@ -533,6 +832,43 @@ mod credential_tests {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct NonAwsOptionsProvider;
|
||||
|
||||
#[async_trait]
|
||||
impl StorageOptionsProvider for NonAwsOptionsProvider {
|
||||
async fn fetch_storage_options(
|
||||
&self,
|
||||
) -> lance_core::Result<Option<HashMap<String, String>>> {
|
||||
Ok(Some(HashMap::from([(
|
||||
"aws_region".to_string(),
|
||||
"us-east-1".to_string(),
|
||||
)])))
|
||||
}
|
||||
|
||||
fn provider_id(&self) -> String {
|
||||
"non-aws-credential-test-provider".to_string()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct CustomPathProvider;
|
||||
|
||||
#[async_trait]
|
||||
impl ObjectStoreProvider for CustomPathProvider {
|
||||
async fn new_store(
|
||||
&self,
|
||||
_base_path: url::Url,
|
||||
_params: &ObjectStoreParams,
|
||||
) -> lance_core::Result<LanceObjectStore> {
|
||||
Err(lance_core::Error::invalid_input("unused test provider"))
|
||||
}
|
||||
|
||||
fn extract_path(&self, _url: &url::Url) -> lance_core::Result<Path> {
|
||||
Ok(Path::from("custom/tenant/path"))
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
struct ObservedCredential {
|
||||
key_id: String,
|
||||
@@ -589,35 +925,28 @@ mod credential_tests {
|
||||
"explicit-secret".to_string(),
|
||||
),
|
||||
]);
|
||||
let params = object_store_params_from_storage_options(storage_options.clone());
|
||||
let resolved_credential = Arc::new(Mutex::new(None));
|
||||
let provider = AtomicAwsStoreProvider {
|
||||
inner: Arc::new(ResolvingProvider {
|
||||
resolved_credential: resolved_credential.clone(),
|
||||
}),
|
||||
};
|
||||
|
||||
// Simulate Lance's environment merge, which adds AWS_SESSION_TOKEN when running in
|
||||
// Lambda. The explicit provider must remain an atomic two-part credential and take
|
||||
// precedence over the mixed storage options.
|
||||
let mut merged_options = storage_options
|
||||
.into_iter()
|
||||
.map(|(key, value)| (AmazonS3ConfigKey::from_str(&key).unwrap(), value))
|
||||
.collect::<HashMap<_, _>>();
|
||||
merged_options.insert(
|
||||
AmazonS3ConfigKey::Token,
|
||||
"lambda-execution-role-token".to_string(),
|
||||
provider
|
||||
.new_store(
|
||||
url::Url::parse("s3://bucket/table").unwrap(),
|
||||
&object_store_params_from_storage_options(storage_options),
|
||||
)
|
||||
.await
|
||||
.unwrap_err();
|
||||
|
||||
assert_eq!(
|
||||
*resolved_credential.lock().unwrap(),
|
||||
Some(ObservedCredential {
|
||||
key_id: "explicit-key".to_string(),
|
||||
token: None,
|
||||
})
|
||||
);
|
||||
|
||||
let (provider, _) = build_aws_credential(
|
||||
Duration::from_secs(60),
|
||||
params.aws_credentials,
|
||||
Some(&merged_options),
|
||||
Some("us-east-1".to_string()),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let credential = provider.get_credential().await.unwrap();
|
||||
|
||||
assert_eq!(credential.key_id, "explicit-key");
|
||||
assert_eq!(credential.secret_key, "explicit-secret");
|
||||
assert_eq!(credential.token, None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -639,9 +968,8 @@ mod credential_tests {
|
||||
];
|
||||
let params = object_store_params_from_storage_options(storage_options.clone());
|
||||
|
||||
let options = atomic_opendal_options(&storage_options, environment)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
let options =
|
||||
atomic_opendal_options(&storage_options, &HashMap::new(), None, environment).unwrap();
|
||||
|
||||
assert!(params.aws_credentials.is_none());
|
||||
assert_eq!(options.get("aws_access_key_id").unwrap(), "explicit-key");
|
||||
@@ -670,18 +998,51 @@ mod credential_tests {
|
||||
|
||||
let options = atomic_opendal_options(
|
||||
&storage_options,
|
||||
&HashMap::new(),
|
||||
None,
|
||||
[("AWS_SESSION_TOKEN".to_string(), "ambient-token".to_string())],
|
||||
)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(options.get("aws_session_token").unwrap(), "explicit-token");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn opendal_dynamic_credential_family_replaces_the_entire_static_family() {
|
||||
let base_options = HashMap::from([
|
||||
("aws_access_key_id".to_string(), "base-key".to_string()),
|
||||
(
|
||||
"aws_secret_access_key".to_string(),
|
||||
"base-secret".to_string(),
|
||||
),
|
||||
("aws_session_token".to_string(), "base-token".to_string()),
|
||||
]);
|
||||
let dynamic_options = HashMap::from([
|
||||
("aws_access_key_id".to_string(), "dynamic-key".to_string()),
|
||||
(
|
||||
"aws_secret_access_key".to_string(),
|
||||
"dynamic-secret".to_string(),
|
||||
),
|
||||
]);
|
||||
|
||||
let options =
|
||||
atomic_opendal_options(&base_options, &dynamic_options, None, std::iter::empty())
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(options.get("aws_access_key_id").unwrap(), "dynamic-key");
|
||||
assert_eq!(
|
||||
options.get("aws_secret_access_key").unwrap(),
|
||||
"dynamic-secret"
|
||||
);
|
||||
assert!(!options.contains_key("aws_session_token"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wholly_ambient_credentials_still_use_the_default_chain() {
|
||||
let options = atomic_opendal_options(
|
||||
&HashMap::new(),
|
||||
&HashMap::new(),
|
||||
None,
|
||||
[
|
||||
("AWS_ACCESS_KEY_ID".to_string(), "ambient-key".to_string()),
|
||||
(
|
||||
@@ -693,7 +1054,10 @@ mod credential_tests {
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert!(options.is_none());
|
||||
assert!(!options.contains_key("aws_access_key_id"));
|
||||
assert!(!options.contains_key("aws_secret_access_key"));
|
||||
assert!(!options.contains_key("aws_session_token"));
|
||||
assert!(!options.contains_key("disable_config_load"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -736,8 +1100,7 @@ mod credential_tests {
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// DirectoryNamespaceBuilder constructs fresh params with only a static accessor. The
|
||||
// OpenDAL selector is included to verify that both paths cross the installed boundary.
|
||||
// DirectoryNamespaceBuilder constructs fresh params with only a static accessor.
|
||||
let params = ObjectStoreParams {
|
||||
storage_options_accessor: Some(Arc::new(StorageOptionsAccessor::with_static_options(
|
||||
HashMap::from([
|
||||
@@ -746,7 +1109,6 @@ mod credential_tests {
|
||||
"aws_secret_access_key".to_string(),
|
||||
"explicit-secret".to_string(),
|
||||
),
|
||||
("use_opendal".to_string(), "true".to_string()),
|
||||
]),
|
||||
))),
|
||||
..Default::default()
|
||||
@@ -763,6 +1125,168 @@ mod credential_tests {
|
||||
assert!(saw_atomic_credentials.load(Ordering::SeqCst));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn opendal_refreshes_the_actual_dynamic_credential_family() {
|
||||
let fetches = Arc::new(AtomicUsize::new(0));
|
||||
let params = ObjectStoreParams {
|
||||
storage_options_accessor: Some(Arc::new(
|
||||
StorageOptionsAccessor::with_initial_and_provider(
|
||||
HashMap::from([
|
||||
("aws_access_key_id".to_string(), "expired-key".to_string()),
|
||||
(
|
||||
"aws_secret_access_key".to_string(),
|
||||
"expired-secret".to_string(),
|
||||
),
|
||||
("expires_at_millis".to_string(), "0".to_string()),
|
||||
("use_opendal".to_string(), "true".to_string()),
|
||||
("aws_region".to_string(), "us-east-1".to_string()),
|
||||
]),
|
||||
Arc::new(RotatingOptionsProvider {
|
||||
fetches: fetches.clone(),
|
||||
}),
|
||||
),
|
||||
)),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
AtomicAwsStoreProvider {
|
||||
inner: Arc::new(AwsStoreProvider),
|
||||
}
|
||||
.new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(fetches.load(Ordering::SeqCst), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn non_aws_dynamic_options_cannot_complete_a_partial_static_family() {
|
||||
let params = ObjectStoreParams {
|
||||
storage_options_accessor: Some(Arc::new(
|
||||
StorageOptionsAccessor::with_initial_and_provider(
|
||||
HashMap::from([
|
||||
("aws_access_key_id".to_string(), "explicit-key".to_string()),
|
||||
("expires_at_millis".to_string(), "0".to_string()),
|
||||
]),
|
||||
Arc::new(NonAwsOptionsProvider),
|
||||
),
|
||||
)),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
let error = AtomicAwsStoreProvider {
|
||||
inner: Arc::new(AwsStoreProvider),
|
||||
}
|
||||
.new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms)
|
||||
.await
|
||||
.unwrap_err();
|
||||
|
||||
assert!(error.to_string().contains("require both"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn complete_dynamic_credentials_replace_a_partial_static_family() {
|
||||
let fetches = Arc::new(AtomicUsize::new(0));
|
||||
let resolved_credential = Arc::new(Mutex::new(None));
|
||||
let params = ObjectStoreParams {
|
||||
storage_options_accessor: Some(Arc::new(
|
||||
StorageOptionsAccessor::with_initial_and_provider(
|
||||
HashMap::from([
|
||||
("aws_access_key_id".to_string(), "stale-key".to_string()),
|
||||
("expires_at_millis".to_string(), "0".to_string()),
|
||||
]),
|
||||
Arc::new(RotatingOptionsProvider {
|
||||
fetches: fetches.clone(),
|
||||
}),
|
||||
),
|
||||
)),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
AtomicAwsStoreProvider {
|
||||
inner: Arc::new(ResolvingProvider {
|
||||
resolved_credential: resolved_credential.clone(),
|
||||
}),
|
||||
}
|
||||
.new_store(url::Url::parse("s3://bucket/table").unwrap(), ¶ms)
|
||||
.await
|
||||
.unwrap_err();
|
||||
|
||||
assert_eq!(fetches.load(Ordering::SeqCst), 1);
|
||||
assert_eq!(
|
||||
*resolved_credential.lock().unwrap(),
|
||||
Some(ObservedCredential {
|
||||
key_id: "refreshed-key".to_string(),
|
||||
token: None,
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn wrapper_delegates_custom_path_extraction() {
|
||||
let provider = AtomicAwsStoreProvider {
|
||||
inner: Arc::new(CustomPathProvider),
|
||||
};
|
||||
|
||||
assert_eq!(
|
||||
provider
|
||||
.extract_path(&url::Url::parse("s3://bucket/original/path").unwrap())
|
||||
.unwrap(),
|
||||
Path::from("custom/tenant/path")
|
||||
);
|
||||
}
|
||||
|
||||
fn local_s3_options() -> HashMap<String, String> {
|
||||
HashMap::from([
|
||||
("aws_access_key_id".to_string(), "explicit-key".to_string()),
|
||||
(
|
||||
"aws_secret_access_key".to_string(),
|
||||
"explicit-secret".to_string(),
|
||||
),
|
||||
("aws_region".to_string(), "us-east-1".to_string()),
|
||||
("aws_endpoint".to_string(), "http://127.0.0.1:9".to_string()),
|
||||
("allow_http".to_string(), "true".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn installing_the_wrapper_does_not_reuse_a_preexisting_store() {
|
||||
let registry = Arc::new(ObjectStoreRegistry::default());
|
||||
let params = object_store_params_from_storage_options(local_s3_options());
|
||||
let url = url::Url::parse("s3://bucket/table").unwrap();
|
||||
let before = registry.get_store(url.clone(), ¶ms).await.unwrap();
|
||||
|
||||
let session = lance::session::Session::new(16, 16, registry.clone());
|
||||
install_atomic_aws_provider(&session);
|
||||
let after = registry.get_store(url, ¶ms).await.unwrap();
|
||||
|
||||
assert!(
|
||||
!Arc::ptr_eq(&before, &after),
|
||||
"the wrapper cache generation must isolate pre-install stores"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn identical_explicit_options_reuse_the_session_store() {
|
||||
let registry = Arc::new(ObjectStoreRegistry::default());
|
||||
let session = lance::session::Session::new(16, 16, registry.clone());
|
||||
install_atomic_aws_provider(&session);
|
||||
let url = url::Url::parse("s3://bucket/table").unwrap();
|
||||
let first_params = object_store_params_from_storage_options(local_s3_options());
|
||||
let second_params = object_store_params_from_storage_options(local_s3_options());
|
||||
|
||||
let first = registry
|
||||
.get_store(url.clone(), &first_params)
|
||||
.await
|
||||
.unwrap();
|
||||
let second = registry.get_store(url, &second_params).await.unwrap();
|
||||
|
||||
assert!(
|
||||
Arc::ptr_eq(&first, &second),
|
||||
"logically identical explicit credentials should hit the session cache"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn dynamic_storage_options_provider_remains_the_credential_authority() {
|
||||
let fetches = Arc::new(AtomicUsize::new(0));
|
||||
|
||||
@@ -2413,10 +2413,15 @@ impl NativeTable {
|
||||
) -> Result<Self> {
|
||||
let mut params = params.unwrap_or_default();
|
||||
|
||||
// Set the session in read params
|
||||
if let Some(sess) = session {
|
||||
params.session(sess);
|
||||
}
|
||||
// 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);
|
||||
params.session(effective_session);
|
||||
|
||||
// patch the params if we have a write store wrapper
|
||||
let params = match write_store_wrapper.clone() {
|
||||
@@ -2636,10 +2641,15 @@ impl NativeTable {
|
||||
// Start with provided params or defaults
|
||||
let mut params = params.unwrap_or_default();
|
||||
|
||||
// Set the session in write params
|
||||
if let Some(sess) = session {
|
||||
params.session = Some(sess);
|
||||
}
|
||||
// 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);
|
||||
params.session = Some(effective_session);
|
||||
|
||||
// Ensure store_params exists and set the storage options provider
|
||||
let store_params = params
|
||||
|
||||
Reference in New Issue
Block a user