From 647d0c6c66da69769f29da7eae410e21386aeb0d Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Tue, 8 Sep 2026 14:33:37 -0400 Subject: [PATCH 1/4] implement custom credential loader into OpenDalStorage::Gcs --- Cargo.lock | 1 + crates/storage/opendal/Cargo.toml | 3 +- crates/storage/opendal/public-api.txt | 12 ++ crates/storage/opendal/src/gcs.rs | 142 +++++++++++++++++- crates/storage/opendal/src/lib.rs | 33 +++- crates/storage/opendal/src/resolving.rs | 45 ++++-- crates/storage/opendal/src/s3.rs | 4 +- .../storage/opendal/tests/file_io_gcs_test.rs | 118 ++++++++++++++- 8 files changed, 327 insertions(+), 31 deletions(-) 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/storage/opendal/Cargo.toml b/crates/storage/opendal/Cargo.toml index 9256c1948f..3c5d7d07f6 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 } 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..c59cb36a21 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,111 @@ 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()) + + if customized_credential_load.is_some() { + // `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 mut builder = cfg.into_builder(); + if let Some(loader) = customized_credential_load { + let chain = ProvideCredentialChain::new().push(Arc::clone(&loader.0)); + builder = 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::*; + + #[test] + fn suppress_clears_every_credential_source() { + let mut cfg = GcsConfig::default(); + cfg.token = Some("vended".to_string()); + cfg.credential = Some("{}".to_string()); + cfg.credential_path = Some("/var/run/creds.json".to_string()); + cfg.service_account = Some("sa@example.invalid".to_string()); + cfg.bucket = "bucket".to_string(); + + suppress_default_credential_sources(&mut cfg); + + assert_eq!(cfg.token, None); + assert_eq!(cfg.credential, None); + assert_eq!(cfg.credential_path, None); + assert_eq!(cfg.service_account, None); + assert!(cfg.disable_vm_metadata); + assert!(cfg.disable_config_load); + // Unrelated settings survive. + assert_eq!(cfg.bucket, "bucket"); + } + + #[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..bf871533fc 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; @@ -84,7 +86,8 @@ fn extract_scheme(path: &str) -> Result<&'static str> { fn build_storage_for_scheme( scheme: &'static str, props: &HashMap, - #[cfg(feature = "opendal-s3")] customized_credential_load: &Option, + #[cfg(feature = "opendal-s3")] s3_credential_load: &Option, + #[cfg(feature = "opendal-gcs")] gcs_credential_load: &Option, ) -> Result { match scheme { #[cfg(feature = "opendal-s3")] @@ -92,7 +95,7 @@ fn build_storage_for_scheme( let config = crate::s3::s3_config_parse(props.clone())?; Ok(OpenDalStorage::S3 { config: Arc::new(config), - customized_credential_load: customized_credential_load.clone(), + customized_credential_load: s3_credential_load.clone(), }) } #[cfg(feature = "opendal-gcs")] @@ -100,6 +103,7 @@ fn build_storage_for_scheme( let config = crate::gcs::gcs_config_parse(props.clone())?; Ok(OpenDalStorage::Gcs { config: Arc::new(config), + customized_credential_load: gcs_credential_load.clone(), }) } #[cfg(feature = "opendal-oss")] @@ -157,7 +161,11 @@ pub struct OpenDalResolvingStorageFactory { /// Custom AWS credential loader for S3 storage. #[cfg(feature = "opendal-s3")] #[serde(skip)] - customized_credential_load: Option, + s3_credential_load: Option, + /// Custom GCS credential loader for GCS storage. + #[cfg(feature = "opendal-gcs")] + #[serde(skip)] + gcs_credential_load: Option, } impl Default for OpenDalResolvingStorageFactory { @@ -171,14 +179,23 @@ impl OpenDalResolvingStorageFactory { pub fn new() -> Self { Self { #[cfg(feature = "opendal-s3")] - customized_credential_load: None, + s3_credential_load: None, + #[cfg(feature = "opendal-gcs")] + gcs_credential_load: None, } } /// 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.s3_credential_load = 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.gcs_credential_load = Some(loader); self } } @@ -190,7 +207,9 @@ impl StorageFactory for OpenDalResolvingStorageFactory { props: config.props().clone(), storages: RwLock::new(HashMap::new()), #[cfg(feature = "opendal-s3")] - customized_credential_load: self.customized_credential_load.clone(), + s3_credential_load: self.s3_credential_load.clone(), + #[cfg(feature = "opendal-gcs")] + gcs_credential_load: self.gcs_credential_load.clone(), })) } } @@ -211,7 +230,11 @@ pub struct OpenDalResolvingStorage { /// Custom AWS credential loader for S3 storage. #[cfg(feature = "opendal-s3")] #[serde(skip)] - customized_credential_load: Option, + s3_credential_load: Option, + /// Custom GCS credential loader for GCS storage. + #[cfg(feature = "opendal-gcs")] + #[serde(skip)] + gcs_credential_load: Option, } impl OpenDalResolvingStorage { @@ -246,7 +269,9 @@ impl OpenDalResolvingStorage { scheme, &self.props, #[cfg(feature = "opendal-s3")] - &self.customized_credential_load, + &self.s3_credential_load, + #[cfg(feature = "opendal-gcs")] + &self.gcs_credential_load, )?; let storage = Arc::new(storage); cache.insert(scheme, storage.clone()); @@ -333,7 +358,9 @@ mod tests { props: HashMap::new(), storages: RwLock::new(HashMap::new()), #[cfg(feature = "opendal-s3")] - customized_credential_load: None, + s3_credential_load: None, + #[cfg(feature = "opendal-gcs")] + gcs_credential_load: None, } } 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..eb338d547c 100644 --- a/crates/storage/opendal/tests/file_io_gcs_test.rs +++ b/crates/storage/opendal/tests/file_io_gcs_test.rs @@ -25,9 +25,13 @@ 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 reqsign_core::Context; static FAKE_GCS_BUCKET: &str = "test-bucket"; @@ -39,12 +43,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 +68,102 @@ 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()) + } + } + + /// 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. + async fn get_file_io_gcs_with_loader(loader: MockCredentialLoader) -> FileIO { + set_up(); + + let gcs_endpoint = get_gcs_endpoint(); + create_bucket(FAKE_GCS_BUCKET, &gcs_endpoint).await.unwrap(); + + FileIOBuilder::new(Arc::new(OpenDalStorageFactory::Gcs { + customized_credential_load: Some(CustomGcsCredentialLoader::new(loader)), + })) + .with_props(vec![ + (GCS_SERVICE_PATH, gcs_endpoint), + (GCS_TOKEN, "stale-baked-in-token".to_string()), + ]) + .build() + } + + #[tokio::test] + async fn test_gcs_with_custom_credential_loader() { + let file_io = get_file_io_gcs_with_loader(MockCredentialLoader(Some( + GcsCredential::with_token(GcsToken { + access_token: "fresh-token".to_string(), + expires_at: None, + }), + ))) + .await; + + file_io + .exists(format!("{}/any", get_gs_path())) + .await + .expect("a loader that vends a credential signs successfully"); + } + + #[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` above, which is + // exactly the credential the loader exists to refresh, would leave an expired token in + // force with no diagnostic. + let file_io = get_file_io_gcs_with_loader(MockCredentialLoader(None)).await; + + 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}" + ); + } + + #[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; From f789d90ff4b516331ba07afc2b9790b20fb3db2a Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Thu, 10 Sep 2026 12:48:28 -0400 Subject: [PATCH 2/4] address comments, updates tests, minor refactor --- crates/storage/opendal/Cargo.toml | 2 +- crates/storage/opendal/src/gcs.rs | 58 ++----- crates/storage/opendal/src/resolving.rs | 163 ++++++++---------- .../storage/opendal/tests/file_io_gcs_test.rs | 130 +++++++++++--- 4 files changed, 201 insertions(+), 152 deletions(-) diff --git a/crates/storage/opendal/Cargo.toml b/crates/storage/opendal/Cargo.toml index 3c5d7d07f6..48f6b106f6 100644 --- a/crates/storage/opendal/Cargo.toml +++ b/crates/storage/opendal/Cargo.toml @@ -66,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/src/gcs.rs b/crates/storage/opendal/src/gcs.rs index c59cb36a21..2fc4e8a528 100644 --- a/crates/storage/opendal/src/gcs.rs +++ b/crates/storage/opendal/src/gcs.rs @@ -111,26 +111,25 @@ pub(crate) fn gcs_config_build( let mut cfg = cfg.clone(); cfg.bucket = bucket.to_string(); - if customized_credential_load.is_some() { - // `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" - ), - )); + 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) } - suppress_default_credential_sources(&mut cfg); - } - - let mut builder = cfg.into_builder(); - if let Some(loader) = customized_credential_load { - let chain = ProvideCredentialChain::new().push(Arc::clone(&loader.0)); - builder = builder.credential_provider_chain(chain); - } + }; Ok(Operator::new(builder).map_err(from_opendal_error)?.finish()) } @@ -170,27 +169,6 @@ impl CustomGcsCredentialLoader { mod tests { use super::*; - #[test] - fn suppress_clears_every_credential_source() { - let mut cfg = GcsConfig::default(); - cfg.token = Some("vended".to_string()); - cfg.credential = Some("{}".to_string()); - cfg.credential_path = Some("/var/run/creds.json".to_string()); - cfg.service_account = Some("sa@example.invalid".to_string()); - cfg.bucket = "bucket".to_string(); - - suppress_default_credential_sources(&mut cfg); - - assert_eq!(cfg.token, None); - assert_eq!(cfg.credential, None); - assert_eq!(cfg.credential_path, None); - assert_eq!(cfg.service_account, None); - assert!(cfg.disable_vm_metadata); - assert!(cfg.disable_config_load); - // Unrelated settings survive. - assert_eq!(cfg.bucket, "bucket"); - } - #[derive(Debug)] struct NoopProvider; diff --git a/crates/storage/opendal/src/resolving.rs b/crates/storage/opendal/src/resolving.rs index bf871533fc..db47fe0cf8 100644 --- a/crates/storage/opendal/src/resolving.rs +++ b/crates/storage/opendal/src/resolving.rs @@ -82,60 +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")] s3_credential_load: &Option, - #[cfg(feature = "opendal-gcs")] gcs_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: s3_credential_load.clone(), - }) - } - #[cfg(feature = "opendal-gcs")] - "gcs" => { - let config = crate::gcs::gcs_config_parse(props.clone())?; - Ok(OpenDalStorage::Gcs { - config: Arc::new(config), - customized_credential_load: gcs_credential_load.clone(), - }) - } - #[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. @@ -158,14 +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)] - s3_credential_load: Option, - /// Custom GCS credential loader for GCS storage. - #[cfg(feature = "opendal-gcs")] - #[serde(skip)] - gcs_credential_load: Option, + loaders: CustomCredentialLoaders, } impl Default for OpenDalResolvingStorageFactory { @@ -178,24 +128,21 @@ impl OpenDalResolvingStorageFactory { /// Create a new resolving storage factory. pub fn new() -> Self { Self { - #[cfg(feature = "opendal-s3")] - s3_credential_load: None, - #[cfg(feature = "opendal-gcs")] - gcs_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.s3_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.gcs_credential_load = Some(loader); + self.loaders.gcs = Some(loader); self } } @@ -206,10 +153,7 @@ impl StorageFactory for OpenDalResolvingStorageFactory { Ok(Arc::new(OpenDalResolvingStorage { props: config.props().clone(), storages: RwLock::new(HashMap::new()), - #[cfg(feature = "opendal-s3")] - s3_credential_load: self.s3_credential_load.clone(), - #[cfg(feature = "opendal-gcs")] - gcs_credential_load: self.gcs_credential_load.clone(), + loaders: self.loaders.clone(), })) } } @@ -227,17 +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)] - s3_credential_load: Option, - /// Custom GCS credential loader for GCS storage. - #[cfg(feature = "opendal-gcs")] - #[serde(skip)] - gcs_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> { @@ -265,15 +257,7 @@ impl OpenDalResolvingStorage { return Ok(storage.clone()); } - let storage = build_storage_for_scheme( - scheme, - &self.props, - #[cfg(feature = "opendal-s3")] - &self.s3_credential_load, - #[cfg(feature = "opendal-gcs")] - &self.gcs_credential_load, - )?; - let storage = Arc::new(storage); + let storage = Arc::new(self.build_storage_for_scheme(scheme)?); cache.insert(scheme, storage.clone()); Ok(storage) } @@ -357,10 +341,7 @@ mod tests { OpenDalResolvingStorage { props: HashMap::new(), storages: RwLock::new(HashMap::new()), - #[cfg(feature = "opendal-s3")] - s3_credential_load: None, - #[cfg(feature = "opendal-gcs")] - gcs_credential_load: None, + loaders: CustomCredentialLoaders::default(), } } diff --git a/crates/storage/opendal/tests/file_io_gcs_test.rs b/crates/storage/opendal/tests/file_io_gcs_test.rs index eb338d547c..c1cbeb8383 100644 --- a/crates/storage/opendal/tests/file_io_gcs_test.rs +++ b/crates/storage/opendal/tests/file_io_gcs_test.rs @@ -31,7 +31,12 @@ mod tests { ProvideCredential, }; use iceberg_test_utils::{get_gcs_endpoint, set_up}; - use reqsign_core::Context; + 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"; @@ -84,48 +89,96 @@ mod tests { } } - /// 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. - async fn get_file_io_gcs_with_loader(loader: MockCredentialLoader) -> FileIO { - set_up(); + /// 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(); - let gcs_endpoint = get_gcs_endpoint(); - create_bucket(FAKE_GCS_BUCKET, &gcs_endpoint).await.unwrap(); + 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, gcs_endpoint), + (GCS_SERVICE_PATH, endpoint), (GCS_TOKEN, "stale-baked-in-token".to_string()), ]) .build() } #[tokio::test] - async fn test_gcs_with_custom_credential_loader() { - let file_io = get_file_io_gcs_with_loader(MockCredentialLoader(Some( - GcsCredential::with_token(GcsToken { + 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, - }), - ))) - .await; + }))), + ); - file_io - .exists(format!("{}/any", get_gs_path())) - .await - .expect("a loader that vends a credential signs successfully"); + 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` above, which is + // 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 file_io = get_file_io_gcs_with_loader(MockCredentialLoader(None)).await; + 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())) @@ -136,6 +189,43 @@ mod tests { .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] From 20ef0ce45ad3888390f740f05a4d682c8a7f8a4e Mon Sep 17 00:00:00 2001 From: Jannik Steinmann Date: Tue, 8 Sep 2026 17:00:52 +0200 Subject: [PATCH 3/4] fix: Fix CI failure due to integration test's Hive install (#3174) Install Hive directly Signed-off-by: Jannik Steinmann (cherry picked from commit 4687d265374bf5dc7bd8cd461f3bd46a23864761) --- dev/hms/Dockerfile | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/dev/hms/Dockerfile b/dev/hms/Dockerfile index 65ebc3d6d5..dff5606bdc 100644 --- a/dev/hms/Dockerfile +++ b/dev/hms/Dockerfile @@ -20,10 +20,12 @@ ENV HADOOP_VERSION=3.1.0 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/* +ADD --chmod=644 --checksum=sha256:a18508b9348af095ea41301e439354dbd449e304ac44c6885b2b4fe78de88126 \ + https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/${HADOOP_VERSION}/hadoop-aws-${HADOOP_VERSION}.jar \ + /opt/hive/lib/hadoop-aws-${HADOOP_VERSION}.jar +ADD --chmod=644 --checksum=sha256:faf78ac4880f56cf52791d84ec1068ce7c66acc4295d580a726104b734c01fcd \ + https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.11.271/aws-java-sdk-bundle-1.11.271.jar \ + /opt/hive/lib/aws-java-sdk-bundle-1.11.271.jar COPY core-site.xml /opt/hadoop/etc/hadoop/core-site.xml From 2cf6128f25a5cc15bc7707142229078bbc55028e Mon Sep 17 00:00:00 2001 From: Kevin Liu Date: Thu, 10 Sep 2026 02:41:14 -0700 Subject: [PATCH 4/4] fix(hms): use `get_table_req` for Hive 4 metastore compatibility (#3187) fix(hms): use get_table_req for Hive 4 metastore compatibility Hive 4.0.1 removed the `get_table` thrift method, so the HMS catalog fails with `Invalid method name: 'get_table'` against any Hive 4 metastore. Switch `load_table`, `table_exists`, and `rename_table` to `get_table_req`, which exists in every Hive release since 2.3 (the IDL the `hive_metastore` crate is generated from), so older metastores keep working. Move the integration test metastore from `apache/hive:3.1.3` (Debian Bullseye, EOL) to `apache/hive:4.2.1` so CI exercises a Hive 4 server, mirroring apache/iceberg-python#3924. Co-authored-by: Claude Fable 5.1 (cherry picked from commit 4aa55ba00c3ef2a4e603f323bfed0187952fe251) --- crates/catalog/hms/src/catalog.rs | 32 +++++++++++++++++++++++-------- dev/hms/Dockerfile | 18 +++++++---------- 2 files changed, 31 insertions(+), 19 deletions(-) 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/dev/hms/Dockerfile b/dev/hms/Dockerfile index dff5606bdc..5e6cfdbfef 100644 --- a/dev/hms/Dockerfile +++ b/dev/hms/Dockerfile @@ -13,20 +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 -ADD --chmod=644 --checksum=sha256:a18508b9348af095ea41301e439354dbd449e304ac44c6885b2b4fe78de88126 \ - https://repo1.maven.org/maven2/org/apache/hadoop/hadoop-aws/${HADOOP_VERSION}/hadoop-aws-${HADOOP_VERSION}.jar \ - /opt/hive/lib/hadoop-aws-${HADOOP_VERSION}.jar -ADD --chmod=644 --checksum=sha256:faf78ac4880f56cf52791d84ec1068ce7c66acc4295d580a726104b734c01fcd \ - https://repo1.maven.org/maven2/com/amazonaws/aws-java-sdk-bundle/1.11.271/aws-java-sdk-bundle-1.11.271.jar \ - /opt/hive/lib/aws-java-sdk-bundle-1.11.271.jar +# 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