github-actions[bot] commented on code in PR #68540:
URL: https://github.com/apache/doris/pull/68540#discussion_r4236178620


##########
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:
   [P2] Preserve scheme-aware Iceberg FileIO selection. This assignment 
overwrites an explicit io-impl and also replaces the Iceberg 1.11 REST default 
ResolvingFileIO. For a native-GCP REST catalog containing gs:// and hdfs:// 
tables, an HDFS FileIO location is parsed as an S3 bucket/key and sent to the 
GCS S3 client, so operations that read that location fail. Preserve the 
selected FileIO for non-GCS locations and cover a mixed-location catalog.



##########
be/src/io/fs/gcs_signed_url_provider.cpp:
##########
@@ -0,0 +1,189 @@
+// 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 <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 "cpp/sync_point.h"
+#include "service/http/http_client.h"
+#include "util/url_coding.h"
+
+namespace doris::io {
+namespace {
+
+constexpr std::string_view IAM_CREDENTIALS_ENDPOINT =
+        "https://iamcredentials.googleapis.com/v1/projects/-/serviceAccounts/";;
+
+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) {
+    TEST_SYNC_POINT_RETURN_WITH_VALUE("GcsV4Signer::sign_blob", Status::OK(), 
access_token,
+                                      service_account, string_to_sign, 
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));

Review Comment:
   [P2] Pass the configured CA to IAM signBlob. When an outbound HTTPS proxy 
uses a private CA supplied only through ca_cert_file_paths, the GCS object 
client and OAuth token provider trust it, but this HttpClient uses libcurl 
default CAs. Uploading a load error log can succeed while signBlob fails TLS 
verification, so the returned error-log path stays BE-local instead of becoming 
a usable URL. Use the selected CA for this request too.



-- 
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]

Reply via email to