diff --git a/Cargo.lock b/Cargo.lock index e08d17a42..47795c052 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4673,6 +4673,7 @@ dependencies = [ "futures", "hex", "hmac 0.12.1", + "http", "indexmap 2.14.0", "jiff", "libloading 0.9.0", diff --git a/crates/paimon/Cargo.toml b/crates/paimon/Cargo.toml index 7fb7d914a..830bfb0f4 100644 --- a/crates/paimon/Cargo.toml +++ b/crates/paimon/Cargo.toml @@ -129,6 +129,8 @@ parquet = { workspace = true, features = ["async", "zstd", "lz4", "snap"] } orc-rust = "0.8.0" async-stream = "0.3.6" reqwest = { version = "0.12", features = ["json", "gzip", "deflate"] } +# Already in the tree via `reqwest`; direct dep for the object storage User-Agent transport. +http = "1" # DLF authentication dependencies base64 = "0.22" hex = "0.4" diff --git a/crates/paimon/src/io/mod.rs b/crates/paimon/src/io/mod.rs index c7eea7c12..445993390 100644 --- a/crates/paimon/src/io/mod.rs +++ b/crates/paimon/src/io/mod.rs @@ -31,14 +31,23 @@ pub use storage::*; feature = "storage-oss", feature = "storage-s3" ))] -fn with_http_transport(op: opendal::Operator) -> opendal::Operator { +fn with_http_transport(op: opendal::Operator, user_agent: &str) -> opendal::Operator { use opendal::{HttpTransporter, OperationContext}; - use opendal_http_transport_reqwest::ReqwestTransport; - let transport = HttpTransporter::new(ReqwestTransport::default()); + let transport = HttpTransporter::new(user_agent::UserAgentTransport::new(user_agent)); op.with_context(OperationContext::new().with_http_transport(transport)) } +#[cfg(any( + feature = "storage-azdls", + feature = "storage-cos", + feature = "storage-gcs", + feature = "storage-obs", + feature = "storage-oss", + feature = "storage-s3" +))] +mod user_agent; + #[cfg(any( feature = "storage-s3", feature = "storage-cos", diff --git a/crates/paimon/src/io/storage.rs b/crates/paimon/src/io/storage.rs index 24885d290..f5ef94836 100644 --- a/crates/paimon/src/io/storage.rs +++ b/crates/paimon/src/io/storage.rs @@ -83,26 +83,31 @@ pub enum Storage { #[cfg(feature = "storage-s3")] S3 { config: Box, + user_agent: String, operators: Mutex>, }, #[cfg(feature = "storage-cos")] Cos { config: Box, + user_agent: String, operators: Mutex>, }, #[cfg(feature = "storage-azdls")] Azdls { config: Box, + user_agent: String, operators: Mutex>, }, #[cfg(feature = "storage-obs")] Obs { config: Box, + user_agent: String, operators: Mutex>, }, #[cfg(feature = "storage-gcs")] Gcs { config: Box, + user_agent: String, operators: Mutex>, }, #[cfg(feature = "storage-hdfs")] @@ -138,41 +143,51 @@ impl Storage { } #[cfg(feature = "storage-s3")] "s3" | "s3a" => { + let user_agent = super::user_agent::storage_user_agent(&props); let config = super::s3_config_parse(props)?; Ok(Self::S3 { config: Box::new(config), + user_agent, operators: Mutex::new(HashMap::new()), }) } #[cfg(feature = "storage-cos")] "cos" | "cosn" => { + let user_agent = super::user_agent::storage_user_agent(&props); let config = super::cos_config_parse(props)?; Ok(Self::Cos { config: Box::new(config), + user_agent, operators: Mutex::new(HashMap::new()), }) } #[cfg(feature = "storage-azdls")] "abfs" | "abfss" | "az" | "azdfs" | "azdls" | "azure" => { + let user_agent = super::user_agent::storage_user_agent(&props); let config = super::azdls_config_parse(props)?; Ok(Self::Azdls { config: Box::new(config), + user_agent, operators: Mutex::new(HashMap::new()), }) } #[cfg(feature = "storage-obs")] "obs" => { + let user_agent = super::user_agent::storage_user_agent(&props); let config = super::obs_config_parse(props)?; Ok(Self::Obs { config: Box::new(config), + user_agent, operators: Mutex::new(HashMap::new()), }) } #[cfg(feature = "storage-gcs")] "gcs" | "gs" => { + let user_agent = super::user_agent::storage_user_agent(&props); let config = super::gcs_config_parse(props)?; Ok(Self::Gcs { config: Box::new(config), + user_agent, operators: Mutex::new(HashMap::new()), }) } @@ -222,45 +237,65 @@ impl Storage { Ok((op, Cow::Borrowed(relative_path))) } #[cfg(feature = "storage-s3")] - Storage::S3 { config, operators } => { + Storage::S3 { + config, + user_agent, + operators, + } => { let (bucket, relative_path) = Self::bucket_and_relative_path(path, "S3", &["s3", "s3a"])?; - let op = Self::cached_s3_operator(config, operators, path, &bucket)?; + let op = Self::cached_s3_operator(config, user_agent, operators, path, &bucket)?; Ok((op, Cow::Borrowed(relative_path))) } #[cfg(feature = "storage-cos")] - Storage::Cos { config, operators } => { + Storage::Cos { + config, + user_agent, + operators, + } => { let (bucket, relative_path) = Self::bucket_and_relative_path(path, "COS", &["cos", "cosn"])?; let op = Self::cached_operator(operators, "COS", &bucket, || { - super::cos_config_build(config, path) + super::cos_config_build(config, path, user_agent) })?; Ok((op, Cow::Borrowed(relative_path))) } #[cfg(feature = "storage-azdls")] - Storage::Azdls { config, operators } => { + Storage::Azdls { + config, + user_agent, + operators, + } => { let relative_path = super::azdls_relative_path(path)?; let cache_key = super::azdls_operator_cache_key(config, path)?; let op = Self::cached_operator(operators, "Azure", &cache_key, || { - super::azdls_config_build(config, path) + super::azdls_config_build(config, path, user_agent) })?; Ok((op, Cow::Borrowed(relative_path))) } #[cfg(feature = "storage-obs")] - Storage::Obs { config, operators } => { + Storage::Obs { + config, + user_agent, + operators, + } => { let (bucket, relative_path) = Self::bucket_and_relative_path(path, "OBS", &["obs"])?; let op = Self::cached_operator(operators, "OBS", &bucket, || { - super::obs_config_build(config, path) + super::obs_config_build(config, path, user_agent) })?; Ok((op, Cow::Borrowed(relative_path))) } #[cfg(feature = "storage-gcs")] - Storage::Gcs { config, operators } => { + Storage::Gcs { + config, + user_agent, + operators, + } => { let (bucket, relative_path) = Self::bucket_and_relative_path(path, "GCS", &["gcs", "gs"])?; let op = Self::cached_operator(operators, "GCS", &bucket, || { - super::gcs_config_build(config, path) + super::gcs_config_build(config, path, user_agent) })?; Ok((op, Cow::Borrowed(relative_path))) } @@ -425,12 +460,13 @@ impl Storage { #[cfg(feature = "storage-s3")] fn cached_s3_operator( config: &S3Config, + user_agent: &str, operators: &Mutex>, path: &str, bucket: &str, ) -> crate::Result { Self::cached_operator(operators, "S3", bucket, || { - super::s3_config_build(config, path) + super::s3_config_build(config, path, user_agent) }) } } @@ -482,6 +518,28 @@ mod scheme_tests { } } + #[cfg(feature = "storage-s3")] + #[test] + fn s3_storage_uses_generic_user_agent_keys() { + let storage = Storage::build(FileIOBuilder::new("s3").with_props([ + ("user-agent.features", "Flink"), + ("user-agent.extended", "vvr"), + ("dlf.access-tracking.extended-info", "uid/123"), + ])) + .unwrap(); + let Storage::S3 { user_agent, .. } = storage else { + panic!("expected S3 storage"); + }; + assert_eq!( + user_agent, + format!( + "paimon-rust/{}(opendal/{};Flink) vvr uid/123", + env!("CARGO_PKG_VERSION"), + opendal::raw::VERSION + ) + ); + } + #[cfg(feature = "storage-cos")] #[test] fn cos_scheme_aliases_are_compatible() { diff --git a/crates/paimon/src/io/storage_azdls.rs b/crates/paimon/src/io/storage_azdls.rs index 6551f6171..9bcde5c45 100644 --- a/crates/paimon/src/io/storage_azdls.rs +++ b/crates/paimon/src/io/storage_azdls.rs @@ -64,11 +64,15 @@ pub(crate) fn azdls_config_parse(props: HashMap) -> Result Result { +pub(crate) fn azdls_config_build( + cfg: &AzdlsStorageConfig, + path: &str, + user_agent: &str, +) -> Result { let (cfg, relative_path) = azdls_config_for_path(cfg, path)?; let builder = cfg.into_builder(); - let op = super::with_http_transport(Operator::new(builder)?); + let op = super::with_http_transport(Operator::new(builder)?, user_agent); debug_assert_eq!( relative_path, @@ -394,14 +398,19 @@ mod tests { fn test_azdls_config_build_hadoop_form() { let cfg = azdls_config_parse(HashMap::new()).unwrap(); - let op = azdls_config_build(&cfg, "abfs://fs@account.dfs.core.windows.net/a/b").unwrap(); + let op = azdls_config_build( + &cfg, + "abfs://fs@account.dfs.core.windows.net/a/b", + "paimon-rust/test", + ) + .unwrap(); assert_eq!(op.info().name(), "fs"); } #[test] fn test_azdls_config_build_fsspec_form_requires_endpoint() { let cfg = azdls_config_parse(HashMap::new()).unwrap(); - let result = azdls_config_build(&cfg, "abfs://fs/a/b"); + let result = azdls_config_build(&cfg, "abfs://fs/a/b", "paimon-rust/test"); assert!(result.is_err()); } @@ -413,7 +422,7 @@ mod tests { )])) .unwrap(); - let op = azdls_config_build(&cfg, "abfs://fs/a/b").unwrap(); + let op = azdls_config_build(&cfg, "abfs://fs/a/b", "paimon-rust/test").unwrap(); assert_eq!(op.info().name(), "fs"); } @@ -439,7 +448,8 @@ mod tests { #[test] fn test_azdls_config_build_missing_filesystem() { let cfg = azdls_config_parse(HashMap::new()).unwrap(); - let result = azdls_config_build(&cfg, "abfs:///path/without/filesystem"); + let result = + azdls_config_build(&cfg, "abfs:///path/without/filesystem", "paimon-rust/test"); assert!(result.is_err()); } } diff --git a/crates/paimon/src/io/storage_cos.rs b/crates/paimon/src/io/storage_cos.rs index 61ea159a7..968841eb1 100644 --- a/crates/paimon/src/io/storage_cos.rs +++ b/crates/paimon/src/io/storage_cos.rs @@ -59,7 +59,7 @@ pub(crate) fn cos_config_parse(props: HashMap) -> Result Result { +pub(crate) fn cos_config_build(cfg: &CosConfig, path: &str, user_agent: &str) -> Result { let url = Url::parse(path).map_err(|_| Error::ConfigInvalid { message: format!("Invalid COS url: {path}"), })?; @@ -69,7 +69,10 @@ pub(crate) fn cos_config_build(cfg: &CosConfig, path: &str) -> Result })?; let builder = cfg.clone().into_builder().bucket(bucket); - Ok(super::with_http_transport(Operator::new(builder)?)) + Ok(super::with_http_transport( + Operator::new(builder)?, + user_agent, + )) } #[cfg(test)] @@ -124,14 +127,14 @@ mod tests { ..Default::default() }; - let op = cos_config_build(&cfg, "cosn://my-bucket/some/path").unwrap(); + let op = cos_config_build(&cfg, "cosn://my-bucket/some/path", "paimon-rust/test").unwrap(); assert_eq!(op.info().name(), "my-bucket"); } #[test] fn test_cos_config_build_missing_bucket() { let cfg = CosConfig::default(); - let result = cos_config_build(&cfg, "cosn:///path/without/bucket"); + let result = cos_config_build(&cfg, "cosn:///path/without/bucket", "paimon-rust/test"); assert!(result.is_err()); } } diff --git a/crates/paimon/src/io/storage_gcs.rs b/crates/paimon/src/io/storage_gcs.rs index 487ede59d..dccb34e9e 100644 --- a/crates/paimon/src/io/storage_gcs.rs +++ b/crates/paimon/src/io/storage_gcs.rs @@ -81,7 +81,7 @@ pub(crate) fn gcs_config_parse(props: HashMap) -> Result Result { +pub(crate) fn gcs_config_build(cfg: &GcsConfig, path: &str, user_agent: &str) -> Result { let url = Url::parse(path).map_err(|_| Error::ConfigInvalid { message: format!("Invalid GCS url: {path}"), })?; @@ -91,7 +91,10 @@ pub(crate) fn gcs_config_build(cfg: &GcsConfig, path: &str) -> Result })?; let builder = cfg.clone().into_builder().bucket(bucket); - Ok(super::with_http_transport(Operator::new(builder)?)) + Ok(super::with_http_transport( + Operator::new(builder)?, + user_agent, + )) } #[cfg(test)] @@ -185,14 +188,14 @@ mod tests { fn test_gcs_config_build_extracts_bucket() { let cfg = GcsConfig::default(); - let op = gcs_config_build(&cfg, "gs://my-bucket/some/path").unwrap(); + let op = gcs_config_build(&cfg, "gs://my-bucket/some/path", "paimon-rust/test").unwrap(); assert_eq!(op.info().name(), "my-bucket"); } #[test] fn test_gcs_config_build_missing_bucket() { let cfg = GcsConfig::default(); - let result = gcs_config_build(&cfg, "gs:///path/without/bucket"); + let result = gcs_config_build(&cfg, "gs:///path/without/bucket", "paimon-rust/test"); assert!(result.is_err()); } } diff --git a/crates/paimon/src/io/storage_obs.rs b/crates/paimon/src/io/storage_obs.rs index a64a97f25..42891d324 100644 --- a/crates/paimon/src/io/storage_obs.rs +++ b/crates/paimon/src/io/storage_obs.rs @@ -53,7 +53,7 @@ pub(crate) fn obs_config_parse(props: HashMap) -> Result Result { +pub(crate) fn obs_config_build(cfg: &ObsConfig, path: &str, user_agent: &str) -> Result { let url = Url::parse(path).map_err(|_| Error::ConfigInvalid { message: format!("Invalid OBS url: {path}"), })?; @@ -63,7 +63,10 @@ pub(crate) fn obs_config_build(cfg: &ObsConfig, path: &str) -> Result })?; let builder = cfg.clone().into_builder().bucket(bucket); - Ok(super::with_http_transport(Operator::new(builder)?)) + Ok(super::with_http_transport( + Operator::new(builder)?, + user_agent, + )) } #[cfg(test)] @@ -120,14 +123,14 @@ mod tests { let mut cfg = ObsConfig::default(); cfg.endpoint = Some("https://obs.cn-north-4.myhuaweicloud.com".to_string()); - let op = obs_config_build(&cfg, "obs://my-bucket/some/path").unwrap(); + let op = obs_config_build(&cfg, "obs://my-bucket/some/path", "paimon-rust/test").unwrap(); assert_eq!(op.info().name(), "my-bucket"); } #[test] fn test_obs_config_build_missing_bucket() { let cfg = ObsConfig::default(); - let result = obs_config_build(&cfg, "obs:///path/without/bucket"); + let result = obs_config_build(&cfg, "obs:///path/without/bucket", "paimon-rust/test"); assert!(result.is_err()); } } diff --git a/crates/paimon/src/io/storage_oss.rs b/crates/paimon/src/io/storage_oss.rs index 83923d707..17b99bea7 100644 --- a/crates/paimon/src/io/storage_oss.rs +++ b/crates/paimon/src/io/storage_oss.rs @@ -61,16 +61,18 @@ pub struct OssStorageConfig { service: OssConfig, retry_count: usize, retry_interval: Duration, + user_agent: String, } /// Parse paimon catalog options into an [`OssStorageConfig`]. /// /// Extracts OSS-related configuration keys (endpoint, access key, secret key, -/// optional security token, and retry settings) from the provided properties. +/// optional security token, retry and User-Agent settings) from the provided properties. /// /// Returns an error if any required configuration key is missing. pub(crate) fn oss_config_parse(mut props: HashMap) -> Result { let mut cfg = OssConfig::default(); + let user_agent = super::user_agent::oss_user_agent(&props); cfg.endpoint = Some( props @@ -107,6 +109,7 @@ pub(crate) fn oss_config_parse(mut props: HashMap) -> Result Result>>>, + headers: HeaderMap, + ) -> StatusCode { + let user_agent = headers + .get(header::USER_AGENT) + .map(|value| value.to_str().unwrap().to_string()) + .unwrap_or_default(); + user_agents.lock().unwrap().push(user_agent); + StatusCode::NOT_FOUND + } + #[test] fn test_oss_config_parse_with_all_keys() { let mut props = required_props(); @@ -292,4 +308,44 @@ mod tests { assert!(op.read("object").await.is_err()); assert_eq!(attempts.load(Ordering::SeqCst), 2); } + + #[tokio::test] + async fn test_oss_sends_configured_user_agent() { + use crate::io::user_agent::{ + DLF_ACCESS_TRACKING_EXTENDED_INFO, OSS_USER_AGENT_EXTENDED, USER_AGENT_EXTENDED, + USER_AGENT_FEATURES, + }; + + let user_agents = Arc::new(Mutex::new(Vec::new())); + let app = Router::new() + .fallback(record_user_agent) + .with_state(user_agents.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + let mut props = required_props(); + props.insert(OSS_ENDPOINT.to_string(), format!("http://{address}")); + props.insert(USER_AGENT_FEATURES.to_string(), "DataFusion".to_string()); + props.insert(USER_AGENT_EXTENDED.to_string(), "overridden".to_string()); + props.insert(OSS_USER_AGENT_EXTENDED.to_string(), "vvr".to_string()); + props.insert( + DLF_ACCESS_TRACKING_EXTENDED_INFO.to_string(), + "uid/123".to_string(), + ); + let mut cfg = oss_config_parse(props).unwrap(); + cfg.service.addressing_style = Some("path".to_string()); + cfg.service.skip_signature = true; + + let op = oss_config_build(&cfg, "oss://bucket/path").unwrap(); + assert!(!op.exists("object").await.unwrap()); + assert_eq!( + *user_agents.lock().unwrap(), + vec![format!( + "paimon-rust/{}(opendal/{};DataFusion) vvr uid/123", + env!("CARGO_PKG_VERSION"), + opendal::raw::VERSION + )] + ); + } } diff --git a/crates/paimon/src/io/storage_s3.rs b/crates/paimon/src/io/storage_s3.rs index 8d77a3659..7aaaa0c3e 100644 --- a/crates/paimon/src/io/storage_s3.rs +++ b/crates/paimon/src/io/storage_s3.rs @@ -128,7 +128,7 @@ pub(crate) fn s3_config_parse(props: HashMap) -> Result Result { +pub(crate) fn s3_config_build(cfg: &S3Config, path: &str, user_agent: &str) -> Result { let url = Url::parse(path).map_err(|_| Error::ConfigInvalid { message: format!("Invalid S3 url: {path}"), })?; @@ -138,7 +138,10 @@ pub(crate) fn s3_config_build(cfg: &S3Config, path: &str) -> Result { })?; let builder = cfg.clone().into_builder().bucket(bucket); - Ok(super::with_http_transport(Operator::new(builder)?)) + Ok(super::with_http_transport( + Operator::new(builder)?, + user_agent, + )) } #[cfg(test)] @@ -309,7 +312,7 @@ mod tests { cfg.endpoint = Some("https://s3.us-east-1.amazonaws.com".to_string()); cfg.region = Some("us-east-1".to_string()); - let op = s3_config_build(&cfg, "s3://my-bucket/some/path").unwrap(); + let op = s3_config_build(&cfg, "s3://my-bucket/some/path", "paimon-rust/test").unwrap(); assert_eq!(op.info().name(), "my-bucket"); } @@ -320,21 +323,21 @@ mod tests { cfg.endpoint = Some("https://s3.us-east-1.amazonaws.com".to_string()); cfg.region = Some("us-east-1".to_string()); - let op = s3_config_build(&cfg, "s3a://my-bucket/some/path").unwrap(); + let op = s3_config_build(&cfg, "s3a://my-bucket/some/path", "paimon-rust/test").unwrap(); assert_eq!(op.info().name(), "my-bucket"); } #[test] fn test_s3_config_build_invalid_url() { let cfg = S3Config::default(); - let result = s3_config_build(&cfg, "not-a-valid-url"); + let result = s3_config_build(&cfg, "not-a-valid-url", "paimon-rust/test"); assert!(result.is_err()); } #[test] fn test_s3_config_build_missing_bucket() { let cfg = S3Config::default(); - let result = s3_config_build(&cfg, "s3:///path/without/bucket"); + let result = s3_config_build(&cfg, "s3:///path/without/bucket", "paimon-rust/test"); assert!(result.is_err()); } diff --git a/crates/paimon/src/io/user_agent.rs b/crates/paimon/src/io/user_agent.rs new file mode 100644 index 000000000..8bb21e2c8 --- /dev/null +++ b/crates/paimon/src/io/user_agent.rs @@ -0,0 +1,291 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Paimon's unified User-Agent for object storage requests: +//! `([;feature...])[ ]`. + +use std::collections::HashMap; + +use http::header::{HeaderValue, USER_AGENT}; +use http::{Request, Response}; +use opendal::{Buffer, HttpBody, HttpTransport}; +use opendal_http_transport_reqwest::ReqwestTransport; + +/// Replaces the default `paimon-rust/` module; shared with REST requests. +pub(crate) const USER_AGENT_MODULE: &str = "user-agent.module"; + +/// Space-separated features, rendered after the transport and separated by `;`. +pub(crate) const USER_AGENT_FEATURES: &str = "user-agent.features"; + +/// Free text after the parentheses. +pub(crate) const USER_AGENT_EXTENDED: &str = "user-agent.extended"; + +/// OSS-only keys, each overriding the matching `user-agent.*` key. +pub(crate) const OSS_USER_AGENT_MODULE: &str = "fs.oss.user.agent.module"; +pub(crate) const OSS_USER_AGENT_FEATURES: &str = "fs.oss.user.agent.features"; +pub(crate) const OSS_USER_AGENT_EXTENDED: &str = "fs.oss.user.agent.extended"; + +/// Set by the DLF data token and appended after the extended part. +pub(crate) const DLF_ACCESS_TRACKING_EXTENDED_INFO: &str = "dlf.access-tracking.extended-info"; + +/// The User-Agent sent when no `user-agent.*` option is set. +pub(crate) fn default_user_agent() -> String { + storage_user_agent(&HashMap::new()) +} + +/// Builds the User-Agent from the `user-agent.*` options. +pub(crate) fn storage_user_agent(props: &HashMap) -> String { + build_user_agent(props, None) +} + +/// Builds the OSS User-Agent, where `fs.oss.user.agent.*` overrides `user-agent.*` per part. +pub(crate) fn oss_user_agent(props: &HashMap) -> String { + build_user_agent( + props, + Some([ + OSS_USER_AGENT_MODULE, + OSS_USER_AGENT_FEATURES, + OSS_USER_AGENT_EXTENDED, + ]), + ) +} + +fn build_user_agent(props: &HashMap, overrides: Option<[&str; 3]>) -> String { + let non_blank = |key: &str| { + props + .get(key) + .map(|value| value.trim()) + .filter(|value| !value.is_empty()) + }; + let part = |index: usize, key: &str| { + overrides + .and_then(|keys| non_blank(keys[index])) + .or_else(|| non_blank(key)) + }; + + let mut user_agent = match part(0, USER_AGENT_MODULE) { + Some(module) => module.to_string(), + None => format!("paimon-rust/{}", env!("CARGO_PKG_VERSION")), + }; + user_agent.push_str("(opendal/"); + user_agent.push_str(opendal::raw::VERSION); + for feature in part(1, USER_AGENT_FEATURES) + .into_iter() + .flat_map(|features| features.split_whitespace()) + { + user_agent.push(';'); + user_agent.push_str(feature); + } + user_agent.push(')'); + + let extended: Vec<&str> = [ + part(2, USER_AGENT_EXTENDED), + non_blank(DLF_ACCESS_TRACKING_EXTENDED_INFO), + ] + .into_iter() + .flatten() + .collect(); + if !extended.is_empty() { + user_agent.push(' '); + user_agent.push_str(&extended.join(" ")); + } + user_agent +} + +/// Sends every request through the shared reqwest client with a User-Agent, unless one is set. +#[derive(Clone)] +pub(crate) struct UserAgentTransport { + user_agent: HeaderValue, + inner: ReqwestTransport, +} + +impl UserAgentTransport { + pub(crate) fn new(user_agent: &str) -> Self { + let user_agent = HeaderValue::from_bytes(user_agent.as_bytes()).unwrap_or_else(|_| { + log::warn!("Invalid object storage User-Agent {user_agent:?}, using the default"); + HeaderValue::from_str(&default_user_agent()) + .expect("the default User-Agent is visible ASCII") + }); + Self { + user_agent, + inner: ReqwestTransport::default(), + } + } + + fn apply(&self, req: &mut Request) { + req.headers_mut() + .entry(USER_AGENT) + .or_insert_with(|| self.user_agent.clone()); + } +} + +impl HttpTransport for UserAgentTransport { + async fn fetch(&self, mut req: Request) -> opendal::Result> { + self.apply(&mut req); + self.inner.fetch(req).await + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn props(entries: &[(&str, &str)]) -> HashMap { + entries + .iter() + .map(|(key, value)| (key.to_string(), value.to_string())) + .collect() + } + + fn default_prefix() -> String { + format!( + "paimon-rust/{}(opendal/{}", + env!("CARGO_PKG_VERSION"), + opendal::raw::VERSION + ) + } + + #[test] + fn test_default_user_agent() { + assert_eq!(default_user_agent(), format!("{})", default_prefix())); + } + + #[test] + fn test_module_override() { + let user_agent = oss_user_agent(&props(&[(OSS_USER_AGENT_MODULE, "MyApp/1.0")])); + assert_eq!( + user_agent, + format!("MyApp/1.0(opendal/{})", opendal::raw::VERSION) + ); + } + + #[test] + fn test_features_are_split_and_joined() { + let user_agent = oss_user_agent(&props(&[(OSS_USER_AGENT_FEATURES, " Flink Paimon ")])); + assert_eq!(user_agent, format!("{};Flink;Paimon)", default_prefix())); + } + + #[test] + fn test_access_tracking_is_appended_to_user_extended() { + let user_agent = oss_user_agent(&props(&[ + (OSS_USER_AGENT_EXTENDED, "vvr"), + (DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123 user/alice"), + ])); + assert_eq!( + user_agent, + format!("{}) vvr uid/123 user/alice", default_prefix()) + ); + } + + #[test] + fn test_access_tracking_only() { + let user_agent = oss_user_agent(&props(&[( + DLF_ACCESS_TRACKING_EXTENDED_INFO, + "uid/123 user/alice", + )])); + assert_eq!( + user_agent, + format!("{}) uid/123 user/alice", default_prefix()) + ); + } + + #[test] + fn test_blank_options_are_ignored() { + let user_agent = oss_user_agent(&props(&[ + (OSS_USER_AGENT_MODULE, " "), + (OSS_USER_AGENT_FEATURES, " "), + (OSS_USER_AGENT_EXTENDED, ""), + (USER_AGENT_MODULE, " "), + (USER_AGENT_FEATURES, ""), + (USER_AGENT_EXTENDED, " "), + (DLF_ACCESS_TRACKING_EXTENDED_INFO, " "), + ])); + assert_eq!(user_agent, default_user_agent()); + } + + #[test] + fn test_generic_keys_apply_to_any_backend() { + let user_agent = storage_user_agent(&props(&[ + (USER_AGENT_MODULE, "MyApp/1.0"), + (USER_AGENT_FEATURES, "Flink"), + (USER_AGENT_EXTENDED, "vvr"), + (DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123"), + (OSS_USER_AGENT_EXTENDED, "ignored"), + ])); + assert_eq!( + user_agent, + format!( + "MyApp/1.0(opendal/{};Flink) vvr uid/123", + opendal::raw::VERSION + ) + ); + } + + #[test] + fn test_oss_falls_back_to_generic_keys() { + let user_agent = oss_user_agent(&props(&[ + (USER_AGENT_FEATURES, "Flink"), + (USER_AGENT_EXTENDED, "vvr"), + (DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123"), + ])); + assert_eq!( + user_agent, + format!("{};Flink) vvr uid/123", default_prefix()) + ); + } + + #[test] + fn test_oss_keys_override_generic_keys_per_part() { + let user_agent = oss_user_agent(&props(&[ + (USER_AGENT_MODULE, "Generic/1.0"), + (USER_AGENT_FEATURES, "Flink"), + (USER_AGENT_EXTENDED, "vvr"), + (OSS_USER_AGENT_FEATURES, "Spark"), + (OSS_USER_AGENT_EXTENDED, "oss/ext"), + (DLF_ACCESS_TRACKING_EXTENDED_INFO, "uid/123"), + ])); + assert_eq!( + user_agent, + format!( + "Generic/1.0(opendal/{};Spark) oss/ext uid/123", + opendal::raw::VERSION + ) + ); + } + + #[test] + fn test_existing_user_agent_is_kept() { + let transport = UserAgentTransport::new("paimon-rust/test"); + + let mut req = Request::new(Buffer::new()); + transport.apply(&mut req); + assert_eq!(req.headers()[USER_AGENT], "paimon-rust/test"); + + let mut req = Request::builder() + .header(USER_AGENT, "caller/1.0") + .body(Buffer::new()) + .unwrap(); + transport.apply(&mut req); + assert_eq!(req.headers()[USER_AGENT], "caller/1.0"); + } + + #[test] + fn test_invalid_user_agent_falls_back_to_default() { + let transport = UserAgentTransport::new("bad\nagent"); + assert_eq!(transport.user_agent, default_user_agent().as_str()); + } +}