From 70cf3bf2eafa47a2d6f7853511c8493cb6828cd4 Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Tue, 29 Sep 2026 00:44:06 +0800 Subject: [PATCH 1/2] [core] Unify the REST User-Agent format REST catalog requests went out without a User-Agent. They now send Paimon's unified User-Agent format, module(transport;features) extended, by default paimon-rust/(reqwest). The parts come from the catalog options user-agent.module, user-agent.features and user-agent.extended, which object storage requests share. A user-set header.User-Agent option still takes precedence and is the only User-Agent sent. The DLF ECS token client sends the default. --- crates/paimon/src/api/auth/dlf_provider.rs | 39 ++++++++ crates/paimon/src/api/mod.rs | 1 + crates/paimon/src/api/rest_api.rs | 6 +- crates/paimon/src/api/rest_client.rs | 101 ++++++++++++++++++++- crates/paimon/src/api/user_agent.rs | 90 ++++++++++++++++++ 5 files changed, 235 insertions(+), 2 deletions(-) create mode 100644 crates/paimon/src/api/user_agent.rs diff --git a/crates/paimon/src/api/auth/dlf_provider.rs b/crates/paimon/src/api/auth/dlf_provider.rs index fc95eea93..18b0ccc61 100644 --- a/crates/paimon/src/api/auth/dlf_provider.rs +++ b/crates/paimon/src/api/auth/dlf_provider.rs @@ -27,6 +27,7 @@ use serde::{Deserialize, Serialize}; use super::base::{AuthProvider, RESTAuthParameter, AUTHORIZATION_HEADER_KEY}; use super::dlf_signer::{DLFRequestSigner, DLFSignerFactory}; +use crate::api::user_agent; use crate::common::{CatalogOptions, Options}; use crate::error::Error; use crate::Result; @@ -474,6 +475,7 @@ impl TokenHTTPClient { let client = Client::builder() .timeout(read_timeout) .connect_timeout(connect_timeout) + .user_agent(user_agent::default_rest_user_agent()) .build() .expect("Failed to create HTTP client"); @@ -520,6 +522,43 @@ impl TokenHTTPClient { mod tests { use super::*; + #[tokio::test] + async fn test_ecs_token_loader_sends_rest_user_agent() { + use axum::http::{header, HeaderMap}; + use axum::routing::get; + use axum::Router; + use std::sync::Mutex; + + let user_agents = Arc::new(Mutex::new(Vec::new())); + let recorded = user_agents.clone(); + let app = Router::new().route( + "/role", + get(move |headers: HeaderMap| async move { + recorded.lock().unwrap().extend( + headers + .get_all(header::USER_AGENT) + .iter() + .map(|value| value.to_str().unwrap().to_string()), + ); + r#"{"AccessKeyId":"ak","AccessKeySecret":"sk","SecurityToken":"st"}"# + }), + ); + 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 loader = DLFECSTokenLoader::new(format!("http://{address}/"), Some("role".to_string())); + loader.load_token().await.unwrap(); + + assert_eq!( + *user_agents.lock().unwrap(), + vec![format!( + "paimon-rust/{}(reqwest)", + env!("CARGO_PKG_VERSION") + )] + ); + } + #[test] fn test_extract_host() { let uri = "http://dlf-abcdfgerrf.net/api/v1"; diff --git a/crates/paimon/src/api/mod.rs b/crates/paimon/src/api/mod.rs index 55a1dd850..cffb8865c 100644 --- a/crates/paimon/src/api/mod.rs +++ b/crates/paimon/src/api/mod.rs @@ -29,6 +29,7 @@ pub mod rest_error; pub mod rest_util; mod api_response; +mod user_agent; // Re-export request types pub use api_request::{ diff --git a/crates/paimon/src/api/rest_api.rs b/crates/paimon/src/api/rest_api.rs index 7e435ae17..b13d91395 100644 --- a/crates/paimon/src/api/rest_api.rs +++ b/crates/paimon/src/api/rest_api.rs @@ -146,7 +146,11 @@ impl RESTApi { // Create auth function first, before making any requests let rest_auth_function = RESTAuthFunction::new(base_headers.clone(), auth_provider); - let mut client = HttpClient::new(uri, Some(rest_auth_function))?; + let mut client = HttpClient::with_user_agent( + uri, + Some(rest_auth_function), + &super::user_agent::rest_user_agent(&options), + )?; let options = if config_required { let warehouse = options.get(CatalogOptions::WAREHOUSE).ok_or_else(|| { diff --git a/crates/paimon/src/api/rest_client.rs b/crates/paimon/src/api/rest_client.rs index 194d73748..beaef849f 100644 --- a/crates/paimon/src/api/rest_client.rs +++ b/crates/paimon/src/api/rest_client.rs @@ -19,8 +19,10 @@ use super::auth::{RESTAuthFunction, RESTAuthParameter}; use super::rest_error::RestError; +use super::user_agent; use crate::Error; use crate::Result; +use reqwest::header::HeaderValue; use serde::de::DeserializeOwned; use std::collections::HashMap; use std::time::Duration; @@ -42,10 +44,29 @@ impl HttpClient { /// # Returns /// A new HttpClient instance. pub fn new(base_url: &str, auth_function: Option) -> Result { + Self::with_user_agent( + base_url, + auth_function, + &user_agent::default_rest_user_agent(), + ) + } + + /// Like [`HttpClient::new`], sending `agent` unless a request sets its own User-Agent. + pub(crate) fn with_user_agent( + base_url: &str, + auth_function: Option, + agent: &str, + ) -> Result { let final_url = Self::normalize_uri(base_url)?; + let agent = HeaderValue::from_str(agent).unwrap_or_else(|_| { + log::warn!("Invalid REST User-Agent {agent:?}, using the default"); + HeaderValue::from_str(&user_agent::default_rest_user_agent()) + .expect("the default User-Agent is visible ASCII") + }); let client = reqwest::Client::builder() .timeout(Duration::from_secs(30)) + .user_agent(agent) .build() .map_err(|e| Error::ConfigInvalid { message: format!("Failed to create HTTP client: {e}"), @@ -283,7 +304,7 @@ mod tests { use async_trait::async_trait; use axum::extract::State; - use axum::http::Uri; + use axum::http::{header, HeaderMap, Uri}; use axum::routing::get; use axum::{Json, Router}; @@ -312,6 +333,84 @@ mod tests { Json(serde_json::json!({})) } + async fn record_user_agents( + State(user_agents): State>>>>, + headers: HeaderMap, + ) -> Json { + let values = headers + .get_all(header::USER_AGENT) + .iter() + .map(|value| value.to_str().unwrap().to_string()) + .collect(); + user_agents.lock().unwrap().push(values); + Json(serde_json::json!({"databases": [], "nextPageToken": null})) + } + + /// The User-Agent headers a REST catalog sends to list databases, with `extra` catalog options. + async fn sent_user_agents(extra: &[(&str, &str)]) -> Vec { + let user_agents = Arc::new(Mutex::new(Vec::new())); + let app = Router::new() + .route("/v1/databases", get(record_user_agents)) + .with_state(user_agents.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let _server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + + let mut options = crate::common::Options::new(); + options.set("uri", format!("http://{address}")); + options.set("token.provider", "bear"); + options.set("token", "token"); + for (key, value) in extra { + options.set(*key, *value); + } + let api = crate::api::rest_api::RESTApi::new(options, false) + .await + .unwrap(); + api.list_databases().await.unwrap(); + + let mut user_agents = user_agents.lock().unwrap(); + assert_eq!(user_agents.len(), 1); + user_agents.pop().unwrap() + } + + #[tokio::test] + async fn test_default_user_agent_is_sent() { + assert_eq!( + sent_user_agents(&[]).await, + vec![format!( + "paimon-rust/{}(reqwest)", + env!("CARGO_PKG_VERSION") + )] + ); + } + + #[tokio::test] + async fn test_user_agent_options_are_sent() { + let options = [ + ("user-agent.features", "Flink"), + ("user-agent.extended", "vvr"), + ]; + assert_eq!( + sent_user_agents(&options).await, + vec![format!( + "paimon-rust/{}(reqwest;Flink) vvr", + env!("CARGO_PKG_VERSION") + )] + ); + } + + #[tokio::test] + async fn test_user_agent_header_option_wins() { + let options = [ + ("header.User-Agent", "starrocks/user"), + ("user-agent.features", "Flink"), + ]; + assert_eq!( + sent_user_agents(&options).await, + vec!["starrocks/user".to_string()] + ); + } + fn canonical(pairs: impl Iterator) -> String { let mut parts: Vec = pairs.collect(); parts.sort(); diff --git a/crates/paimon/src/api/user_agent.rs b/crates/paimon/src/api/user_agent.rs new file mode 100644 index 000000000..d8c9ee206 --- /dev/null +++ b/crates/paimon/src/api/user_agent.rs @@ -0,0 +1,90 @@ +// 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 REST requests: +//! `([;feature...])[ ]`. + +use crate::common::Options; + +/// Replaces the default `paimon-rust/` module; shared with object storage 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"; + +/// The User-Agent of REST requests without options, `paimon-rust/(reqwest)`. +pub(crate) fn default_rest_user_agent() -> String { + rest_user_agent(&Options::new()) +} + +/// Builds the REST User-Agent from the `user-agent.*` catalog options. +pub(crate) fn rest_user_agent(options: &Options) -> String { + let non_blank = |key: &str| { + options + .get(key) + .map(|value| value.trim()) + .filter(|value| !value.is_empty()) + }; + + let mut user_agent = match non_blank(USER_AGENT_MODULE) { + Some(module) => module.to_string(), + None => format!("paimon-rust/{}", env!("CARGO_PKG_VERSION")), + }; + user_agent.push_str("(reqwest"); + for feature in non_blank(USER_AGENT_FEATURES) + .into_iter() + .flat_map(|features| features.split_whitespace()) + { + user_agent.push(';'); + user_agent.push_str(feature); + } + user_agent.push(')'); + if let Some(extended) = non_blank(USER_AGENT_EXTENDED) { + user_agent.push(' '); + user_agent.push_str(extended); + } + user_agent +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_rest_user_agent() { + let version = env!("CARGO_PKG_VERSION"); + assert_eq!( + default_rest_user_agent(), + format!("paimon-rust/{version}(reqwest)") + ); + + let mut options = Options::new(); + options.set(USER_AGENT_FEATURES, " Flink PVFS "); + options.set(USER_AGENT_EXTENDED, " vvr "); + assert_eq!( + rest_user_agent(&options), + format!("paimon-rust/{version}(reqwest;Flink;PVFS) vvr") + ); + + options.set(USER_AGENT_MODULE, "MyApp/1.0"); + options.set(USER_AGENT_EXTENDED, " "); + assert_eq!(rest_user_agent(&options), "MyApp/1.0(reqwest;Flink;PVFS)"); + } +} From 7979517f404a207d4929dc74602359b0638ea397 Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Fri, 2 Oct 2026 00:47:23 +0800 Subject: [PATCH 2/2] [core] Rebuild the REST User-Agent after merging the server config The REST client was built from the initial options only, so user-agent.* options delivered through /v1/config were ignored. RESTApi::new now rebuilds the default User-Agent from the merged options after bootstrap; an explicit header.User-Agent still takes precedence per request. --- crates/paimon/src/api/rest_api.rs | 2 + crates/paimon/src/api/rest_client.rs | 68 +++++++++++++++++++++++----- 2 files changed, 59 insertions(+), 11 deletions(-) diff --git a/crates/paimon/src/api/rest_api.rs b/crates/paimon/src/api/rest_api.rs index b13d91395..7a59f77f7 100644 --- a/crates/paimon/src/api/rest_api.rs +++ b/crates/paimon/src/api/rest_api.rs @@ -182,6 +182,8 @@ impl RESTApi { let rest_auth_function = RESTAuthFunction::new(base_headers, auth_provider); client.set_auth_function(rest_auth_function); + // Server config may carry user-agent.* options; header.User-Agent still wins per request. + client.set_user_agent(&super::user_agent::rest_user_agent(&merged))?; merged } else { diff --git a/crates/paimon/src/api/rest_client.rs b/crates/paimon/src/api/rest_client.rs index beaef849f..5448d47cf 100644 --- a/crates/paimon/src/api/rest_client.rs +++ b/crates/paimon/src/api/rest_client.rs @@ -57,26 +57,32 @@ impl HttpClient { auth_function: Option, agent: &str, ) -> Result { - let final_url = Self::normalize_uri(base_url)?; + Ok(HttpClient { + client: Self::build_client(agent)?, + base_url: Self::normalize_uri(base_url)?, + auth_function, + }) + } + + /// Replace the default User-Agent, e.g. after options are merged with the server config. + pub(crate) fn set_user_agent(&mut self, agent: &str) -> Result<()> { + self.client = Self::build_client(agent)?; + Ok(()) + } + + fn build_client(agent: &str) -> Result { let agent = HeaderValue::from_str(agent).unwrap_or_else(|_| { log::warn!("Invalid REST User-Agent {agent:?}, using the default"); HeaderValue::from_str(&user_agent::default_rest_user_agent()) .expect("the default User-Agent is visible ASCII") }); - - let client = reqwest::Client::builder() + reqwest::Client::builder() .timeout(Duration::from_secs(30)) .user_agent(agent) .build() .map_err(|e| Error::ConfigInvalid { message: format!("Failed to create HTTP client: {e}"), - })?; - - Ok(HttpClient { - client, - base_url: final_url, - auth_function, - }) + }) } /// Normalize and validate a URI. @@ -348,9 +354,20 @@ mod tests { /// The User-Agent headers a REST catalog sends to list databases, with `extra` catalog options. async fn sent_user_agents(extra: &[(&str, &str)]) -> Vec { + sent_user_agents_with_config(extra, None).await + } + + /// Like [`sent_user_agents`], bootstrapping from a server returning `config` when set. + async fn sent_user_agents_with_config( + extra: &[(&str, &str)], + config: Option, + ) -> Vec { let user_agents = Arc::new(Mutex::new(Vec::new())); + let config_required = config.is_some(); + let config = config.unwrap_or_else(|| serde_json::json!({})); let app = Router::new() .route("/v1/databases", get(record_user_agents)) + .route("/v1/config", get(move || async move { Json(config) })) .with_state(user_agents.clone()); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); @@ -358,12 +375,13 @@ mod tests { let mut options = crate::common::Options::new(); options.set("uri", format!("http://{address}")); + options.set("warehouse", "warehouse"); options.set("token.provider", "bear"); options.set("token", "token"); for (key, value) in extra { options.set(*key, *value); } - let api = crate::api::rest_api::RESTApi::new(options, false) + let api = crate::api::rest_api::RESTApi::new(options, config_required) .await .unwrap(); api.list_databases().await.unwrap(); @@ -411,6 +429,34 @@ mod tests { ); } + #[tokio::test] + async fn test_user_agent_options_from_server_config_are_sent() { + let config = serde_json::json!({ + "defaults": {"user-agent.features": "ServerFeature"}, + "overrides": {"user-agent.extended": "catalog-tag"}, + }); + assert_eq!( + sent_user_agents_with_config(&[], Some(config)).await, + vec![format!( + "paimon-rust/{}(reqwest;ServerFeature) catalog-tag", + env!("CARGO_PKG_VERSION") + )] + ); + } + + #[tokio::test] + async fn test_user_agent_header_option_wins_over_server_config() { + let config = serde_json::json!({ + "defaults": {"user-agent.features": "ServerFeature"}, + "overrides": {"user-agent.extended": "catalog-tag"}, + }); + assert_eq!( + sent_user_agents_with_config(&[("header.User-Agent", "starrocks/user")], Some(config)) + .await, + vec!["starrocks/user".to_string()] + ); + } + fn canonical(pairs: impl Iterator) -> String { let mut parts: Vec = pairs.collect(); parts.sort();