use std::sync::Arc;
use datafusion::error::{DataFusionError, Result as DataFusionResult};
use datafusion::execution::object_store::{DefaultObjectStoreRegistry, ObjectStoreRegistry};
use object_store::ObjectStore;
use url::Url;
#[derive(Debug, Default)]
pub struct LazyCloudObjectStoreRegistry {
inner: DefaultObjectStoreRegistry,
}
impl LazyCloudObjectStoreRegistry {
pub fn new() -> Self {
Self::default()
}
}
impl ObjectStoreRegistry for LazyCloudObjectStoreRegistry {
fn register_store(
&self,
url: &Url,
store: Arc<dyn ObjectStore>,
) -> Option<Arc<dyn ObjectStore>> {
self.inner.register_store(url, store)
}
fn get_store(&self, url: &Url) -> DataFusionResult<Arc<dyn ObjectStore>> {
if let Ok(store) = self.inner.get_store(url) {
return Ok(store);
}
if matches!(url.scheme(), "s3" | "s3a") {
let bucket = url.host_str().unwrap_or_default();
if bucket.is_empty() {
return Err(DataFusionError::Execution(format!(
"object-store url {url} has no bucket"
)));
}
let store = crate::build_s3_object_store(bucket).map_err(|error| {
DataFusionError::Execution(format!(
"cannot build an S3 object store for bucket '{bucket}': {error}"
))
})?;
let store: Arc<dyn ObjectStore> = store;
self.inner.register_store(url, Arc::clone(&store));
return Ok(store);
}
self.inner.get_store(url)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn an_unregistered_s3_bucket_resolves() {
let registry = LazyCloudObjectStoreRegistry::new();
let url = Url::parse("s3://krishiv-bench/tpch/sf100/lineitem/").expect("url");
assert!(
registry.get_store(&url).is_ok(),
"an s3 bucket must resolve without prior registration"
);
}
#[test]
fn repeated_lookups_reuse_one_store() {
let registry = LazyCloudObjectStoreRegistry::new();
let url = Url::parse("s3://krishiv-bench/a").expect("url");
let first = registry.get_store(&url).expect("first lookup");
let second = registry.get_store(&url).expect("second lookup");
assert!(
Arc::ptr_eq(&first, &second),
"the store must be cached, not rebuilt per lookup"
);
}
#[test]
fn explicit_registration_wins() {
let registry = LazyCloudObjectStoreRegistry::new();
let url = Url::parse("s3://explicit-bucket/").expect("url");
let installed: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
registry.register_store(&url, Arc::clone(&installed));
let resolved = registry.get_store(&url).expect("lookup");
assert!(
Arc::ptr_eq(&installed, &resolved),
"an explicitly registered store must win over lazy construction"
);
}
#[test]
fn unknown_schemes_keep_the_default_error() {
let registry = LazyCloudObjectStoreRegistry::new();
let url = Url::parse("ftp://somewhere/path").expect("url");
assert!(registry.get_store(&url).is_err());
}
}