github-actions[bot] commented on code in PR #68540:
URL: https://github.com/apache/doris/pull/68540#discussion_r4226469640
##########
fe/be-java-extensions/hadoop-deps/pom.xml:
##########
@@ -73,6 +73,11 @@ under the License.
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
</dependency>
+ <dependency>
+ <groupId>com.google.cloud.bigdataoss</groupId>
+ <artifactId>gcs-connector</artifactId>
+ <classifier>shaded</classifier>
+ </dependency>
Review Comment:
[P1] Package the GCS connector with the isolated JNI scanners. Native GCP
config now sets `fs.gs.impl` to `GoogleHadoopFileSystem`, but
`hadoop-hudi-scanner` and `paimon-scanner` each exclude every `hadoop-deps`
transitive and have no direct `gcs-connector` dependency. Their plugin
classloaders cannot see the BE system classpath, so a Hudi MOR/forced-JNI or
Paimon read of `gs://` data fails when Hadoop tries to load the configured
filesystem class. Add the shaded connector to both scanner runtime directories
and cover those JNI reads.
##########
cloud/src/common/auth/obj_credential.cpp:
##########
@@ -0,0 +1,239 @@
+// 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.
+
+#include "common/auth/obj_credential.h"
+
+#include <gen_cpp/cloud.pb.h>
+
+#include <string_view>
+
+#include "cpp/obj-client/auth/obj_credential.h"
+
+namespace doris {
+namespace {
+
+GcpCredentialConfig to_gcp_credential_config(const cloud::GcpCredentialPB&
credential) {
+ GcpCredentialConfig config;
+ switch (credential.credential_provider_type()) {
+ case cloud::GcpCredentialPB::DEFAULT:
+ config.provider_type = GcpCredentialProviderType::Default;
+ break;
+ case cloud::GcpCredentialPB::COMPUTE_ENGINE:
+ config.provider_type = GcpCredentialProviderType::ComputeEngine;
+ break;
+ }
+ config.impersonation_service_account =
credential.impersonation_service_account();
+ return config;
+}
+
+} // namespace
+
+void convert_obj_credential(const cloud::ObjectStoreInfoPB& source,
CredentialConfig* credential) {
+ if (!source.has_credential() || !source.credential().has_gcp_credential())
{
+ *credential = std::monostate {};
+ return;
+ }
+ *credential =
to_gcp_credential_config(source.credential().gcp_credential());
+}
+
+} // namespace doris
+
+namespace doris::cloud {
+namespace {
+
+// Extension checklist for a new passwordless provider:
+// 1. add a provider-specific validate_<provider>_obj_credential();
+// 2. dispatch it from validate_obj_credential(); and
+// 3. reject requests that set more than one provider credential.
+// Generic create/alter normalization does not need provider-specific branches.
+ObjCredentialProvider to_obj_credential_provider(const ObjectStoreInfoPB& obj)
{
+ if (!obj.has_provider()) {
+ return ObjCredentialProvider::Unknown;
+ }
+ switch (obj.provider()) {
+ case ObjectStoreInfoPB::S3:
+ return ObjCredentialProvider::Aws;
+ case ObjectStoreInfoPB::AZURE:
+ return ObjCredentialProvider::Azure;
+ case ObjectStoreInfoPB::BOS:
+ return ObjCredentialProvider::Bos;
+ case ObjectStoreInfoPB::COS:
+ return ObjCredentialProvider::Cos;
+ case ObjectStoreInfoPB::GCP:
+ return ObjCredentialProvider::Gcp;
+ case ObjectStoreInfoPB::OBS:
+ return ObjCredentialProvider::Obs;
+ case ObjectStoreInfoPB::OSS:
+ return ObjCredentialProvider::Oss;
+ case ObjectStoreInfoPB::TOS:
+ return ObjCredentialProvider::Tos;
+ case ObjectStoreInfoPB::UNKONWN:
+ return ObjCredentialProvider::Unknown;
+ }
+ return ObjCredentialProvider::Unknown;
+}
+
+ObjCredentialValidationContext credential_validation_context(const
ObjectStoreInfoPB& obj) {
+ return ObjCredentialValidationContext {
+ .provider = to_obj_credential_provider(obj),
+ .ak = obj.ak(),
+ .sk = obj.sk(),
+ .has_aws_role_arn = !obj.role_arn().empty(),
+ .has_aws_external_id = !obj.external_id().empty(),
+ .has_aws_credential_provider = obj.has_cred_provider_type(),
+ .has_encrypted_access_keys = obj.has_encryption_info(),
+ };
+}
+
+// A provider-native credential supersedes both common static credentials and
+// provider-specific authentication fields outside the credential envelope.
+void clear_legacy_credentials(ObjectStoreInfoPB* obj) {
+ obj->clear_ak();
+ obj->clear_sk();
+ obj->clear_encryption_info();
+ obj->clear_cred_provider_type();
+ obj->clear_role_arn();
+ obj->clear_external_id();
+}
+
+std::optional<std::string> validate_required_storage_fields(const
ObjectStoreInfoPB& obj,
+ std::string_view
provider) {
+ if (obj.bucket().empty() || obj.endpoint().empty() ||
obj.region().empty()) {
+ return std::string(provider) + " storage conf requires bucket,
endpoint and region";
+ }
+ return std::nullopt;
+}
+
+std::optional<std::string> validate_gcp_obj_credential(const
ObjectStoreInfoPB& obj) {
+ const auto& credential = obj.credential().gcp_credential();
+ if (!credential.has_credential_provider_type()) {
+ return "GCP credential provider type is required";
+ }
+ if (auto error = validate_required_storage_fields(obj, "GCP");
error.has_value()) {
+ return error;
+ }
+
+ CredentialConfig config;
+ convert_obj_credential(obj, &config);
+ return validate_obj_credential_config(config,
credential_validation_context(obj));
+}
+
+} // namespace
+
+bool has_obj_credential(const ObjectStoreInfoPB& obj) {
+ return obj.has_credential();
+}
+
+std::optional<std::string> validate_obj_credential(const ObjectStoreInfoPB&
obj) {
+ if (!obj.has_credential()) {
+ return "credential is not set";
+ }
+
+ // Add one provider-specific branch for every new field in
+ // ObjectStoreCredentialPB. Reject multiple fields before dispatch when a
+ // second provider is introduced. The common visitor validates only the
+ // normalized credential itself.
+ if (obj.credential().has_gcp_credential()) {
+ return validate_gcp_obj_credential(obj);
+ }
+ return "unsupported object storage credential";
+}
+
+std::optional<std::string> validate_obj_authentication(const
ObjectStoreInfoPB& obj) {
+ if (has_obj_credential(obj)) {
+ return validate_obj_credential(obj);
+ }
+ if (obj.has_role_arn() && obj.role_arn().empty()) {
+ return "AWS role ARN cannot be empty";
+ }
+
+ const auto context = credential_validation_context(obj);
+ if (context.has_static_credentials() && context.has_aws_authentication()) {
+ return "access keys cannot be combined with AWS role or credential
provider";
+ }
+ if (!context.has_aws_authentication()) {
+ return std::nullopt;
+ }
+ if (!obj.has_provider() || obj.provider() != ObjectStoreInfoPB::S3) {
+ return "AWS role and credential provider require provider=S3";
+ }
+ if (obj.has_external_id() && obj.external_id().size() > 0 &&
obj.role_arn().empty()) {
+ return "AWS external ID requires a role ARN";
+ }
+ return std::nullopt;
+}
+
+std::optional<std::string>
validate_and_normalize_obj_credential(ObjectStoreInfoPB* obj) {
+ if (auto error = validate_obj_credential(*obj); error.has_value()) {
+ return error;
+ }
+ clear_legacy_credentials(obj);
+ return std::nullopt;
+}
+
+std::optional<std::string> apply_obj_credential(const ObjectStoreInfoPB&
update,
+ ObjectStoreInfoPB* target) {
+ if (!has_obj_credential(update)) {
+ return "credential is not set";
+ }
+ if (credential_validation_context(update).has_conflicting_credentials()) {
+ return "credentials cannot be combined with other credentials";
+ }
+
+ // Merge into a copy first so a rejected alter never partially mutates
target.
+ ObjectStoreInfoPB candidate = *target;
+ copy_obj_credential(update, &candidate);
+ if (update.credential().has_gcp_credential()) {
+ const auto& patch = update.credential().gcp_credential();
+ if (!patch.has_credential_provider_type() &&
!patch.has_impersonation_service_account()) {
+ return "GCP credential update must specify a provider type or
service account";
+ }
+ auto* credential =
candidate.mutable_credential()->mutable_gcp_credential();
+ if (target->credential().has_gcp_credential()) {
+ credential->CopyFrom(target->credential().gcp_credential());
+ credential->MergeFrom(patch);
+ }
+ // Only a new native credential needs a default; an existing source is
+ // preserved when the patch changes just the impersonation target.
+ if (!credential->has_credential_provider_type()) {
+ credential->set_credential_provider_type(GcpCredentialPB::DEFAULT);
+ }
+ if (credential->impersonation_service_account().empty()) {
Review Comment:
[P2] Keep an empty impersonation clear from replacing HMAC vault
credentials. On an existing `provider=GCP` vault with HMAC AK/SK, `ALTER
STORAGE VAULT` with only `gs.impersonation_service_account=""` reaches this
branch as a native credential patch without a provider type. It becomes DEFAULT
here, then `clear_legacy_credentials` erases the working AK/SK and the ALTER
commits. Reads now use ADC and can fail for a bucket accessible only with the
HMAC keys. Treat the empty account as a clear only when the vault already has a
native credential, or reject this patch on HMAC vaults; cover the
ALTER/readback sequence.
##########
be/src/io/fs/gcs_signed_url_provider.cpp:
##########
@@ -0,0 +1,212 @@
+// 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.
+
+#include "io/fs/gcs_signed_url_provider.h"
+
+#include <fmt/format.h>
+#include <rapidjson/document.h>
+
+#include <algorithm>
+#include <cstdlib>
+#include <exception>
+#include <string_view>
+#include <utility>
+
+#include "cpp/obj-client/auth/gcp/gcp_token_provider.h"
+#include "cpp/obj-client/auth/gcp/gcs_signed_url.h"
+#include "service/http/http_client.h"
+#include "util/string_util.h"
+#include "util/url_coding.h"
+
+namespace doris::io {
+namespace {
+
+constexpr std::string_view IAM_CREDENTIALS_ENDPOINT =
+ "https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/";
+constexpr std::string_view DEFAULT_METADATA_HOST = "metadata.google.internal";
+constexpr std::string_view METADATA_SERVICE_ACCOUNT_EMAIL_PATH =
+ "/computeMetadata/v1/instance/service-accounts/default/email";
+constexpr int64_t METADATA_REQUEST_TIMEOUT_MS = 1000;
+
+std::string iam_error_message(const rapidjson::Document& document) {
+ if (!document.IsObject() || !document.HasMember("error") ||
!document["error"].IsObject()) {
+ return "unparseable error response";
+ }
+ const auto& error = document["error"];
+ std::string message;
+ if (error.HasMember("status") && error["status"].IsString()) {
+ message = error["status"].GetString();
+ }
+ if (error.HasMember("message") && error["message"].IsString()) {
+ if (!message.empty()) {
+ message.append(": ");
+ }
+ message.append(error["message"].GetString());
+ }
+ if (message.empty()) {
+ return "error response did not contain status or message";
+ }
+ constexpr size_t MAX_ERROR_MESSAGE_SIZE = 1024;
+ if (message.size() > MAX_ERROR_MESSAGE_SIZE) {
+ message.resize(MAX_ERROR_MESSAGE_SIZE);
+ message.append("...");
+ }
+ return message;
+}
+
+bool is_unreserved(unsigned char c) {
+ return (c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c
<= '9') || c == '-' ||
+ c == '_' || c == '.' || c == '~';
+}
+
+std::string percent_encode(std::string_view value) {
+ constexpr char HEX[] = "0123456789ABCDEF";
+ std::string encoded;
+ encoded.reserve(value.size());
+ for (unsigned char c : value) {
+ if (is_unreserved(c)) {
+ encoded.push_back(static_cast<char>(c));
+ continue;
+ }
+ encoded.push_back('%');
+ encoded.push_back(HEX[c >> 4]);
+ encoded.push_back(HEX[c & 0x0F]);
+ }
+ return encoded;
+}
+
+Status call_iam_sign_blob(std::string_view access_token, std::string_view
service_account,
+ std::string_view string_to_sign, int64_t
request_timeout_ms,
+ std::string* signature) {
+ std::string encoded_payload;
+ base64_encode(std::string(string_to_sign), &encoded_payload);
+ std::string request_body = fmt::format(R"({{"payload":"{}"}})",
encoded_payload);
+ std::string endpoint =
+ std::string(IAM_CREDENTIALS_ENDPOINT) +
percent_encode(service_account) + ":signBlob";
+
+ HttpClient client;
+ // Keep the response body for non-2xx replies so IAM permission and
+ // service-account errors are actionable to operators.
+ RETURN_IF_ERROR(client.init(endpoint, false,
HttpClient::AuthTokenMode::NONE));
+ client.set_authorization("Bearer " + std::string(access_token));
+ client.set_content_type("application/json");
+ client.set_timeout_ms(request_timeout_ms > 0 ? request_timeout_ms : 10000);
+
+ std::string response;
+ RETURN_IF_ERROR(client.execute_post_request(request_body, &response));
+
+ rapidjson::Document document;
+ document.Parse(response.data(), response.size());
+ const auto http_status = client.get_http_status();
+ if (http_status < 200 || http_status >= 300) {
+ return Status::HttpError("IAM signBlob failed with HTTP {}: {}",
http_status,
+ iam_error_message(document));
+ }
+ if (document.HasParseError() || !document.IsObject() ||
!document.HasMember("signedBlob") ||
+ !document["signedBlob"].IsString()) {
+ return Status::InternalError("IAM signBlob returned an invalid
response");
+ }
+ if (!base64_decode(document["signedBlob"].GetString(), signature)) {
+ return Status::InternalError("IAM signBlob returned an invalid base64
signature");
+ }
+ return Status::OK();
+}
+
+Status fetch_metadata_service_account_email(int64_t request_timeout_ms,
std::string* email) {
+ const char* configured_host = std::getenv("GCE_METADATA_HOST");
+ std::string_view metadata_host = configured_host != nullptr &&
configured_host[0] != '\0'
+ ? configured_host
+ : DEFAULT_METADATA_HOST;
+ HttpClient client;
+ RETURN_IF_ERROR(client.init(
+ fmt::format("http://{}{}", metadata_host,
METADATA_SERVICE_ACCOUNT_EMAIL_PATH), true,
+ HttpClient::AuthTokenMode::NONE));
+ client.set_header("Metadata-Flavor", "Google");
+ client.set_timeout_ms(request_timeout_ms > 0
+ ? std::min(request_timeout_ms,
METADATA_REQUEST_TIMEOUT_MS)
+ : METADATA_REQUEST_TIMEOUT_MS);
+
+ std::string response;
+ RETURN_IF_ERROR(client.execute(&response));
+ auto resolved_email = trim(response);
+ if (!is_valid_gcp_service_account_email(resolved_email)) {
+ return Status::InternalError(
+ "GCP metadata server returned an invalid service account
email");
+ }
+ email->assign(resolved_email);
+ return Status::OK();
+}
+
+} // namespace
+
+Status generate_gcs_v4_signed_url(const GcsV4SignedUrlProviderOptions& options,
+ const GcpCredentialConfig& credential,
+ const std::shared_ptr<GcpTokenProvider>&
token_provider,
+ std::string* signed_url) {
+ if (signed_url == nullptr) {
+ return Status::InvalidArgument("signed_url output must not be null");
+ }
+ signed_url->clear();
+
+ std::string signer_email;
+ if (!credential.impersonation_service_account.empty()) {
+ signer_email = credential.impersonation_service_account;
+ } else if (credential.provider_type ==
GcpCredentialProviderType::ComputeEngine) {
+ RETURN_IF_ERROR(
+
fetch_metadata_service_account_email(options.request_timeout_ms,
&signer_email));
+ } else {
+ return Status::InvalidArgument(
+ "GCS V4 signing with DEFAULT credentials requires "
+ "gs.impersonation_service_account; use COMPUTE_ENGINE to
resolve the VM service "
Review Comment:
[P2] Resolve the signer for DEFAULT ADC on Compute Engine. When `DEFAULT`
resolves to a VM service account with `iam.serviceAccounts.signBlob`, this
branch rejects the request before fetching its email or invoking IAM, although
`COMPUTE_ENGINE` on the same VM takes that path. `RuntimeState` has already
uploaded the load error log, but returns a BE-local link instead of a GCS
signed link, which can become inaccessible after BE recycling or when its web
port is private. Resolve the ADC service-account identity for this case and
cover the load error-log flow.
##########
fe/fe-connector/fe-connector-iceberg/src/main/java/org/apache/doris/connector/iceberg/IcebergCatalogFactory.java:
##########
@@ -419,6 +422,33 @@ public static Map<String, String>
buildCatalogProperties(IcebergCatalogPropertie
// s3tables: bespoke instantiation. Preserve the skeleton's
base+impl routing.
break;
}
+ chosenS3.ifPresent(storage ->
storage.toBackendProperties().ifPresent(backend -> {
+ Optional<GcsAuth> auth = GcsAuthResolver.resolve(backend.toMap());
+ if (auth.filter(GcsAuth::isAnonymous).isPresent()) {
+ opts.put(AwsClientProperties.CLIENT_CREDENTIALS_PROVIDER,
+
"software.amazon.awssdk.auth.credentials.AnonymousCredentialsProvider");
+ }
+ auth.flatMap(GcsAuth::getNativeCredential).ifPresent(credential ->
{
+ // HadoopCatalog resolves its namespace filesystem
independently of S3FileIO.
+ // Native GCS credentials configure fs.gs.*, so compatibility
warehouse schemes
+ // must select that same filesystem before HadoopCatalog
initializes.
+ String warehouse =
opts.get(CatalogProperties.WAREHOUSE_LOCATION);
+ if (IcebergCatalogProperties.TYPE_HADOOP.equals(flavor) &&
warehouse != null) {
+ if (warehouse.regionMatches(true, 0, "s3://", 0, 5)) {
+ opts.put(CatalogProperties.WAREHOUSE_LOCATION, "gs://"
+ warehouse.substring(5));
+ } else if (warehouse.regionMatches(true, 0, "s3a://", 0,
6)) {
+ opts.put(CatalogProperties.WAREHOUSE_LOCATION, "gs://"
+ warehouse.substring(6));
+ }
+ }
+ putS3FileIODialect(opts, storage);
+ opts.put("provider", "GCP");
+ opts.put(GcpCredential.CREDENTIAL_PROVIDER_TYPE,
credential.getCredentialProviderType().name());
+ putIfNotBlank(opts,
GcpCredential.IMPERSONATION_SERVICE_ACCOUNT,
+ credential.getImpersonationServiceAccount());
+ opts.put(CatalogProperties.FILE_IO_IMPL,
"org.apache.iceberg.aws.s3.S3FileIO");
+ opts.put(S3FileIOProperties.CLIENT_FACTORY,
GcpS3FileIOAwsClientFactory.class.getName());
+ });
+ }));
Review Comment:
[P1] Make the native GCP Iceberg client factory available to the BE metadata
scanner. A `table$files` or `table$data_files` scan serializes an Iceberg
manifest task containing the table's S3FileIO. On the BE, `rows()` reopens the
manifest and Iceberg resolves this configured factory using the isolated
`iceberg-metadata-scanner` classloader, but that plugin contains neither
`GcpS3FileIOAwsClientFactory` nor its FE GCS dependencies. The scan fails
loading the factory even when ordinary table scans work. Supply a BE-loadable
factory with its dependencies or reconstruct scanner-local FileIO, and cover a
native GCP metadata-table scan.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]