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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

32 changes: 24 additions & 8 deletions crates/catalog/hms/src/catalog.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,8 @@ use std::sync::Arc;
use anyhow::anyhow;
use async_trait::async_trait;
use hive_metastore::{
ThriftHiveMetastoreClient, ThriftHiveMetastoreClientBuilder,
ThriftHiveMetastoreGetDatabaseException, ThriftHiveMetastoreGetTableException,
GetTableRequest, ThriftHiveMetastoreClient, ThriftHiveMetastoreClientBuilder,
ThriftHiveMetastoreGetDatabaseException, ThriftHiveMetastoreGetTableReqException,
};
use iceberg::io::{FileIO, FileIOBuilder, StorageFactory};
use iceberg::spec::{TableMetadata, TableMetadataBuilder};
Expand Down Expand Up @@ -566,10 +566,15 @@ impl Catalog for HmsCatalog {
let hive_table = self
.client
.0
.get_table(db_name.clone().into(), table.name.clone().into())
.get_table_req(GetTableRequest {
db_name: db_name.clone().into(),
tbl_name: table.name.clone().into(),
..Default::default()
})
.await
.map(from_thrift_exception)
.map_err(from_thrift_error)??;
.map_err(from_thrift_error)??
.table;

let metadata_location = get_metadata_location(&hive_table.parameters)?;

Expand Down Expand Up @@ -641,12 +646,18 @@ impl Catalog for HmsCatalog {
let resp = self
.client
.0
.get_table(db_name.into(), table_name.into())
.get_table_req(GetTableRequest {
db_name: db_name.into(),
tbl_name: table_name.into(),
..Default::default()
})
.await;

match resp {
Ok(MaybeException::Ok(_)) => Ok(true),
Ok(MaybeException::Exception(ThriftHiveMetastoreGetTableException::O2(_))) => Ok(false),
Ok(MaybeException::Exception(ThriftHiveMetastoreGetTableReqException::O2(_))) => {
Ok(false)
}
Ok(MaybeException::Exception(exception)) => Err(Error::new(
ErrorKind::Unexpected,
"Operation failed for hitting thrift error".to_string(),
Expand Down Expand Up @@ -678,10 +689,15 @@ impl Catalog for HmsCatalog {
let mut tbl = self
.client
.0
.get_table(src_dbname.clone().into(), src_tbl_name.clone().into())
.get_table_req(GetTableRequest {
db_name: src_dbname.clone().into(),
tbl_name: src_tbl_name.clone().into(),
..Default::default()
})
.await
.map(from_thrift_exception)
.map_err(from_thrift_error)??;
.map_err(from_thrift_error)??
.table;

tbl.db_name = Some(dest_dbname.into());
tbl.table_name = Some(dest_tbl_name.into());
Expand Down
5 changes: 3 additions & 2 deletions crates/storage/opendal/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ opendal-all = [

opendal-azdls = ["opendal/services-azdls"]
opendal-fs = ["opendal/services-fs"]
opendal-gcs = ["opendal/services-gcs"]
opendal-gcs = ["opendal/services-gcs", "reqsign-google", "reqsign-core"]
opendal-hf = ["opendal/services-hf"]
opendal-memory = ["opendal/services-memory"]
opendal-oss = ["opendal/services-oss"]
Expand All @@ -57,6 +57,7 @@ iceberg = { workspace = true }
opendal = { workspace = true }
reqsign-aws-v4 = { version = "3.0.0", optional = true }
reqsign-core = { version = "3.0.0", optional = true }
reqsign-google = { version = "3.0.0", optional = true }
serde = { workspace = true }
typetag = { workspace = true }
url = { workspace = true }
Expand All @@ -65,4 +66,4 @@ url = { workspace = true }
async-trait = { workspace = true }
iceberg_test_utils = { path = "../../test_utils", features = ["tests"] }
reqwest = { workspace = true }
tokio = { workspace = true, features = ["macros"] }
tokio = { workspace = true, features = ["io-util", "macros", "net"] }
12 changes: 12 additions & 0 deletions crates/storage/opendal/public-api.txt
Original file line number Diff line number Diff line change
@@ -1,11 +1,14 @@
pub mod iceberg_storage_opendal
pub use iceberg_storage_opendal::AwsCredential
pub use iceberg_storage_opendal::GcsCredential
pub use iceberg_storage_opendal::GcsToken
pub use iceberg_storage_opendal::ProvideCredential
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::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>
pub iceberg_storage_opendal::OpenDalStorage::Hf
pub iceberg_storage_opendal::OpenDalStorage::Hf::config: alloc::sync::Arc<opendal_service_hf::config::HfConfig>
pub iceberg_storage_opendal::OpenDalStorage::LocalFs
Expand Down Expand Up @@ -39,6 +42,7 @@ pub enum iceberg_storage_opendal::OpenDalStorageFactory
pub iceberg_storage_opendal::OpenDalStorageFactory::Azdls
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>
pub iceberg_storage_opendal::OpenDalStorageFactory::Hf
pub iceberg_storage_opendal::OpenDalStorageFactory::Memory
pub iceberg_storage_opendal::OpenDalStorageFactory::Oss
Expand All @@ -61,6 +65,13 @@ 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::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
impl core::clone::Clone for iceberg_storage_opendal::CustomGcsCredentialLoader
pub fn iceberg_storage_opendal::CustomGcsCredentialLoader::clone(&self) -> Self
impl core::fmt::Debug for iceberg_storage_opendal::CustomGcsCredentialLoader
pub fn iceberg_storage_opendal::CustomGcsCredentialLoader::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
pub struct iceberg_storage_opendal::OpenDalResolvingStorage
impl core::fmt::Debug for iceberg_storage_opendal::OpenDalResolvingStorage
pub fn iceberg_storage_opendal::OpenDalResolvingStorage::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
Expand All @@ -83,6 +94,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_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
pub fn iceberg_storage_opendal::OpenDalResolvingStorageFactory::clone(&self) -> iceberg_storage_opendal::OpenDalResolvingStorageFactory
Expand Down
120 changes: 115 additions & 5 deletions crates/storage/opendal/src/gcs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,14 +17,20 @@
//! Google Cloud Storage properties

use std::collections::HashMap;
use std::sync::Arc;

use iceberg::io::{
GCS_ALLOW_ANONYMOUS, GCS_CREDENTIALS_JSON, GCS_DISABLE_CONFIG_LOAD, GCS_DISABLE_VM_METADATA,
GCS_NO_AUTH, GCS_SERVICE_PATH, GCS_TOKEN,
};
use iceberg::{Error, ErrorKind, Result};
use opendal::Operator;
use opendal::services::GcsConfig;
use opendal::{Configurator, Operator};
use reqsign_core::{ProvideCredential, ProvideCredentialChain, ProvideCredentialDyn};
/// GCS credentials: either a service account or an OAuth2 access token.
pub use reqsign_google::Credential as GcsCredential;
/// An OAuth2 access token, the form a catalog-vended GCS credential takes.
pub use reqsign_google::Token as GcsToken;
use url::Url;

use crate::utils::{from_opendal_error, is_truthy};
Expand Down Expand Up @@ -70,8 +76,30 @@ pub(crate) fn gcs_config_parse(mut m: HashMap<String, String>) -> Result<GcsConf
Ok(cfg)
}

/// Clear every credential source opendal would consult on its own.
///
/// opendal's GCS service *prepends* a caller-supplied credential chain to its own
/// (unlike its S3 service, which replaces it), and [`ProvideCredentialChain`] swallows a
/// provider's error and falls through to the next entry. So without this, a custom loader
/// that fails is indistinguishable from one that was never installed: signing quietly
/// succeeds using whatever else the chain can reach, including the very static
/// `gcs.oauth2.token` the loader exists to replace, or ambient credentials on a GCE
/// instance.
fn suppress_default_credential_sources(cfg: &mut GcsConfig) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is there a reason this gets done as part of building the operator instead of during config initialization (e.g. crate::gcs::gcs_config_parse)?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

so this is because config initialization is not actually a required step, it happens in the OpenDalStorageFactory::build function, and one could technically skip that and the whole OpenDalStorageFactory step, go straight to building an OpenDalStorage::Gcs{ ... } themselves, meaning this step gets skipped.

Arguably that's fine, since we're realistically the only ones using this interface and this is the only location where it happens, but this way just ensures that the suppression actually happens. Maybe even we should make the default credential suppression optional so that customers can decide to use it as a fallback, but I'm okay with saying that if they specify vended-credentials in their catalog connection, then we WILL use vended creds, and error if the creds don't work for some reason.

cfg.token = None;
cfg.credential = None;
cfg.credential_path = None;
cfg.service_account = None;
cfg.disable_vm_metadata = true;
cfg.disable_config_load = true;
}

/// Build a new OpenDAL [`Operator`] based on a provided [`GcsConfig`].
pub(crate) fn gcs_config_build(cfg: &GcsConfig, path: &str) -> Result<Operator> {
pub(crate) fn gcs_config_build(
cfg: &GcsConfig,
customized_credential_load: &Option<CustomGcsCredentialLoader>,
path: &str,
) -> Result<Operator> {
let url = Url::parse(path)?;
let bucket = url.host_str().ok_or_else(|| {
Error::new(
Expand All @@ -82,7 +110,89 @@ pub(crate) fn gcs_config_build(cfg: &GcsConfig, path: &str) -> Result<Operator>

let mut cfg = cfg.clone();
cfg.bucket = bucket.to_string();
Ok(Operator::from_config(cfg)
.map_err(from_opendal_error)?
.finish())

let builder = match customized_credential_load {
None => cfg.into_builder(),
Some(loader) => {
// `skip_signature` bypasses the signer entirely, so the loader would never be
// consulted. Refuse rather than silently ignore one of the two.
if cfg.skip_signature {
return Err(Error::new(
ErrorKind::DataInvalid,
format!(
"a custom GCS credential loader cannot be combined with \
{GCS_NO_AUTH} or {GCS_ALLOW_ANONYMOUS}, which disable request signing"
),
));
}
suppress_default_credential_sources(&mut cfg);
let chain = ProvideCredentialChain::new().push(Arc::clone(&loader.0));
cfg.into_builder().credential_provider_chain(chain)
}
};

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

/// Custom GCS credential loader.
///
/// Wraps any [`ProvideCredential`] implementation for use with the GCS storage backend.
/// Use [`CustomGcsCredentialLoader::new`] to create one, then pass it to
/// [`OpenDalStorageFactory::Gcs`](crate::OpenDalStorageFactory).
///
/// Installing a loader suppresses every credential source opendal would otherwise reach
/// for, so the loader is the sole authority on credentials. It must therefore succeed on
/// its own: there is no application default fallback behind it.
pub struct CustomGcsCredentialLoader(Arc<dyn ProvideCredentialDyn<Credential = GcsCredential>>);

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

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

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

#[cfg(test)]
mod tests {
use super::*;

#[derive(Debug)]
struct NoopProvider;

impl ProvideCredential for NoopProvider {
type Credential = GcsCredential;

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

#[test]
fn loader_rejects_unsigned_requests() {
let mut cfg = GcsConfig::default();
cfg.skip_signature = true;
let loader = Some(CustomGcsCredentialLoader::new(NoopProvider));

let err = gcs_config_build(&cfg, &loader, "gs://bucket/key").unwrap_err();
assert!(
err.to_string().contains("disable request signing"),
"unexpected error: {err}"
);
}
}
33 changes: 28 additions & 5 deletions crates/storage/opendal/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ cfg_if! {
cfg_if! {
if #[cfg(feature = "opendal-gcs")] {
mod gcs;
use gcs::*;
pub use gcs::*;
use opendal::services::GcsConfig;
}
}
Expand Down Expand Up @@ -97,6 +97,14 @@ cfg_if! {
}
}

/// Trait for types that can asynchronously supply credentials to a custom credential
/// loader, such as [`CustomAwsCredentialLoader`] or [`CustomGcsCredentialLoader`].
///
/// Downstream implementors must name the trait through this re-export: a loader built
/// against a differently-versioned `reqsign-core` will not satisfy the bound.
#[cfg(any(feature = "opendal-s3", feature = "opendal-gcs"))]
pub use reqsign_core::ProvideCredential;

mod resolving;
pub use resolving::{OpenDalResolvingStorage, OpenDalResolvingStorageFactory};

Expand All @@ -121,7 +129,11 @@ pub enum OpenDalStorageFactory {
},
/// GCS storage factory.
#[cfg(feature = "opendal-gcs")]
Gcs,
Gcs {
/// Custom GCS credential loader.
#[serde(skip)]
customized_credential_load: Option<gcs::CustomGcsCredentialLoader>,
},
/// OSS storage factory.
#[cfg(feature = "opendal-oss")]
Oss,
Expand Down Expand Up @@ -152,8 +164,11 @@ impl StorageFactory for OpenDalStorageFactory {
customized_credential_load: customized_credential_load.clone(),
})),
#[cfg(feature = "opendal-gcs")]
OpenDalStorageFactory::Gcs => Ok(Arc::new(OpenDalStorage::Gcs {
OpenDalStorageFactory::Gcs {
customized_credential_load,
} => Ok(Arc::new(OpenDalStorage::Gcs {
config: gcs_config_parse(config.props().clone())?.into(),
customized_credential_load: customized_credential_load.clone(),
})),
#[cfg(feature = "opendal-oss")]
OpenDalStorageFactory::Oss => Ok(Arc::new(OpenDalStorage::Oss {
Expand Down Expand Up @@ -216,6 +231,9 @@ pub enum OpenDalStorage {
Gcs {
/// GCS configuration.
config: Arc<GcsConfig>,
/// Custom GCS credential loader.
#[serde(skip)]
customized_credential_load: Option<gcs::CustomGcsCredentialLoader>,
},
/// OSS storage variant.
#[cfg(feature = "opendal-oss")]
Expand Down Expand Up @@ -310,8 +328,11 @@ impl OpenDalStorage {
}
}
#[cfg(feature = "opendal-gcs")]
OpenDalStorage::Gcs { config } => {
let operator = gcs_config_build(config, path)?;
OpenDalStorage::Gcs {
config,
customized_credential_load,
} => {
let operator = gcs_config_build(config, customized_credential_load, path)?;
let prefix = format!("gs://{}/", operator.info().name());
if path.starts_with(&prefix) {
(operator, &path[prefix.len()..])
Expand Down Expand Up @@ -697,6 +718,7 @@ mod tests {
fn test_relativize_path_gcs() {
let storage = OpenDalStorage::Gcs {
config: Arc::new(GcsConfig::default()),
customized_credential_load: None,
};

assert_eq!(
Expand All @@ -712,6 +734,7 @@ mod tests {
fn test_relativize_path_gcs_invalid_scheme() {
let storage = OpenDalStorage::Gcs {
config: Arc::new(GcsConfig::default()),
customized_credential_load: None,
};

assert!(
Expand Down
Loading
Loading