This is an automated email from the ASF dual-hosted git repository. diqiu50 pushed a commit to branch dell-1.3 in repository https://gitbox.apache.org/repos/asf/gravitino.git
commit cef52ce77b2e5c250575f11b22ba4a8bd0833541 Author: diqiu50 <[email protected]> AuthorDate: Mon Aug 24 22:30:39 2026 +0800 [Cherry-pick to branch-1.3] [#12554] improvement(trino-connector): Address more review findings on the Iceberg REST routing - Reload a catalog when its discovered Iceberg REST URI changes, not just when Gravitino's lastModifiedTime advances - Scope the manual gravitino.iceberg.rest-uri override to a metalake so it no longer misroutes other metalakes in multi-metalake mode - Fix the deprecated aux-service config WARN firing on every discovery poll instead of once per key - Redact secrets embedded in the catalog config JSON, not just top-level SQL assignments, when logging CREATE CATALOG - Let the shared trino-connector module pick its own JDK toolchain by trinoVersion so minSupportedTrinoVersion can track the actual minimum supported version - Dedupe the aux-service prefix constant duplicated across two IT classes - Clarify a misleading comment about HashMap put-order priority (cherry picked from commit 3ec12a854b008c117bf06692b5c61e5e5395f1a2) --- build.gradle.kts | 11 +-- .../client/TestIcebergRestServiceUri.java | 82 +++++++++++++++++++ .../auxiliary/AuxiliaryServiceManager.java | 15 +++- docs/trino-connector/configuration.md | 5 ++ gradle.properties | 2 +- .../gravitino/integration/test/util/BaseIT.java | 4 +- .../web/rest/TestIcebergRESTServiceOperations.java | 10 +++ .../integration/test/TrinoConnectorIT.java | 2 - .../integration/test/TrinoQueryITBase.java | 7 +- trino-connector/trino-connector/build.gradle.kts | 30 +++++++ .../gravitino/trino/connector/GravitinoConfig.java | 25 ++++-- .../connector/catalog/CatalogConnectorManager.java | 8 +- .../trino/connector/catalog/CatalogRegister.java | 23 ++++++ .../iceberg/IcebergCatalogPropertyConverter.java | 6 +- .../catalog/iceberg/IcebergConnectorAdapter.java | 25 +++++- .../trino/connector/TestGravitinoConfig.java | 27 +++++- .../connector/catalog/TestCatalogRegister.java | 18 ++++ .../iceberg/TestIcebergConnectorAdapterReload.java | 95 ++++++++++++++++++++++ 18 files changed, 360 insertions(+), 35 deletions(-) diff --git a/build.gradle.kts b/build.gradle.kts index 08a0128000..ac214df4e9 100644 --- a/build.gradle.kts +++ b/build.gradle.kts @@ -416,16 +416,7 @@ subprojects { java { toolchain { - // Some JDK vendors like Homebrew installed OpenJDK 17 have problems in building trino-connector: - // It will cause tests of Trino-connector hanging forever on macOS, to avoid this issue and - // other vendor-related problems, Gravitino will use the specified AMAZON OpenJDK 17 to build - // Trino-connector on macOS. - if (project.name == "trino-connector") { - if (OperatingSystem.current().isMacOsX) { - vendor.set(JvmVendorSpec.AMAZON) - } - languageVersion.set(JavaLanguageVersion.of(17)) - } else if (compatibleWithJDK8(project)) { + if (compatibleWithJDK8(project)) { languageVersion.set(JavaLanguageVersion.of(17)) sourceCompatibility = JavaVersion.VERSION_1_8 targetCompatibility = JavaVersion.VERSION_1_8 diff --git a/clients/client-java/src/test/java/org/apache/gravitino/client/TestIcebergRestServiceUri.java b/clients/client-java/src/test/java/org/apache/gravitino/client/TestIcebergRestServiceUri.java new file mode 100644 index 0000000000..2186118e99 --- /dev/null +++ b/clients/client-java/src/test/java/org/apache/gravitino/client/TestIcebergRestServiceUri.java @@ -0,0 +1,82 @@ +/* + * 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. + */ +package org.apache.gravitino.client; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.google.common.collect.ImmutableMap; +import java.util.Optional; +import org.apache.gravitino.dto.responses.ErrorResponse; +import org.apache.gravitino.dto.responses.IcebergRESTServiceResponse; +import org.apache.hc.core5.http.HttpStatus; +import org.apache.hc.core5.http.Method; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +public class TestIcebergRestServiceUri extends TestBase { + + private static final String ICEBERG_REST_PATH = withSlash("api/system/iceberg-rest"); + + @Test + public void testIcebergRestServiceUriPresent() throws JsonProcessingException { + IcebergRESTServiceResponse resp = + new IcebergRESTServiceResponse("http://irc-host:9001/iceberg"); + buildMockResource( + Method.GET, + ICEBERG_REST_PATH, + ImmutableMap.of("metalake", "test"), + null, + resp, + HttpStatus.SC_OK); + + Optional<String> uri = client.icebergRestServiceUri("test"); + + Assertions.assertTrue(uri.isPresent()); + Assertions.assertEquals("http://irc-host:9001/iceberg", uri.get()); + } + + @Test + public void testIcebergRestServiceUriAbsent() throws JsonProcessingException { + IcebergRESTServiceResponse resp = new IcebergRESTServiceResponse(null); + buildMockResource( + Method.GET, + ICEBERG_REST_PATH, + ImmutableMap.of("metalake", "test"), + null, + resp, + HttpStatus.SC_OK); + + Optional<String> uri = client.icebergRestServiceUri("test"); + + Assertions.assertFalse(uri.isPresent()); + } + + @Test + public void testIcebergRestServiceUriPropagatesServerError() throws JsonProcessingException { + ErrorResponse errResp = ErrorResponse.internalError("internal error"); + buildMockResource( + Method.GET, + ICEBERG_REST_PATH, + ImmutableMap.of("metalake", "test"), + null, + errResp, + HttpStatus.SC_INTERNAL_SERVER_ERROR); + + Assertions.assertThrows(RuntimeException.class, () -> client.icebergRestServiceUri("test")); + } +} diff --git a/core/src/main/java/org/apache/gravitino/auxiliary/AuxiliaryServiceManager.java b/core/src/main/java/org/apache/gravitino/auxiliary/AuxiliaryServiceManager.java index 83971c9953..9c67a93583 100644 --- a/core/src/main/java/org/apache/gravitino/auxiliary/AuxiliaryServiceManager.java +++ b/core/src/main/java/org/apache/gravitino/auxiliary/AuxiliaryServiceManager.java @@ -32,6 +32,8 @@ import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.ServiceLoader; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.function.Consumer; import java.util.stream.Collectors; import org.apache.commons.lang3.StringUtils; @@ -53,6 +55,9 @@ public class AuxiliaryServiceManager { private static final Splitter splitter = Splitter.on(","); private static final Joiner DOT = Joiner.on("."); + // Tracks which deprecated keys have already been warned about, since getAuxServiceConfig + // re-extracts the config on every call. + private static final Set<String> warnedDeprecatedAuxServiceKeys = ConcurrentHashMap.newKeySet(); private Map<String, GravitinoAuxiliaryService> auxServices = new HashMap<>(); private Map<String, IsolatedClassLoader> auxServiceClassLoaders = new HashMap<>(); @@ -268,10 +273,12 @@ public class AuxiliaryServiceManager { .forEach( entry -> { String extractKey = entry.getKey().substring(GRAVITINO_AUX_SERVICE_PREFIX.length()); - LOG.warn( - "The configuration {} is deprecated(still working), please use gravitino.{} instead.", - entry.getKey(), - extractKey); + if (warnedDeprecatedAuxServiceKeys.add(entry.getKey())) { + LOG.warn( + "The configuration {} is deprecated(still working), please use gravitino.{} instead.", + entry.getKey(), + extractKey); + } serviceConfigs.put(extractKey, entry.getValue()); }); splitter diff --git a/docs/trino-connector/configuration.md b/docs/trino-connector/configuration.md index 9e6747ef67..7e6afd07f7 100644 --- a/docs/trino-connector/configuration.md +++ b/docs/trino-connector/configuration.md @@ -42,6 +42,11 @@ before. See [Iceberg catalog](./catalog-iceberg.md#how-trino-reaches-the-catalog Multi-metalake mode (`gravitino.use-single-metalake=false`) is supported on Trino connector versions 440-445 and 469-478. On versions 446-468, a warning is logged and the connector initializes, but the mode is not fully supported and some operations may fail. +**Note:** In multi-metalake mode, `gravitino.iceberg.rest-uri` is only honored when scoped to a +metalake, as `gravitino.iceberg.rest-uri.<metalake_name>` — the unscoped form is ignored, since a +single Iceberg REST server serves exactly one metalake and applying it to every metalake would +misroute the others. The unscoped form remains valid in single-metalake mode. + ## Connecting to a TLS-enabled coordinator The Gravitino Trino connector registers catalogs by connecting back to the Trino coordinator over diff --git a/gradle.properties b/gradle.properties index 4c1ad7e606..05f7d7d451 100644 --- a/gradle.properties +++ b/gradle.properties @@ -48,4 +48,4 @@ skipDockerTests = true enableFuse = false # The minimum supported Trino version. -minSupportedTrinoVersion= 435 +minSupportedTrinoVersion= 440 diff --git a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java index ef45535de9..6a4e41c935 100644 --- a/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java +++ b/integration-test-common/src/test/java/org/apache/gravitino/integration/test/util/BaseIT.java @@ -122,6 +122,8 @@ public class BaseIT { public static final String DOWNLOAD_SQLITE_JDBC_DRIVER_URL = "https://repo1.maven.org/maven2/org/xerial/sqlite-jdbc/3.42.0.0/sqlite-jdbc-3.42.0.0.jar"; + public static final String GRAVITINO_ICEBERG_REST_PREFIX = "gravitino.iceberg-rest."; + public static final Map<String, Pattern> SUPPORTED_CLEAN_CONFLICTS_DRIVER_TYPES = ImmutableMap.of( "mysql", Pattern.compile("mysql-connector-java-([\\d.]+)\\.jar"), @@ -688,7 +690,7 @@ public class BaseIT { protected String getIcebergRestServiceUri() { JettyServerConfig jettyServerConfig = - JettyServerConfig.fromConfig(serverConfig, String.format("gravitino.iceberg-rest.")); + JettyServerConfig.fromConfig(serverConfig, GRAVITINO_ICEBERG_REST_PREFIX); return String.format( "http://%s:%d/iceberg/", jettyServerConfig.getHost(), jettyServerConfig.getHttpPort()); } diff --git a/server/src/test/java/org/apache/gravitino/server/web/rest/TestIcebergRESTServiceOperations.java b/server/src/test/java/org/apache/gravitino/server/web/rest/TestIcebergRESTServiceOperations.java index f4bb9eddb7..5a559bf059 100644 --- a/server/src/test/java/org/apache/gravitino/server/web/rest/TestIcebergRESTServiceOperations.java +++ b/server/src/test/java/org/apache/gravitino/server/web/rest/TestIcebergRESTServiceOperations.java @@ -134,6 +134,16 @@ public class TestIcebergRESTServiceOperations { assertEquals("http://irc-host:19001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } + @Test + public void testMalformedPortFallsBackToDefaultPort() { + IcebergRESTServiceOperations ops = + newOps( + true, + withDynamicProvider(ImmutableMap.of("host", "irc-host", "httpPort", "not-a-number")), + "gravitino-host"); + assertEquals("http://irc-host:9001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); + } + @Test public void testWildcardHostFallsBackToRequestServerName() { IcebergRESTServiceOperations ops = diff --git a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java index 968e26fb30..7b04f21da6 100644 --- a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java +++ b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoConnectorIT.java @@ -77,8 +77,6 @@ public class TrinoConnectorIT extends BaseIT { private static final ContainerSuite containerSuite = ContainerSuite.getInstance(); - private static final String GRAVITINO_ICEBERG_REST_PREFIX = "gravitino.iceberg-rest."; - public static String metalakeName = GravitinoITUtils.genRandomName("TrinoIT_metalake").toLowerCase(); public static String catalogName = diff --git a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java index e5be0c11c4..405df92076 100644 --- a/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java +++ b/trino-connector/integration-test/src/test/java/org/apache/gravitino/trino/connector/integration/test/TrinoQueryITBase.java @@ -47,8 +47,6 @@ import org.slf4j.LoggerFactory; public class TrinoQueryITBase { private static final Logger LOG = LoggerFactory.getLogger(TrinoQueryITBase.class); - private static final String GRAVITINO_ICEBERG_REST_PREFIX = "gravitino.iceberg-rest."; - // Auto start docker containers and Gravitino server protected static boolean autoStart = true; @@ -95,9 +93,10 @@ public class TrinoQueryITBase { // so the auxiliary service has to run and serve this test's metalake. baseIT.enableIcebergAuxRestService( ImmutableMap.of( - GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.ICEBERG_REST_CATALOG_CONFIG_PROVIDER, + BaseIT.GRAVITINO_ICEBERG_REST_PREFIX + + IcebergConstants.ICEBERG_REST_CATALOG_CONFIG_PROVIDER, IcebergConstants.DYNAMIC_ICEBERG_CATALOG_CONFIG_PROVIDER_NAME, - GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.GRAVITINO_METALAKE, + BaseIT.GRAVITINO_ICEBERG_REST_PREFIX + IcebergConstants.GRAVITINO_METALAKE, metalakeName)); baseIT.startIntegrationTest(); gravitinoClient = baseIT.getGravitinoClient(); diff --git a/trino-connector/trino-connector/build.gradle.kts b/trino-connector/trino-connector/build.gradle.kts index 69ed4e7593..898fc62a4d 100644 --- a/trino-connector/trino-connector/build.gradle.kts +++ b/trino-connector/trino-connector/build.gradle.kts @@ -17,6 +17,9 @@ * under the License. */ +import net.ltgt.gradle.errorprone.errorprone +import org.gradle.internal.os.OperatingSystem + plugins { id("java") id("idea") @@ -30,6 +33,33 @@ val minSupportedTrinoVersionProperty = providers.gradleProperty("minSupportedTri val trinoVersionProperty = providers.gradleProperty("trinoVersion").orElse(minSupportedTrinoVersionProperty) val trinoVersion = trinoVersionProperty.map { it.trim().toInt() }.get() +if (trinoVersion >= 440) { + // Trino 440+'s trino-spi is compiled for JDK 21+, so this module needs the same JDK 24 + // toolchain the versioned trino-connector-<range> modules use for the same trinoVersion. + // Error Prone is incompatible with that toolchain, so it is disabled here too, matching those + // modules' own override. + java { + toolchain.languageVersion.set(JavaLanguageVersion.of(24)) + } + tasks.withType<JavaCompile>().configureEach { + options.errorprone.isEnabled.set(false) + options.release.set(17) + } +} else { + java { + toolchain { + // Some JDK vendors like Homebrew installed OpenJDK 17 have problems in building + // trino-connector: It will cause tests of Trino-connector hanging forever on macOS, to + // avoid this issue and other vendor-related problems, Gravitino will use the specified + // AMAZON OpenJDK 17 to build Trino-connector on macOS. + if (OperatingSystem.current().isMacOsX) { + vendor.set(JvmVendorSpec.AMAZON) + } + languageVersion.set(JavaLanguageVersion.of(17)) + } + } +} + dependencies { implementation(project(":catalogs:catalog-common")) implementation(project(":clients:client-java-runtime", configuration = "shadow")) diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java index b86fa5c5dc..5cdeba701e 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConfig.java @@ -742,15 +742,28 @@ public class GravitinoConfig { } /** - * Retrieves the manually configured {@code gravitino.iceberg.rest-uri}, if any. Unlike the - * discovered endpoint, this is plain local file configuration and is therefore identical and - * valid on every node — coordinator and workers alike. + * Retrieves the manually configured Iceberg REST server endpoint for the given metalake, if any. + * Unlike the discovered endpoint, this is plain local file configuration and is therefore + * identical and valid on every node — coordinator and workers alike. * + * <p>{@code gravitino.iceberg.rest-uri.<metalake>} is checked first. The unscoped {@code + * gravitino.iceberg.rest-uri} is honored only in single-metalake mode, where it is unambiguous; + * in multi-metalake mode it is ignored, since a single Iceberg REST server serves exactly one + * metalake and applying it to every metalake would misroute the others. + * + * @param metalake the metalake to resolve the override for * @return the manually configured Iceberg REST server endpoint, or an empty string when unset */ - public String getManualIcebergRestUri() { - return config.getOrDefault( - GRAVITINO_ICEBERG_REST_URI.key, GRAVITINO_ICEBERG_REST_URI.defaultValue); + public String getManualIcebergRestUri(String metalake) { + String scopedValue = config.get(GRAVITINO_ICEBERG_REST_URI.key + "." + metalake); + if (StringUtils.isNotBlank(scopedValue)) { + return scopedValue; + } + if (singleMetalakeMode()) { + return config.getOrDefault( + GRAVITINO_ICEBERG_REST_URI.key, GRAVITINO_ICEBERG_REST_URI.defaultValue); + } + return GRAVITINO_ICEBERG_REST_URI.defaultValue; } /** diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java index ce9dc5b7da..b1a809ca7b 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorManager.java @@ -41,6 +41,7 @@ import org.apache.gravitino.client.GravitinoMetalake; import org.apache.gravitino.exceptions.NoSuchMetalakeException; import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; +import org.apache.gravitino.trino.connector.catalog.iceberg.IcebergConnectorAdapter; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; import org.apache.gravitino.trino.connector.security.GravitinoAuthProvider; import org.slf4j.Logger; @@ -329,7 +330,12 @@ public class CatalogConnectorManager { private void reloadCatalog(GravitinoCatalog catalog) { String catalogFullName = getTrinoCatalogName(catalog); GravitinoCatalog oldCatalog = catalogConnectors.get(catalogFullName).getCatalog(); - if (catalog.getLastModifiedTime() <= oldCatalog.getLastModifiedTime()) { + // The discovered Iceberg REST endpoint is embedded into the catalog independently of + // Gravitino's own lastModifiedTime, so it is checked as a separate reload trigger. + boolean icebergRestUriChanged = + IcebergConnectorAdapter.hasDiscoveredIcebergRestUriChanged(catalog, oldCatalog, config); + if (catalog.getLastModifiedTime() <= oldCatalog.getLastModifiedTime() + && !icebergRestUriChanged) { return; } diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java index 4189a251cc..fd971c3bce 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogRegister.java @@ -72,6 +72,14 @@ public class CatalogRegister { "\"([^\"]*(?:credential|token|secret|password)[^\"]*)\"\\s*=\\s*'([^']*)'", Pattern.CASE_INSENSITIVE); + // Matches "key":"value" style secret assignments inside the serialized GravitinoCatalog JSON + // that GRAVITINO_DYNAMIC_CONNECTOR_CATALOG_CONFIG carries (e.g. "jdbc-password":"..." or + // "s3-secret-key":"..."), which SECRET_PROPERTY_PATTERN's SQL-assignment shape does not match. + private static final Pattern SECRET_JSON_PROPERTY_PATTERN = + Pattern.compile( + "\"([^\"]*(?:credential|token|secret|password)[^\"]*)\"\\s*:\\s*\"([^\"]*)\"", + Pattern.CASE_INSENSITIVE); + private Connection connection; private boolean isStarted = false; private String catalogStoreDirectory; @@ -342,6 +350,10 @@ public class CatalogRegister { @VisibleForTesting static String redactSecrets(String createCatalogCommand) { + return redactJsonSecrets(redactSqlSecrets(createCatalogCommand)); + } + + private static String redactSqlSecrets(String createCatalogCommand) { Matcher matcher = SECRET_PROPERTY_PATTERN.matcher(createCatalogCommand); StringBuffer redacted = new StringBuffer(); while (matcher.find()) { @@ -352,6 +364,17 @@ public class CatalogRegister { return redacted.toString(); } + private static String redactJsonSecrets(String createCatalogCommand) { + Matcher matcher = SECRET_JSON_PROPERTY_PATTERN.matcher(createCatalogCommand); + StringBuffer redacted = new StringBuffer(); + while (matcher.find()) { + matcher.appendReplacement( + redacted, Matcher.quoteReplacement("\"" + matcher.group(1) + "\":\"***\"")); + } + matcher.appendTail(redacted); + return redacted.toString(); + } + private String generateDropCatalogCommand(String name) { return String.format("DROP CATALOG %s", name); } diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java index 14aed17709..ec0237b3b7 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergCatalogPropertyConverter.java @@ -131,7 +131,8 @@ public class IcebergCatalogPropertyConverter extends CatalogPropertyConverter { throw new UnsupportedOperationException("Unsupported backend type: " + backend); } Map<String, String> config = new HashMap<>(); - // The order of put operations determines the priority of parameters. + // Later put/putAll calls override earlier ones for the same key; this is call-order + // precedence, unrelated to HashMap's (unspecified) iteration order. config.putAll(super.gravitinoToEngineProperties(properties)); config.putAll(stringStringMap); config.put("fs.hadoop.enabled", "true"); @@ -169,7 +170,8 @@ public class IcebergCatalogPropertyConverter extends CatalogPropertyConverter { } Map<String, String> config = new HashMap<>(); - // The order of put operations determines the priority of parameters. + // Later put/putAll calls override earlier ones for the same key; this is call-order + // precedence, unrelated to HashMap's (unspecified) iteration order. config.putAll(buildStorageProperties(catalog.getProperties())); config.put(TRINO_ICEBERG_REST_VENDED_CREDENTIALS, "true"); if (gravitinoConfig.isForwardUser()) { diff --git a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java index 85df20d944..44c06503f8 100644 --- a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java +++ b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/iceberg/IcebergConnectorAdapter.java @@ -88,7 +88,7 @@ public class IcebergConnectorAdapter implements CatalogConnectorAdapter { // discovered endpoint is coordinator-only knowledge, so it is read from the catalog's own // properties, where the coordinator embeds it at registration time (see // embedDiscoveredIcebergRestUri), rather than from GravitinoConfig directly. - String restUri = config.getManualIcebergRestUri(); + String restUri = config.getManualIcebergRestUri(catalog.getMetalake()); if (StringUtils.isBlank(restUri)) { restUri = catalog.getProperty(DISCOVERED_ICEBERG_REST_URI_PROPERTY, ""); } @@ -147,6 +147,29 @@ public class IcebergConnectorAdapter implements CatalogConnectorAdapter { catalog.getLastModifiedTime()); } + /** + * Returns whether the Iceberg REST server endpoint currently discovered for {@code catalog}'s + * metalake differs from the one embedded into {@code oldCatalog} the last time it was registered. + * The catalog's Gravitino {@code lastModifiedTime} does not advance when only Iceberg REST + * availability changes, so callers that use it to decide whether to reload a catalog need this as + * a separate signal. + * + * @param catalog the freshly loaded catalog + * @param oldCatalog the previously registered version of the same catalog + * @param config the connector configuration holding the discovered endpoints + * @return {@code true} if {@code catalog} is a lakehouse-iceberg catalog and the discovered + * endpoint for its metalake differs from what {@code oldCatalog} was registered with + */ + public static boolean hasDiscoveredIcebergRestUriChanged( + GravitinoCatalog catalog, GravitinoCatalog oldCatalog, GravitinoConfig config) { + if (!ICEBERG_PROVIDER.equals(catalog.getProvider())) { + return false; + } + String currentUri = config.getDiscoveredIcebergRestUri(catalog.getMetalake()); + String previousUri = oldCatalog.getProperty(DISCOVERED_ICEBERG_REST_URI_PROPERTY, ""); + return !currentUri.equals(previousUri); + } + @Override public String internalConnectorName() { return CONNECTOR_ICEBERG; diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java index 5d03d2e507..d2792d75df 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConfig.java @@ -305,7 +305,7 @@ public class TestGravitinoConfig { GravitinoConfig config = new GravitinoConfig(ImmutableMap.of("gravitino.metalake", "user_001")); // Nothing configured, and nothing discovered yet. - assertEquals("", config.getManualIcebergRestUri()); + assertEquals("", config.getManualIcebergRestUri("user_001")); assertEquals("", config.getDiscoveredIcebergRestUri("user_001")); assertTrue(config.getIcebergRestCatalogConfig().isEmpty()); } @@ -324,8 +324,8 @@ public class TestGravitinoConfig { "client_id:client_secret"); GravitinoConfig config = new GravitinoConfig(configMap); - // A manually configured URI is plain local config, not scoped to any metalake. - assertEquals("http://127.0.0.1:9001/iceberg", config.getManualIcebergRestUri()); + // The unscoped URI is honored as-is in single-metalake mode. + assertEquals("http://127.0.0.1:9001/iceberg", config.getManualIcebergRestUri("user_001")); Map<String, String> restCatalogConfig = config.getIcebergRestCatalogConfig(); assertEquals(2, restCatalogConfig.size()); @@ -334,6 +334,27 @@ public class TestGravitinoConfig { "client_id:client_secret", restCatalogConfig.get("iceberg.rest-catalog.oauth2.credential")); } + @Test + public void testIcebergRestConfigScopedToMetalakeInMultiMetalakeMode() { + ImmutableMap<String, String> configMap = + ImmutableMap.of( + "gravitino.metalake", + "metalake_a", + "gravitino.use-single-metalake", + "false", + "gravitino.iceberg.rest-uri", + "http://unscoped:9001/iceberg", + "gravitino.iceberg.rest-uri.metalake_a", + "http://metalake-a:9001/iceberg"); + GravitinoConfig config = new GravitinoConfig(configMap); + + // The scoped key wins for the metalake it names. + assertEquals("http://metalake-a:9001/iceberg", config.getManualIcebergRestUri("metalake_a")); + // The unscoped key is ignored in multi-metalake mode, since it would otherwise misroute every + // metalake other than the one the Iceberg REST server actually serves. + assertEquals("", config.getManualIcebergRestUri("metalake_b")); + } + @Test public void testDiscoveredIcebergRestUriIsPerMetalake() { GravitinoConfig config = new GravitinoConfig(ImmutableMap.of("gravitino.metalake", "user_001")); diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogRegister.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogRegister.java index 488cfc7ed6..e4248de016 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogRegister.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogRegister.java @@ -465,4 +465,22 @@ public class TestCatalogRegister { assertTrue( redacted.contains("\"gravitino.iceberg.rest-catalog.uri\"='http://irc-host:9001/iceberg'")); } + + @Test + public void testRedactSecretsMasksJsonEmbeddedSecrets() { + String command = + "CREATE CATALOG c USING gravitino WITH ( " + + "\"__gravitino.dynamic.connector.catalog.config\"=" + + "'{\"name\":\"hive_catalog\",\"properties\":" + + "{\"jdbc-password\":\"hunter2\",\"s3-secret-key\":\"abc123\",\"jdbc-user\":\"admin\"}}')"; + + String redacted = CatalogRegister.redactSecrets(command); + + assertFalse(redacted.contains("hunter2")); + assertFalse(redacted.contains("abc123")); + assertTrue(redacted.contains("\"jdbc-password\":\"***\"")); + assertTrue(redacted.contains("\"s3-secret-key\":\"***\"")); + // Non-secret properties must survive redaction unchanged. + assertTrue(redacted.contains("\"jdbc-user\":\"admin\"")); + } } diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergConnectorAdapterReload.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergConnectorAdapterReload.java new file mode 100644 index 0000000000..d90057f24d --- /dev/null +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergConnectorAdapterReload.java @@ -0,0 +1,95 @@ +/* + * 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. + */ +package org.apache.gravitino.trino.connector.catalog.iceberg; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import com.google.common.collect.ImmutableMap; +import org.apache.gravitino.trino.connector.GravitinoConfig; +import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; +import org.junit.jupiter.api.Test; + +public class TestIcebergConnectorAdapterReload { + + private static GravitinoConfig configFor(String metalake, String discoveredUri) { + GravitinoConfig config = + new GravitinoConfig( + ImmutableMap.of( + "gravitino.uri", "http://127.0.0.1:8090", "gravitino.metalake", metalake)); + if (discoveredUri != null) { + config.setDiscoveredIcebergRestUri(metalake, discoveredUri); + } + return config; + } + + private static GravitinoCatalog icebergCatalog(String embeddedUri) { + ImmutableMap<String, String> properties = + embeddedUri == null + ? ImmutableMap.of() + : ImmutableMap.of("__gravitino.iceberg.rest-uri", embeddedUri); + return new GravitinoCatalog("test", "lakehouse-iceberg", "iceberg_catalog", properties, 5L); + } + + @Test + public void testUnchangedWhenDiscoveredUriMatchesEmbeddedUri() { + GravitinoConfig config = configFor("test", "http://irc-host:9001/iceberg"); + GravitinoCatalog freshCatalog = icebergCatalog(null); + GravitinoCatalog registeredCatalog = icebergCatalog("http://irc-host:9001/iceberg"); + + assertFalse( + IcebergConnectorAdapter.hasDiscoveredIcebergRestUriChanged( + freshCatalog, registeredCatalog, config)); + } + + @Test + public void testChangedWhenIcebergRestBecomesAvailable() { + GravitinoConfig config = configFor("test", "http://irc-host:9001/iceberg"); + GravitinoCatalog freshCatalog = icebergCatalog(null); + GravitinoCatalog registeredCatalog = icebergCatalog(null); + + assertTrue( + IcebergConnectorAdapter.hasDiscoveredIcebergRestUriChanged( + freshCatalog, registeredCatalog, config)); + } + + @Test + public void testChangedWhenIcebergRestBecomesUnavailable() { + GravitinoConfig config = configFor("test", null); + GravitinoCatalog freshCatalog = icebergCatalog(null); + GravitinoCatalog registeredCatalog = icebergCatalog("http://irc-host:9001/iceberg"); + + assertTrue( + IcebergConnectorAdapter.hasDiscoveredIcebergRestUriChanged( + freshCatalog, registeredCatalog, config)); + } + + @Test + public void testNonIcebergCatalogNeverReportsChanged() { + GravitinoConfig config = configFor("test", "http://irc-host:9001/iceberg"); + GravitinoCatalog freshCatalog = + new GravitinoCatalog("test", "hive", "hive_catalog", ImmutableMap.of(), 5L); + GravitinoCatalog registeredCatalog = + new GravitinoCatalog("test", "hive", "hive_catalog", ImmutableMap.of(), 5L); + + assertFalse( + IcebergConnectorAdapter.hasDiscoveredIcebergRestUriChanged( + freshCatalog, registeredCatalog, config)); + } +}
