diff --git a/Cargo.lock b/Cargo.lock index 78810b6331..c0f38f36c1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4046,6 +4046,7 @@ dependencies = [ "opendal", "reqsign-aws-v4", "reqsign-core", + "reqsign-google", "reqwest 0.12.28", "serde", "tokio", diff --git a/crates/catalog/hms/src/catalog.rs b/crates/catalog/hms/src/catalog.rs index 0f15890c77..4cd4e76427 100644 --- a/crates/catalog/hms/src/catalog.rs +++ b/crates/catalog/hms/src/catalog.rs @@ -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}; @@ -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)?; @@ -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(), @@ -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()); diff --git a/crates/storage/opendal/Cargo.toml b/crates/storage/opendal/Cargo.toml index 9256c1948f..48f6b106f6 100644 --- a/crates/storage/opendal/Cargo.toml +++ b/crates/storage/opendal/Cargo.toml @@ -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"] @@ -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 } @@ -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"] } diff --git a/crates/storage/opendal/public-api.txt b/crates/storage/opendal/public-api.txt index fc1ed7cf78..39e438f931 100644 --- a/crates/storage/opendal/public-api.txt +++ b/crates/storage/opendal/public-api.txt @@ -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 pub iceberg_storage_opendal::OpenDalStorage::Gcs pub iceberg_storage_opendal::OpenDalStorage::Gcs::config: alloc::sync::Arc +pub iceberg_storage_opendal::OpenDalStorage::Gcs::customized_credential_load: core::option::Option pub iceberg_storage_opendal::OpenDalStorage::Hf pub iceberg_storage_opendal::OpenDalStorage::Hf::config: alloc::sync::Arc pub iceberg_storage_opendal::OpenDalStorage::LocalFs @@ -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 pub iceberg_storage_opendal::OpenDalStorageFactory::Hf pub iceberg_storage_opendal::OpenDalStorageFactory::Memory pub iceberg_storage_opendal::OpenDalStorageFactory::Oss @@ -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 + '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 @@ -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 diff --git a/crates/storage/opendal/src/gcs.rs b/crates/storage/opendal/src/gcs.rs index 47ca52ccc6..2fc4e8a528 100644 --- a/crates/storage/opendal/src/gcs.rs +++ b/crates/storage/opendal/src/gcs.rs @@ -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}; @@ -70,8 +76,30 @@ pub(crate) fn gcs_config_parse(mut m: HashMap) -> Result Result { +pub(crate) fn gcs_config_build( + cfg: &GcsConfig, + customized_credential_load: &Option, + path: &str, +) -> Result { let url = Url::parse(path)?; let bucket = url.host_str().ok_or_else(|| { Error::new( @@ -82,7 +110,89 @@ pub(crate) fn gcs_config_build(cfg: &GcsConfig, path: &str) -> Result 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>); + +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 + 'static) -> Self { + Self(Arc::new(provider) as Arc>) + } +} + +#[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> { + 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}" + ); + } } diff --git a/crates/storage/opendal/src/lib.rs b/crates/storage/opendal/src/lib.rs index 52f4e68ed3..d69868befb 100644 --- a/crates/storage/opendal/src/lib.rs +++ b/crates/storage/opendal/src/lib.rs @@ -69,7 +69,7 @@ cfg_if! { cfg_if! { if #[cfg(feature = "opendal-gcs")] { mod gcs; - use gcs::*; + pub use gcs::*; use opendal::services::GcsConfig; } } @@ -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}; @@ -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, + }, /// OSS storage factory. #[cfg(feature = "opendal-oss")] Oss, @@ -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 { @@ -216,6 +231,9 @@ pub enum OpenDalStorage { Gcs { /// GCS configuration. config: Arc, + /// Custom GCS credential loader. + #[serde(skip)] + customized_credential_load: Option, }, /// OSS storage variant. #[cfg(feature = "opendal-oss")] @@ -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()..]) @@ -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!( @@ -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!( diff --git a/crates/storage/opendal/src/resolving.rs b/crates/storage/opendal/src/resolving.rs index 86993220a8..db47fe0cf8 100644 --- a/crates/storage/opendal/src/resolving.rs +++ b/crates/storage/opendal/src/resolving.rs @@ -34,6 +34,8 @@ use serde::{Deserialize, Serialize}; use url::Url; use crate::OpenDalStorage; +#[cfg(feature = "opendal-gcs")] +use crate::gcs::CustomGcsCredentialLoader; #[cfg(feature = "opendal-s3")] use crate::s3::CustomAwsCredentialLoader; @@ -80,58 +82,16 @@ fn extract_scheme(path: &str) -> Result<&'static str> { parse_scheme(url.scheme()) } -/// Build an [`OpenDalStorage`] variant for the given scheme and config properties. -fn build_storage_for_scheme( - scheme: &'static str, - props: &HashMap, - #[cfg(feature = "opendal-s3")] customized_credential_load: &Option, -) -> Result { - match scheme { - #[cfg(feature = "opendal-s3")] - "s3" => { - let config = crate::s3::s3_config_parse(props.clone())?; - Ok(OpenDalStorage::S3 { - config: Arc::new(config), - customized_credential_load: customized_credential_load.clone(), - }) - } - #[cfg(feature = "opendal-gcs")] - "gcs" => { - let config = crate::gcs::gcs_config_parse(props.clone())?; - Ok(OpenDalStorage::Gcs { - config: Arc::new(config), - }) - } - #[cfg(feature = "opendal-oss")] - "oss" => { - let config = crate::oss::oss_config_parse(props.clone())?; - Ok(OpenDalStorage::Oss { - config: Arc::new(config), - }) - } - #[cfg(feature = "opendal-azdls")] - "azdls" => { - let config = crate::azdls::azdls_config_parse(props.clone())?; - Ok(OpenDalStorage::Azdls { - config: Arc::new(config), - }) - } - #[cfg(feature = "opendal-fs")] - "file" => Ok(OpenDalStorage::LocalFs), - #[cfg(feature = "opendal-memory")] - "memory" => Ok(OpenDalStorage::Memory(crate::memory::memory_config_build()?)), - #[cfg(feature = "opendal-hf")] - "hf" => { - let config = crate::hf::hf_config_parse(props.clone())?; - Ok(OpenDalStorage::Hf { - config: Arc::new(config), - }) - } - unsupported => Err(Error::new( - ErrorKind::FeatureUnsupported, - format!("Unsupported storage scheme: {unsupported}"), - )), - } +/// The custom credential loaders a resolving storage carries, one per service that supports one. +/// +/// Bundled so that a service does not have to thread another `#[cfg]`-gated field through the +/// factory, the storage, and both of their constructors. +#[derive(Clone, Debug, Default)] +struct CustomCredentialLoaders { + #[cfg(feature = "opendal-s3")] + s3: Option, + #[cfg(feature = "opendal-gcs")] + gcs: Option, } /// A resolving storage factory that creates [`OpenDalResolvingStorage`] instances. @@ -154,10 +114,8 @@ fn build_storage_for_scheme( /// ``` #[derive(Clone, Debug, Serialize, Deserialize)] pub struct OpenDalResolvingStorageFactory { - /// Custom AWS credential loader for S3 storage. - #[cfg(feature = "opendal-s3")] #[serde(skip)] - customized_credential_load: Option, + loaders: CustomCredentialLoaders, } impl Default for OpenDalResolvingStorageFactory { @@ -170,15 +128,21 @@ impl OpenDalResolvingStorageFactory { /// Create a new resolving storage factory. pub fn new() -> Self { Self { - #[cfg(feature = "opendal-s3")] - customized_credential_load: None, + loaders: CustomCredentialLoaders::default(), } } /// Set a custom AWS credential loader for S3 storage. #[cfg(feature = "opendal-s3")] pub fn with_s3_credential_loader(mut self, loader: CustomAwsCredentialLoader) -> Self { - self.customized_credential_load = Some(loader); + self.loaders.s3 = Some(loader); + self + } + + /// Set a custom GCS credential loader for GCS storage. + #[cfg(feature = "opendal-gcs")] + pub fn with_gcs_credential_loader(mut self, loader: CustomGcsCredentialLoader) -> Self { + self.loaders.gcs = Some(loader); self } } @@ -189,8 +153,7 @@ impl StorageFactory for OpenDalResolvingStorageFactory { Ok(Arc::new(OpenDalResolvingStorage { props: config.props().clone(), storages: RwLock::new(HashMap::new()), - #[cfg(feature = "opendal-s3")] - customized_credential_load: self.customized_credential_load.clone(), + loaders: self.loaders.clone(), })) } } @@ -208,13 +171,65 @@ pub struct OpenDalResolvingStorage { /// Cache of canonical scheme to storage mappings. #[serde(skip, default)] storages: RwLock>>, - /// Custom AWS credential loader for S3 storage. - #[cfg(feature = "opendal-s3")] + // Every field of the bundle is feature-gated, so with no loader-capable service enabled it + // is empty and nothing below reads it. + #[allow(dead_code)] #[serde(skip)] - customized_credential_load: Option, + loaders: CustomCredentialLoaders, } impl OpenDalResolvingStorage { + /// Build an [`OpenDalStorage`] variant for `scheme` out of this storage's props. + fn build_storage_for_scheme(&self, scheme: &'static str) -> Result { + match scheme { + #[cfg(feature = "opendal-s3")] + "s3" => { + let config = crate::s3::s3_config_parse(self.props.clone())?; + Ok(OpenDalStorage::S3 { + config: Arc::new(config), + customized_credential_load: self.loaders.s3.clone(), + }) + } + #[cfg(feature = "opendal-gcs")] + "gcs" => { + let config = crate::gcs::gcs_config_parse(self.props.clone())?; + Ok(OpenDalStorage::Gcs { + config: Arc::new(config), + customized_credential_load: self.loaders.gcs.clone(), + }) + } + #[cfg(feature = "opendal-oss")] + "oss" => { + let config = crate::oss::oss_config_parse(self.props.clone())?; + Ok(OpenDalStorage::Oss { + config: Arc::new(config), + }) + } + #[cfg(feature = "opendal-azdls")] + "azdls" => { + let config = crate::azdls::azdls_config_parse(self.props.clone())?; + Ok(OpenDalStorage::Azdls { + config: Arc::new(config), + }) + } + #[cfg(feature = "opendal-fs")] + "file" => Ok(OpenDalStorage::LocalFs), + #[cfg(feature = "opendal-memory")] + "memory" => Ok(OpenDalStorage::Memory(crate::memory::memory_config_build()?)), + #[cfg(feature = "opendal-hf")] + "hf" => { + let config = crate::hf::hf_config_parse(self.props.clone())?; + Ok(OpenDalStorage::Hf { + config: Arc::new(config), + }) + } + unsupported => Err(Error::new( + ErrorKind::FeatureUnsupported, + format!("Unsupported storage scheme: {unsupported}"), + )), + } + } + /// Resolve the storage for the given path by extracting the canonical scheme and /// returning the cached or newly-created [`OpenDalStorage`]. fn resolve(&self, path: &str) -> Result> { @@ -242,13 +257,7 @@ impl OpenDalResolvingStorage { return Ok(storage.clone()); } - let storage = build_storage_for_scheme( - scheme, - &self.props, - #[cfg(feature = "opendal-s3")] - &self.customized_credential_load, - )?; - let storage = Arc::new(storage); + let storage = Arc::new(self.build_storage_for_scheme(scheme)?); cache.insert(scheme, storage.clone()); Ok(storage) } @@ -332,8 +341,7 @@ mod tests { OpenDalResolvingStorage { props: HashMap::new(), storages: RwLock::new(HashMap::new()), - #[cfg(feature = "opendal-s3")] - customized_credential_load: None, + loaders: CustomCredentialLoaders::default(), } } diff --git a/crates/storage/opendal/src/s3.rs b/crates/storage/opendal/src/s3.rs index 4b3893b39c..cbb450e999 100644 --- a/crates/storage/opendal/src/s3.rs +++ b/crates/storage/opendal/src/s3.rs @@ -29,9 +29,7 @@ use opendal::services::S3Config; use opendal::{Configurator, Operator}; /// AWS credentials: access key ID, secret access key, and optional session token. pub use reqsign_aws_v4::Credential as AwsCredential; -/// Trait for types that can asynchronously supply [`AwsCredential`] to a [`CustomAwsCredentialLoader`]. -pub use reqsign_core::ProvideCredential; -use reqsign_core::{ProvideCredentialChain, ProvideCredentialDyn}; +use reqsign_core::{ProvideCredential, ProvideCredentialChain, ProvideCredentialDyn}; use url::Url; use crate::utils::{from_opendal_error, is_truthy}; diff --git a/crates/storage/opendal/tests/file_io_gcs_test.rs b/crates/storage/opendal/tests/file_io_gcs_test.rs index 5e04491131..c1cbeb8383 100644 --- a/crates/storage/opendal/tests/file_io_gcs_test.rs +++ b/crates/storage/opendal/tests/file_io_gcs_test.rs @@ -25,9 +25,18 @@ mod tests { use std::sync::Arc; use bytes::Bytes; - use iceberg::io::{FileIO, FileIOBuilder, GCS_NO_AUTH, GCS_SERVICE_PATH}; - use iceberg_storage_opendal::OpenDalStorageFactory; + use iceberg::io::{FileIO, FileIOBuilder, GCS_NO_AUTH, GCS_SERVICE_PATH, GCS_TOKEN}; + use iceberg_storage_opendal::{ + CustomGcsCredentialLoader, GcsCredential, GcsToken, OpenDalStorageFactory, + ProvideCredential, + }; use iceberg_test_utils::{get_gcs_endpoint, set_up}; + use opendal::services::GcsConfig; + use opendal::{Configurator, Operator}; + use reqsign_core::{Context, ProvideCredentialChain}; + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + use tokio::net::TcpListener; + use tokio::sync::oneshot; static FAKE_GCS_BUCKET: &str = "test-bucket"; @@ -39,12 +48,14 @@ mod tests { // A bucket must exist for FileIO create_bucket(FAKE_GCS_BUCKET, &gcs_endpoint).await.unwrap(); - FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Gcs)) - .with_props(vec![ - (GCS_SERVICE_PATH, gcs_endpoint), - (GCS_NO_AUTH, "true".to_string()), - ]) - .build() + FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Gcs { + customized_credential_load: None, + })) + .with_props(vec![ + (GCS_SERVICE_PATH, gcs_endpoint), + (GCS_NO_AUTH, "true".to_string()), + ]) + .build() } // Create a bucket against the emulated GCS storage server. @@ -62,6 +73,187 @@ mod tests { format!("gs://{FAKE_GCS_BUCKET}") } + /// A loader that hands back whatever it was constructed with, standing in for a catalog that + /// vends a credential or fails to. + #[derive(Debug)] + struct MockCredentialLoader(Option); + + impl ProvideCredential for MockCredentialLoader { + type Credential = GcsCredential; + + async fn provide_credential( + &self, + _ctx: &Context, + ) -> reqsign_core::Result> { + Ok(self.0.clone()) + } + } + + /// Serves one request with a 404 and reports the `Authorization` header it saw, or `None` if + /// the header was absent. + /// + /// fake-gcs-server ignores credentials outright, so it cannot witness *which* token signed a + /// request. This stands in for it wherever that is the property under test. Returns the + /// endpoint to point `gcs.service.path` at. + async fn serve_one_request() -> (String, oneshot::Receiver>) { + let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind"); + let endpoint = format!("http://{}", listener.local_addr().expect("local addr")); + let (tx, rx) = oneshot::channel(); + + tokio::spawn(async move { + let (mut socket, _) = listener.accept().await.expect("accept"); + + // A GCS metadata read carries no body, so the header block is the whole request. + let mut request = Vec::new(); + let mut buf = [0u8; 1024]; + while !request.windows(4).any(|w| w == b"\r\n\r\n") { + match socket.read(&mut buf).await.expect("read") { + 0 => break, + n => request.extend_from_slice(&buf[..n]), + } + } + socket + .write_all(b"HTTP/1.1 404 Not Found\r\ncontent-length: 0\r\n\r\n") + .await + .expect("write response"); + + let request = String::from_utf8_lossy(&request); + let authorization = request + .lines() + .filter_map(|line| line.split_once(':')) + .find(|(name, _)| name.eq_ignore_ascii_case("authorization")) + .map(|(_, value)| value.trim().to_string()); + let _ = tx.send(authorization); + }); + + (endpoint, rx) + } + + /// Builds a `FileIO` whose only credential source should be `loader`, while also passing the + /// static `gcs.oauth2.token` prop that a vending catalog would supply. That prop is the + /// credential a loader exists to displace, so leaving it set is what makes the assertions + /// below meaningful. + fn file_io_with_loader(endpoint: String, loader: MockCredentialLoader) -> FileIO { + FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Gcs { + customized_credential_load: Some(CustomGcsCredentialLoader::new(loader)), + })) + .with_props(vec![ + (GCS_SERVICE_PATH, endpoint), + (GCS_TOKEN, "stale-baked-in-token".to_string()), + ]) + .build() + } + + #[tokio::test] + async fn test_gcs_custom_credential_loader_signs_with_its_token() { + let (endpoint, authorization) = serve_one_request().await; + let file_io = file_io_with_loader( + endpoint, + MockCredentialLoader(Some(GcsCredential::with_token(GcsToken { + access_token: "fresh-token".to_string(), + expires_at: None, + }))), + ); + + assert!( + !file_io + .exists(format!("{}/any", get_gs_path())) + .await + .expect("the loader's credential signs successfully"), + "the stub server answers 404" + ); + + // The loader's token reached the wire, and the static prop did not. + assert_eq!( + authorization.await.expect("the stub server saw a request"), + Some("Bearer fresh-token".to_string()) + ); + } + + #[tokio::test] + async fn test_gcs_custom_credential_loader_has_no_fallback() { + // OpenDAL's GCS service *prepends* a custom chain to its own rather than replacing it, + // and the chain swallows a provider's miss. So a loader that comes up empty must still + // fail the request: falling through to the static `gcs.oauth2.token` prop, which is + // exactly the credential the loader exists to refresh, would leave an expired token in + // force with no diagnostic. + let (endpoint, mut authorization) = serve_one_request().await; + let file_io = file_io_with_loader(endpoint, MockCredentialLoader(None)); + + let err = file_io + .exists(format!("{}/any", get_gs_path())) + .await + .expect_err("a loader that vends nothing must not fall back"); + assert!( + err.to_string() + .contains("failed to load signing credential"), + "unexpected error: {err}" + ); + // Signing fails before any I/O, so the stub server was never reached at all. + assert_eq!( + authorization.try_recv(), + Err(oneshot::error::TryRecvError::Empty) + ); + } + + /// The control for [`test_gcs_custom_credential_loader_has_no_fallback`]: the same empty + /// loader and the same stale config token, but skipping the suppression that + /// `gcs_config_build` applies by manually constructing the credential chain. Here the request + /// succeeds, signed with the credential the loader was supposed to displace. + #[tokio::test] + async fn test_opendal_gcs_chain_falls_through_to_the_config_token() { + let (endpoint, authorization) = serve_one_request().await; + + // Disabling the ambient sources leaves the chain with exactly two entries that can + // produce anything: the empty loader, then the token. + let mut cfg = GcsConfig::default(); + cfg.bucket = FAKE_GCS_BUCKET.to_string(); + cfg.endpoint = Some(endpoint); + cfg.token = Some("stale-baked-in-token".to_string()); + cfg.disable_vm_metadata = true; + cfg.disable_config_load = true; + + let chain = ProvideCredentialChain::new().push(MockCredentialLoader(None)); + let operator = Operator::new(cfg.into_builder().credential_provider_chain(chain)) + .expect("operator builds") + .finish(); + + assert!( + !operator.exists("any").await.expect("the request is signed"), + "the stub server answers 404" + ); + assert_eq!( + authorization.await.expect("the stub server saw a request"), + Some("Bearer stale-baked-in-token".to_string()) + ); + } + + #[tokio::test] + async fn test_gcs_custom_credential_loader_rejects_unsigned_requests() { + set_up(); + + // `gcs.no-auth` skips signing entirely, which would silently sideline the loader. + let file_io = FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Gcs { + customized_credential_load: Some(CustomGcsCredentialLoader::new(MockCredentialLoader( + None, + ))), + })) + .with_props(vec![ + (GCS_SERVICE_PATH, get_gcs_endpoint()), + (GCS_NO_AUTH, "true".to_string()), + ]) + .build(); + + let err = file_io + .exists(format!("{}/any", get_gs_path())) + .await + .expect_err("a loader combined with gcs.no-auth is a configuration error"); + assert!( + err.to_string().contains("disable request signing"), + "unexpected error: {err}" + ); + } + #[tokio::test] async fn gcs_exists() { let file_io = get_file_io_gcs().await; diff --git a/dev/hms/Dockerfile b/dev/hms/Dockerfile index 65ebc3d6d5..5e6cfdbfef 100644 --- a/dev/hms/Dockerfile +++ b/dev/hms/Dockerfile @@ -13,18 +13,16 @@ # See the License for the specific language governing permissions and # limitations under the License. -FROM apache/hive:3.1.3 - -ENV AWSSDK_VERSION=2.20.18 -ENV HADOOP_VERSION=3.1.0 +FROM apache/hive:4.2.1 USER root -RUN apt-get update -qq && apt-get -qq -y install curl && \ - curl https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/${HADOOP_VERSION}/hadoop-aws-${HADOOP_VERSION}.jar -Lo /opt/hive/lib/hadoop-aws-${HADOOP_VERSION}.jar && \ - curl https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.11.271/aws-java-sdk-bundle-1.11.271.jar -Lo /opt/hive/lib/aws-java-sdk-bundle-1.11.271.jar && \ - apt-get clean && rm -rf /var/lib/apt/lists/* +# Link the hadoop-aws and AWS SDK jars that the image ships into the metastore classpath +RUN ln -s /opt/hadoop/share/hadoop/tools/lib/hadoop-aws-*.jar /opt/hive/lib/ && \ + ln -s /opt/hadoop/share/hadoop/tools/lib/bundle-*.jar /opt/hive/lib/ -COPY core-site.xml /opt/hadoop/etc/hadoop/core-site.xml +# The entrypoint links this directory into the Hive config directory, over its own core-site.xml +ENV HIVE_CUSTOM_CONF_DIR=/opt/hive/custom-conf +COPY core-site.xml ${HIVE_CUSTOM_CONF_DIR}/core-site.xml USER hive