Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
612 changes: 370 additions & 242 deletions Cargo.lock

Large diffs are not rendered by default.

5 changes: 4 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -117,8 +117,11 @@ mockito = "1"
motore-macros = "0.4.3"
murmur3 = "0.5.2"
once_cell = "1.20"
opendal = { version = "0.57", default-features = false, features = [
opendal = { version = "0.59", default-features = false, features = [
"executors-tokio",
# 0.59 made the HTTP transport opt-in. `iceberg-storage-opendal` installs it,
# since every object-store backend is unusable without one.
"http-transport-reqwest",
"services-memory",
# `iceberg-storage-opendal` wraps every operator in these two layers.
"layers-retry",
Expand Down
7 changes: 6 additions & 1 deletion crates/storage/opendal/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,11 @@ opendal-all = [
"opendal-hf",
]

opendal-azdls = ["opendal/services-azdls"]
opendal-azdls = [
"opendal/services-azdls",
"reqsign-azure-storage",
"reqsign-core",
]
opendal-fs = ["opendal/services-fs"]
opendal-gcs = ["opendal/services-gcs", "reqsign-google", "reqsign-core"]
opendal-hf = ["opendal/services-hf"]
Expand All @@ -56,6 +60,7 @@ futures = { workspace = true }
iceberg = { workspace = true }
opendal = { workspace = true }
reqsign-aws-v4 = { version = "3.0.0", optional = true }
reqsign-azure-storage = { version = "3.2.1", optional = true }
reqsign-core = { version = "3.0.0", optional = true }
reqsign-google = { version = "3.0.0", optional = true }
serde = { workspace = true }
Expand Down
36 changes: 35 additions & 1 deletion crates/storage/opendal/public-api.txt
Original file line number Diff line number Diff line change
@@ -1,11 +1,36 @@
pub mod iceberg_storage_opendal
pub use iceberg_storage_opendal::AwsCredential
pub use iceberg_storage_opendal::AzdlsCredential
pub use iceberg_storage_opendal::GcsCredential
pub use iceberg_storage_opendal::GcsToken
pub use iceberg_storage_opendal::ProvideCredential
pub enum iceberg_storage_opendal::AzureStorageScheme
pub iceberg_storage_opendal::AzureStorageScheme::Abfs
pub iceberg_storage_opendal::AzureStorageScheme::Abfss
pub iceberg_storage_opendal::AzureStorageScheme::Wasb
pub iceberg_storage_opendal::AzureStorageScheme::Wasbs
impl iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::as_http_scheme(&self) -> &str
impl core::clone::Clone for iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::clone(&self) -> iceberg_storage_opendal::AzureStorageScheme
impl core::cmp::PartialEq for iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::eq(&self, other: &iceberg_storage_opendal::AzureStorageScheme) -> bool
impl core::fmt::Debug for iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
impl core::fmt::Display for iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
impl core::marker::StructuralPartialEq for iceberg_storage_opendal::AzureStorageScheme
impl core::str::traits::FromStr for iceberg_storage_opendal::AzureStorageScheme
pub type iceberg_storage_opendal::AzureStorageScheme::Err = iceberg::error::Error
pub fn iceberg_storage_opendal::AzureStorageScheme::from_str(s: &str) -> iceberg::error::Result<Self>
impl serde_core::ser::Serialize for iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer
impl<'de> serde_core::de::Deserialize<'de> for iceberg_storage_opendal::AzureStorageScheme
pub fn iceberg_storage_opendal::AzureStorageScheme::deserialize<__D>(__deserializer: __D) -> core::result::Result<Self, <__D as serde_core::de::Deserializer>::Error> where __D: serde_core::de::Deserializer<'de>
pub enum iceberg_storage_opendal::OpenDalStorage
pub iceberg_storage_opendal::OpenDalStorage::Azdls
pub iceberg_storage_opendal::OpenDalStorage::Azdls::config: alloc::sync::Arc<opendal_service_azdls::config::AzdlsConfig>
pub iceberg_storage_opendal::OpenDalStorage::Azdls::customized_credential_load: core::option::Option<iceberg_storage_opendal::CustomAzdlsCredentialLoader>
pub iceberg_storage_opendal::OpenDalStorage::Gcs
pub iceberg_storage_opendal::OpenDalStorage::Gcs::config: alloc::sync::Arc<opendal_service_gcs::config::GcsConfig>
pub iceberg_storage_opendal::OpenDalStorage::Gcs::customized_credential_load: core::option::Option<iceberg_storage_opendal::CustomGcsCredentialLoader>
Expand Down Expand Up @@ -40,6 +65,7 @@ impl<'de> serde_core::de::Deserialize<'de> for iceberg_storage_opendal::OpenDalS
pub fn iceberg_storage_opendal::OpenDalStorage::deserialize<__D>(__deserializer: __D) -> core::result::Result<Self, <__D as serde_core::de::Deserializer>::Error> where __D: serde_core::de::Deserializer<'de>
pub enum iceberg_storage_opendal::OpenDalStorageFactory
pub iceberg_storage_opendal::OpenDalStorageFactory::Azdls
pub iceberg_storage_opendal::OpenDalStorageFactory::Azdls::customized_credential_load: core::option::Option<iceberg_storage_opendal::CustomAzdlsCredentialLoader>
pub iceberg_storage_opendal::OpenDalStorageFactory::Fs
pub iceberg_storage_opendal::OpenDalStorageFactory::Gcs
pub iceberg_storage_opendal::OpenDalStorageFactory::Gcs::customized_credential_load: core::option::Option<iceberg_storage_opendal::CustomGcsCredentialLoader>
Expand All @@ -60,11 +86,18 @@ impl<'de> serde_core::de::Deserialize<'de> for iceberg_storage_opendal::OpenDalS
pub fn iceberg_storage_opendal::OpenDalStorageFactory::deserialize<__D>(__deserializer: __D) -> core::result::Result<Self, <__D as serde_core::de::Deserializer>::Error> where __D: serde_core::de::Deserializer<'de>
pub struct iceberg_storage_opendal::CustomAwsCredentialLoader(_)
impl iceberg_storage_opendal::CustomAwsCredentialLoader
pub fn iceberg_storage_opendal::CustomAwsCredentialLoader::new(provider: impl reqsign_core::api::ProvideCredential<Credential = reqsign_aws_v4::credential::Credential> + 'static) -> Self
pub fn iceberg_storage_opendal::CustomAwsCredentialLoader::new(provider: impl reqsign_core::api::ProvideCredential<Credential = reqsign_aws_core::credential::Credential> + 'static) -> Self
impl core::clone::Clone for iceberg_storage_opendal::CustomAwsCredentialLoader
pub fn iceberg_storage_opendal::CustomAwsCredentialLoader::clone(&self) -> Self
impl core::fmt::Debug for iceberg_storage_opendal::CustomAwsCredentialLoader
pub fn iceberg_storage_opendal::CustomAwsCredentialLoader::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct iceberg_storage_opendal::CustomAzdlsCredentialLoader(_)
impl iceberg_storage_opendal::CustomAzdlsCredentialLoader
pub fn iceberg_storage_opendal::CustomAzdlsCredentialLoader::new(provider: impl reqsign_core::api::ProvideCredential<Credential = reqsign_azure_storage::credential::Credential> + 'static) -> Self
impl core::clone::Clone for iceberg_storage_opendal::CustomAzdlsCredentialLoader
pub fn iceberg_storage_opendal::CustomAzdlsCredentialLoader::clone(&self) -> Self
impl core::fmt::Debug for iceberg_storage_opendal::CustomAzdlsCredentialLoader
pub fn iceberg_storage_opendal::CustomAzdlsCredentialLoader::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct iceberg_storage_opendal::CustomGcsCredentialLoader(_)
impl iceberg_storage_opendal::CustomGcsCredentialLoader
pub fn iceberg_storage_opendal::CustomGcsCredentialLoader::new(provider: impl reqsign_core::api::ProvideCredential<Credential = reqsign_google::credential::Credential> + 'static) -> Self
Expand Down Expand Up @@ -94,6 +127,7 @@ pub fn iceberg_storage_opendal::OpenDalResolvingStorage::deserialize<__D>(__dese
pub struct iceberg_storage_opendal::OpenDalResolvingStorageFactory
impl iceberg_storage_opendal::OpenDalResolvingStorageFactory
pub fn iceberg_storage_opendal::OpenDalResolvingStorageFactory::new() -> Self
pub fn iceberg_storage_opendal::OpenDalResolvingStorageFactory::with_azdls_credential_loader(self, loader: iceberg_storage_opendal::CustomAzdlsCredentialLoader) -> Self
pub fn iceberg_storage_opendal::OpenDalResolvingStorageFactory::with_gcs_credential_loader(self, loader: iceberg_storage_opendal::CustomGcsCredentialLoader) -> Self
pub fn iceberg_storage_opendal::OpenDalResolvingStorageFactory::with_s3_credential_loader(self, loader: iceberg_storage_opendal::CustomAwsCredentialLoader) -> Self
impl core::clone::Clone for iceberg_storage_opendal::OpenDalResolvingStorageFactory
Expand Down
136 changes: 129 additions & 7 deletions crates/storage/opendal/src/azdls.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
use std::collections::HashMap;
use std::fmt::Display;
use std::str::FromStr;
use std::sync::Arc;

use iceberg::io::{
ADLS_ACCOUNT_KEY, ADLS_ACCOUNT_NAME, ADLS_AUTHORITY_HOST, ADLS_CLIENT_ID, ADLS_CLIENT_SECRET,
Expand All @@ -26,6 +27,8 @@ use iceberg::io::{
use iceberg::{Error, ErrorKind, Result};
use opendal::Configurator;
use opendal::services::AzdlsConfig;
pub use reqsign_azure_storage::Credential as AzdlsCredential;
use reqsign_core::{ProvideCredential, ProvideCredentialChain, ProvideCredentialDyn};
use serde::{Deserialize, Serialize};
use url::Url;

Expand Down Expand Up @@ -90,12 +93,13 @@ pub(crate) fn azdls_config_parse(mut properties: HashMap<String, String>) -> Res
/// `abfss://<myfs>@<myaccount>.dfs.core.windows.net/mydir/myfile.parquet`.
pub(crate) fn azdls_create_operator<'a>(
absolute_path: &'a str,
customized_credential_load: &Option<CustomAzdlsCredentialLoader>,
config: &AzdlsConfig,
) -> Result<(opendal::Operator, &'a str)> {
let path = absolute_path.parse::<AzureStoragePath>()?;
match_path_with_config(&path, config)?;

let op = azdls_config_build(config, &path)?;
let op = azdls_config_build(config, customized_credential_load, &path)?;

// Paths to files in ADLS tend to be written in fully qualified form,
// including their filesystem and account name.
Expand Down Expand Up @@ -192,18 +196,24 @@ pub(crate) fn match_path_with_config(path: &AzureStoragePath, config: &AzdlsConf
Ok(())
}

fn azdls_config_build(config: &AzdlsConfig, path: &AzureStoragePath) -> Result<opendal::Operator> {
fn azdls_config_build(
config: &AzdlsConfig,
customized_credential_load: &Option<CustomAzdlsCredentialLoader>,
path: &AzureStoragePath,
) -> Result<opendal::Operator> {
let mut builder = config.clone().into_builder();

if config.endpoint.is_none() {
// If no endpoint is provided, we construct it from the fully-qualified path.
builder = builder.endpoint(&path.as_endpoint());
}
builder = builder.filesystem(&path.filesystem);
if let Some(loader) = customized_credential_load {
let chain = ProvideCredentialChain::new().push(Arc::clone(&loader.0));
builder = builder.credential_provider_chain(chain);
}

Ok(opendal::Operator::new(builder)
.map_err(from_opendal_error)?
.finish())
opendal::Operator::new(builder).map_err(from_opendal_error)
}

/// Represents a fully qualified path to blob/ file in Azure Storage.
Expand Down Expand Up @@ -317,13 +327,49 @@ fn validate_storage_and_scheme(
}
}

/// Custom AZDLS credential loader.
///
/// Wraps any [`ProvideCredential`] implementation for use with the AZDLS storage backend.
/// Use [`CustomAzdlsCredentialLoader::new`] to create one, then pass it to
/// [`OpenDalStorageFactory::Azdls`](crate::OpenDalStorageFactory).
///
/// Installing a loader replaces the default credential chain opendal would otherwise reach
/// for, so the loader is the sole authority on credentials.
pub struct CustomAzdlsCredentialLoader(Arc<dyn ProvideCredentialDyn<Credential = AzdlsCredential>>);

impl Clone for CustomAzdlsCredentialLoader {
fn clone(&self) -> Self {
Self(Arc::clone(&self.0))
}
}

impl std::fmt::Debug for CustomAzdlsCredentialLoader {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("CustomAzdlsCredentialLoader")
.finish_non_exhaustive()
}
}

impl CustomAzdlsCredentialLoader {
/// Create a new custom AZDLS credential loader from any [`ProvideCredential`] implementation.
pub fn new(provider: impl ProvideCredential<Credential = AzdlsCredential> + 'static) -> Self {
Self(Arc::new(provider)
as Arc<
dyn ProvideCredentialDyn<Credential = AzdlsCredential>,
>)
}
}

#[cfg(test)]
mod tests {
use std::collections::HashMap;

use opendal::services::AzdlsConfig;

use super::{AzureStoragePath, AzureStorageScheme, azdls_config_parse, azdls_create_operator};
use super::{
AzdlsCredential, AzureStoragePath, AzureStorageScheme, CustomAzdlsCredentialLoader,
ProvideCredential, azdls_config_parse, azdls_create_operator,
};

#[test]
fn test_azdls_config_parse() {
Expand Down Expand Up @@ -466,7 +512,7 @@ mod tests {
];

for (name, input, expected) in test_cases {
let result = azdls_create_operator(input.0, &input.1);
let result = azdls_create_operator(input.0, &None, &input.1);
match expected {
Some((expected_filesystem, expected_path)) => {
assert!(result.is_ok(), "Test case {name} failed: {result:?}");
Expand All @@ -482,6 +528,82 @@ mod tests {
}
}

#[derive(Debug)]
struct StubProvider;

impl ProvideCredential for StubProvider {
type Credential = AzdlsCredential;

async fn provide_credential(
&self,
_ctx: &reqsign_core::Context,
) -> reqsign_core::Result<Option<Self::Credential>> {
Ok(None)
}
}

#[test]
fn test_azdls_create_operator_accepts_credential_loader() {
let config = AzdlsConfig {
account_name: Some("myaccount".to_string()),
endpoint: Some("https://myaccount.dfs.core.windows.net".to_string()),
..Default::default()
};
let loader = Some(CustomAzdlsCredentialLoader::new(StubProvider));

let (op, relative_path) = azdls_create_operator(
"abfss://myfs@myaccount.dfs.core.windows.net/path/to/file.parquet",
&loader,
&config,
)
.expect("operator builds with a custom credential loader");

// Installing a loader must not disturb filesystem or path resolution.
assert_eq!(op.info().name(), "myfs");
assert_eq!(relative_path, "/path/to/file.parquet");
}

#[tokio::test]
async fn test_azdls_credential_loader_replaces_configured_credentials() {
// ADLS replaces the provider chain rather than prepending to it, so a loader that
// yields nothing must fail signing outright rather than falling back to the account key
// below. Signing with that key instead would mean the loader was never installed, and
// the request would get far enough to fail for some other reason.
let config = AzdlsConfig {
account_name: Some("myaccount".to_string()),
account_key: Some("dGVzdGtleQ==".to_string()),
endpoint: Some("https://myaccount.dfs.core.windows.net".to_string()),
..Default::default()
};
let loader = Some(CustomAzdlsCredentialLoader::new(StubProvider));

let (op, relative_path) = azdls_create_operator(
"abfss://myfs@myaccount.dfs.core.windows.net/path/to/file.parquet",
&loader,
&config,
)
.expect("operator builds with a custom credential loader");

let err = op
.exists(relative_path)
.await
.expect_err("a loader that yields nothing must not fall back to the account key");
assert!(
err.to_string().contains("signing credential"),
"unexpected error: {err}"
);
}

#[test]
fn test_custom_azdls_credential_loader_debug_is_redacted() {
let loader = CustomAzdlsCredentialLoader::new(StubProvider);
assert_eq!(
format!("{loader:?}"),
"CustomAzdlsCredentialLoader { .. }",
"the loader's Debug must not reach into the provider it wraps"
);
}

#[test]
fn test_azure_storage_path_parse() {
let test_cases = vec![
Expand Down
4 changes: 1 addition & 3 deletions crates/storage/opendal/src/fs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,5 @@ pub(crate) fn fs_config_build() -> Result<Operator> {
let mut cfg = FsConfig::default();
cfg.root = Some("/".to_string());

Ok(Operator::from_config(cfg)
.map_err(from_opendal_error)?
.finish())
Operator::from_config(cfg).map_err(from_opendal_error)
}
2 changes: 1 addition & 1 deletion crates/storage/opendal/src/gcs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -131,7 +131,7 @@ pub(crate) fn gcs_config_build(
}
};

Ok(Operator::new(builder).map_err(from_opendal_error)?.finish())
Operator::new(builder).map_err(from_opendal_error)
}

/// Custom GCS credential loader.
Expand Down
4 changes: 1 addition & 3 deletions crates/storage/opendal/src/hf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -197,9 +197,7 @@ pub(crate) fn hf_config_build<'a>(
.ok_or_else(|| Error::new(ErrorKind::DataInvalid, format!("Invalid hf url: {path}")))?;
let relative_path = &path[path.len() - parsed.path.len()..];

let op = Operator::from_config(hf_cfg)
.map_err(from_opendal_error)?
.finish();
let op = Operator::from_config(hf_cfg).map_err(from_opendal_error)?;
Ok((op, relative_path))
}

Expand Down
Loading
Loading