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 fcbcfb2ede24512d6aad6e4093cbfce02cba93a7 Author: diqiu50 <[email protected]> AuthorDate: Mon Aug 24 20:44:41 2026 +0800 [Cherry-pick to branch-1.3] [#12554] improvement(trino-connector): Fix IRC discovery misrouting on static-config-provider and other review findings - Only report the Iceberg REST server's endpoint when it runs the dynamic config provider; the default static provider serves catalogs unrelated to Gravitino catalog names, so reporting it there would misroute every Iceberg catalog after upgrade. - Log discovery failures at WARN on state transitions instead of always at DEBUG, so a persistent failure is visible without spamming logs for the expected old-server case. - Honor the deprecated gravitino.auxService.iceberg-rest.* config form when resolving the reported host/port, matching what the Iceberg REST server itself sees. - Redact credential/token/secret/password-looking properties before logging the CREATE CATALOG statement. - Add log lines for the previously-silent provider/metalake-mismatch null-return paths. - Add test coverage for CatalogConnectorManager's discovery polling and CatalogRegister's embed-then-serialize wiring, plus a couple of gaps around manual-vs-discovered precedence and the non-Iceberg guard. (cherry picked from commit b22c372e821359002f009d0fb280274893161f5f) --- .../auxiliary/AuxiliaryServiceManager.java | 18 +++++ .../auxiliary/TestAuxiliaryServiceManager.java | 39 +++++++++ .../web/rest/IcebergRESTServiceOperations.java | 62 +++++++++++---- .../web/rest/TestIcebergRESTServiceOperations.java | 93 ++++++++++++++-------- .../connector/catalog/CatalogConnectorManager.java | 26 +++++- .../trino/connector/catalog/CatalogRegister.java | 36 ++++++++- .../catalog/TestCatalogConnectorManager.java | 53 ++++++++++++ .../connector/catalog/TestCatalogRegister.java | 88 ++++++++++++++++++++ .../TestIcebergCatalogPropertyConverter.java | 40 ++++++++++ 9 files changed, 401 insertions(+), 54 deletions(-) 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 c1c0c14235..83971c9953 100644 --- a/core/src/main/java/org/apache/gravitino/auxiliary/AuxiliaryServiceManager.java +++ b/core/src/main/java/org/apache/gravitino/auxiliary/AuxiliaryServiceManager.java @@ -109,6 +109,24 @@ public class AuxiliaryServiceManager { return IsolatedClassLoader.buildClassLoader(classPaths); } + /** + * Returns the effective configuration for one auxiliary service, with keys stripped of the {@code + * gravitino.<name>.} prefix (or the deprecated {@code gravitino.auxService.<name>.} prefix, which + * is still honored) — the same map the service itself receives in {@link + * GravitinoAuxiliaryService#serviceInit}. Callers outside this service (e.g. reporting an + * auxiliary service's own endpoint) should read its config through this method rather than the + * raw {@code gravitino.<name>.} prefix directly, so both config forms are honored consistently. + * + * @param gravitinoConfig the Gravitino server's configuration + * @param auxServiceName the auxiliary service's short name, e.g. {@code iceberg-rest} + * @return the auxiliary service's own configuration, with keys unprefixed + */ + public static Map<String, String> getAuxServiceConfig( + Config gravitinoConfig, String auxServiceName) { + Map<String, String> serviceConfigs = extractAuxiliaryServiceConfigs(gravitinoConfig); + return MapUtils.getPrefixMap(serviceConfigs, DOT.join(auxServiceName, "")); + } + @VisibleForTesting static String getValidPath(String auxServiceName, String pathString) { Path path = Paths.get(pathString); diff --git a/core/src/test/java/org/apache/gravitino/auxiliary/TestAuxiliaryServiceManager.java b/core/src/test/java/org/apache/gravitino/auxiliary/TestAuxiliaryServiceManager.java index 5df4b897f2..0148321038 100644 --- a/core/src/test/java/org/apache/gravitino/auxiliary/TestAuxiliaryServiceManager.java +++ b/core/src/test/java/org/apache/gravitino/auxiliary/TestAuxiliaryServiceManager.java @@ -139,6 +139,45 @@ public class TestAuxiliaryServiceManager { Assertions.assertFalse(spyAuxManager.isAuxServiceRegistered("lance-rest")); } + @Test + void testGetAuxServiceConfig() { + DummyConfig config = + DummyConfig.of( + ImmutableMap.of( + AuxiliaryServiceManager.GRAVITINO_AUX_SERVICE_PREFIX + + AuxiliaryServiceManager.AUX_SERVICE_NAMES, + "iceberg-rest", + "gravitino.iceberg-rest.host", + "irc-host", + "gravitino.iceberg-rest.httpPort", + "9001")); + + Map<String, String> resolved = + AuxiliaryServiceManager.getAuxServiceConfig(config, "iceberg-rest"); + + Assertions.assertEquals("irc-host", resolved.get("host")); + Assertions.assertEquals("9001", resolved.get("httpPort")); + } + + @Test + void testGetAuxServiceConfigHonorsDeprecatedPrefix() { + // gravitino.auxService.<name>.<key> is deprecated but still honored; the effective config a + // caller reads through this method must match what the service itself receives. + DummyConfig config = + DummyConfig.of( + ImmutableMap.of( + AuxiliaryServiceManager.GRAVITINO_AUX_SERVICE_PREFIX + + AuxiliaryServiceManager.AUX_SERVICE_NAMES, + "iceberg-rest", + AuxiliaryServiceManager.GRAVITINO_AUX_SERVICE_PREFIX + "iceberg-rest.host", + "irc-host")); + + Map<String, String> resolved = + AuxiliaryServiceManager.getAuxServiceConfig(config, "iceberg-rest"); + + Assertions.assertEquals("irc-host", resolved.get("host")); + } + @Test void testAuxiliaryServiceConfigs() { Map<String, String> m = diff --git a/server/src/main/java/org/apache/gravitino/server/web/rest/IcebergRESTServiceOperations.java b/server/src/main/java/org/apache/gravitino/server/web/rest/IcebergRESTServiceOperations.java index f296e068db..ce672875e8 100644 --- a/server/src/main/java/org/apache/gravitino/server/web/rest/IcebergRESTServiceOperations.java +++ b/server/src/main/java/org/apache/gravitino/server/web/rest/IcebergRESTServiceOperations.java @@ -20,6 +20,7 @@ package org.apache.gravitino.server.web.rest; import com.codahale.metrics.annotation.ResponseMetered; import com.codahale.metrics.annotation.Timed; +import java.util.Map; import javax.servlet.http.HttpServletRequest; import javax.ws.rs.Consumes; import javax.ws.rs.GET; @@ -30,12 +31,13 @@ import javax.ws.rs.core.Context; import javax.ws.rs.core.MediaType; import javax.ws.rs.core.Response; import org.apache.commons.lang3.StringUtils; -import org.apache.gravitino.Config; import org.apache.gravitino.GravitinoEnv; import org.apache.gravitino.auxiliary.AuxiliaryServiceManager; import org.apache.gravitino.dto.responses.IcebergRESTServiceResponse; import org.apache.gravitino.metrics.MetricNames; import org.apache.gravitino.server.web.Utils; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * Reports the endpoint of the Gravitino Iceberg REST server, so that clients which already connect @@ -46,16 +48,25 @@ import org.apache.gravitino.server.web.Utils; @Produces(MediaType.APPLICATION_JSON) public class IcebergRESTServiceOperations { + private static final Logger LOG = LoggerFactory.getLogger(IcebergRESTServiceOperations.class); + // Matches gravitino.auxService.names / AuxiliaryServiceManager's registration key. private static final String AUX_SERVICE_NAME = "iceberg-rest"; - private static final String CONFIG_PREFIX = "gravitino.iceberg-rest."; + // Keys below are read from AuxiliaryServiceManager.getAuxServiceConfig, which already strips + // the gravitino.iceberg-rest. (or deprecated gravitino.auxService.iceberg-rest.) prefix, so + // they must NOT be re-prefixed here. + // The provider name used by the Iceberg REST server itself; see + // IcebergConstants.ICEBERG_REST_CATALOG_CONFIG_PROVIDER and DynamicIcebergConfigProvider. The + // server module cannot depend on iceberg-common/catalog-common, hence the literal here. + private static final String CATALOG_CONFIG_PROVIDER_KEY = "catalog-config-provider"; + private static final String DYNAMIC_CATALOG_CONFIG_PROVIDER_NAME = "dynamic-config-provider"; // The post-strip key used by the Iceberg REST server itself; see // IcebergConstants.GRAVITINO_METALAKE and DynamicIcebergConfigProvider. - private static final String SERVED_METALAKE_KEY = CONFIG_PREFIX + "gravitino-metalake"; - private static final String HOST_KEY = CONFIG_PREFIX + "host"; - private static final String HTTP_PORT_KEY = CONFIG_PREFIX + "httpPort"; - private static final String HTTPS_PORT_KEY = CONFIG_PREFIX + "httpsPort"; - private static final String ENABLE_HTTPS_KEY = CONFIG_PREFIX + "enableHttps"; + private static final String SERVED_METALAKE_KEY = "gravitino-metalake"; + private static final String HOST_KEY = "host"; + private static final String HTTP_PORT_KEY = "httpPort"; + private static final String HTTPS_PORT_KEY = "httpsPort"; + private static final String ENABLE_HTTPS_KEY = "enableHttps"; // Match IcebergConfig.DEFAULT_ICEBERG_REST_SERVICE_HTTP_PORT/HTTPS_PORT: the server module // cannot depend on iceberg-common, and JettyServerConfig's own defaults are the Gravitino // server's (8090/8433), not the Iceberg REST server's — reading raw values with these @@ -91,8 +102,12 @@ public class IcebergRESTServiceOperations { return GravitinoEnv.getInstance().auxServiceManager(); } - Config getConfig() { - return GravitinoEnv.getInstance().config(); + // Resolved through AuxiliaryServiceManager.getAuxServiceConfig rather than reading + // gravitino.iceberg-rest.* directly, so the deprecated gravitino.auxService.iceberg-rest.* + // config form is honored too — the same precedence the Iceberg REST server itself sees. + Map<String, String> getIcebergRestServiceConfig() { + return AuxiliaryServiceManager.getAuxServiceConfig( + GravitinoEnv.getInstance().config(), AUX_SERVICE_NAME); } HttpServletRequest getHttpRequest() { @@ -104,17 +119,34 @@ public class IcebergRESTServiceOperations { return null; } - Config config = getConfig(); - String servedMetalake = config.getRawString(SERVED_METALAKE_KEY, ""); + Map<String, String> config = getIcebergRestServiceConfig(); + String provider = config.getOrDefault(CATALOG_CONFIG_PROVIDER_KEY, ""); + if (!DYNAMIC_CATALOG_CONFIG_PROVIDER_NAME.equals(provider)) { + // Only the dynamic catalog config provider maps Iceberg REST catalog names onto Gravitino + // catalogs; the default static provider serves statically-declared catalogs unrelated to + // Gravitino catalog names, so routing at it would 404 on every request. + LOG.debug( + "Iceberg REST service does not use the dynamic catalog config provider " + + "(catalog-config-provider={}); not reporting its endpoint for auto-discovery.", + provider); + return null; + } + + String servedMetalake = config.getOrDefault(SERVED_METALAKE_KEY, ""); if (StringUtils.isNotBlank(metalake) && StringUtils.isNotBlank(servedMetalake) && !servedMetalake.equals(metalake)) { // The Iceberg REST server serves exactly one metalake. Routing a different metalake's // catalogs at it would 404 on every request, so report it as unavailable instead. + LOG.debug( + "Iceberg REST service serves metalake {}, not the requested metalake {}; not " + + "reporting its endpoint for auto-discovery.", + servedMetalake, + metalake); return null; } - String host = config.getRawString(HOST_KEY, DEFAULT_HOST); + String host = config.getOrDefault(HOST_KEY, DEFAULT_HOST); if (isWildcardHost(host)) { // The Iceberg REST server binds to all interfaces, so it has no single externally // reachable address of its own. The caller already reached this Gravitino server at some @@ -123,7 +155,7 @@ public class IcebergRESTServiceOperations { // gravitino.iceberg.rest-uri manually. host = getHttpRequest().getServerName(); } - boolean enableHttps = Boolean.parseBoolean(config.getRawString(ENABLE_HTTPS_KEY, "false")); + boolean enableHttps = Boolean.parseBoolean(config.getOrDefault(ENABLE_HTTPS_KEY, "false")); String scheme = enableHttps ? "https" : "http"; int port = parsePort( @@ -133,8 +165,8 @@ public class IcebergRESTServiceOperations { return String.format("%s://%s:%d/iceberg", scheme, host, port); } - private static int parsePort(Config config, String key, int defaultPort) { - String value = config.getRawString(key, ""); + private static int parsePort(Map<String, String> config, String key, int defaultPort) { + String value = config.getOrDefault(key, ""); if (StringUtils.isBlank(value)) { return defaultPort; } 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 d295f75ada..f4bb9eddb7 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 @@ -27,26 +27,25 @@ import com.google.common.collect.ImmutableMap; import java.util.Map; import javax.servlet.http.HttpServletRequest; import javax.ws.rs.core.Response; -import org.apache.gravitino.Config; import org.apache.gravitino.auxiliary.AuxiliaryServiceManager; import org.apache.gravitino.dto.responses.IcebergRESTServiceResponse; import org.junit.jupiter.api.Test; public class TestIcebergRESTServiceOperations { - private static class DummyConfig extends Config { - static DummyConfig of(Map<String, String> m) { - DummyConfig config = new DummyConfig(); - config.loadFromMap(m, k -> true); - return config; - } + private static final String DYNAMIC_PROVIDER = "dynamic-config-provider"; + + private static Map<String, String> withDynamicProvider(Map<String, String> extra) { + return ImmutableMap.<String, String>builder() + .put("catalog-config-provider", DYNAMIC_PROVIDER) + .putAll(extra) + .buildKeepingLast(); } private IcebergRESTServiceOperations newOps( boolean registered, Map<String, String> icebergConfig, String requestServerName) { AuxiliaryServiceManager auxServiceManager = mock(AuxiliaryServiceManager.class); when(auxServiceManager.isAuxServiceRegistered("iceberg-rest")).thenReturn(registered); - Config config = DummyConfig.of(icebergConfig); HttpServletRequest request = mock(HttpServletRequest.class); when(request.getServerName()).thenReturn(requestServerName); @@ -57,8 +56,8 @@ public class TestIcebergRESTServiceOperations { } @Override - Config getConfig() { - return config; + Map<String, String> getIcebergRestServiceConfig() { + return icebergConfig; } @Override @@ -78,13 +77,34 @@ public class TestIcebergRESTServiceOperations { assertNull(uriOf(ops.getIcebergRestServiceUri("test"))); } + @Test + public void testReturnsNullWhenNotUsingDynamicConfigProvider() { + // Regression test: the default (static) catalog config provider serves statically-declared + // catalogs unrelated to Gravitino catalog names, so it must never be reported for + // auto-discovery — even with a blank/absent catalog-config-provider, which is what the + // common no-provider-configured deployment looks like. + IcebergRESTServiceOperations ops = + newOps(true, ImmutableMap.of("host", "irc-host"), "gravitino-host"); + assertNull(uriOf(ops.getIcebergRestServiceUri("test"))); + } + + @Test + public void testReturnsNullWhenStaticConfigProviderExplicit() { + IcebergRESTServiceOperations ops = + newOps( + true, + ImmutableMap.of("catalog-config-provider", "static-config-provider", "host", "h"), + "gravitino-host"); + assertNull(uriOf(ops.getIcebergRestServiceUri("test"))); + } + @Test public void testDefaultPortIsTheIcebergRestDefaultNotTheGravitinoServerDefault() { - // Regression test: without an explicit gravitino.iceberg-rest.httpPort, the reported port - // must be the Iceberg REST server's own default (9001), not the Gravitino webserver's - // default (8090) that JettyServerConfig would otherwise fall back to. + // Regression test: without an explicit httpPort, the reported port must be the Iceberg REST + // server's own default (9001), not the Gravitino webserver's default (8090) that + // JettyServerConfig would otherwise fall back to. IcebergRESTServiceOperations ops = - newOps(true, ImmutableMap.of("gravitino.iceberg-rest.host", "irc-host"), "gravitino-host"); + newOps(true, withDynamicProvider(ImmutableMap.of("host", "irc-host")), "gravitino-host"); assertEquals("http://irc-host:9001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } @@ -93,9 +113,10 @@ public class TestIcebergRESTServiceOperations { IcebergRESTServiceOperations ops = newOps( true, - ImmutableMap.of( - "gravitino.iceberg-rest.host", "irc-host", - "gravitino.iceberg-rest.enableHttps", "true"), + withDynamicProvider( + ImmutableMap.of( + "host", "irc-host", + "enableHttps", "true")), "gravitino-host"); assertEquals("https://irc-host:9433/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } @@ -105,9 +126,10 @@ public class TestIcebergRESTServiceOperations { IcebergRESTServiceOperations ops = newOps( true, - ImmutableMap.of( - "gravitino.iceberg-rest.host", "irc-host", - "gravitino.iceberg-rest.httpPort", "19001"), + withDynamicProvider( + ImmutableMap.of( + "host", "irc-host", + "httpPort", "19001")), "gravitino-host"); assertEquals("http://irc-host:19001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } @@ -116,16 +138,15 @@ public class TestIcebergRESTServiceOperations { public void testWildcardHostFallsBackToRequestServerName() { IcebergRESTServiceOperations ops = newOps( - true, - ImmutableMap.of("gravitino.iceberg-rest.host", "0.0.0.0"), - "host.docker.internal"); + true, withDynamicProvider(ImmutableMap.of("host", "0.0.0.0")), "host.docker.internal"); assertEquals( "http://host.docker.internal:9001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } @Test public void testBlankHostIsTreatedAsWildcard() { - IcebergRESTServiceOperations ops = newOps(true, ImmutableMap.of(), "gravitino-host"); + IcebergRESTServiceOperations ops = + newOps(true, withDynamicProvider(ImmutableMap.of()), "gravitino-host"); assertEquals("http://gravitino-host:9001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } @@ -134,9 +155,10 @@ public class TestIcebergRESTServiceOperations { IcebergRESTServiceOperations ops = newOps( true, - ImmutableMap.of( - "gravitino.iceberg-rest.host", "irc-host", - "gravitino.iceberg-rest.gravitino-metalake", "prod"), + withDynamicProvider( + ImmutableMap.of( + "host", "irc-host", + "gravitino-metalake", "prod")), "gravitino-host"); assertNull(uriOf(ops.getIcebergRestServiceUri("test"))); } @@ -146,9 +168,10 @@ public class TestIcebergRESTServiceOperations { IcebergRESTServiceOperations ops = newOps( true, - ImmutableMap.of( - "gravitino.iceberg-rest.host", "irc-host", - "gravitino.iceberg-rest.gravitino-metalake", "test"), + withDynamicProvider( + ImmutableMap.of( + "host", "irc-host", + "gravitino-metalake", "test")), "gravitino-host"); assertEquals("http://irc-host:9001/iceberg", uriOf(ops.getIcebergRestServiceUri("test"))); } @@ -158,16 +181,18 @@ public class TestIcebergRESTServiceOperations { IcebergRESTServiceOperations ops = newOps( true, - ImmutableMap.of( - "gravitino.iceberg-rest.host", "irc-host", - "gravitino.iceberg-rest.gravitino-metalake", "prod"), + withDynamicProvider( + ImmutableMap.of( + "host", "irc-host", + "gravitino-metalake", "prod")), "gravitino-host"); assertEquals("http://irc-host:9001/iceberg", uriOf(ops.getIcebergRestServiceUri(""))); } @Test public void testResponseIsNotCacheable() { - IcebergRESTServiceOperations ops = newOps(true, ImmutableMap.of(), "gravitino-host"); + IcebergRESTServiceOperations ops = + newOps(true, withDynamicProvider(ImmutableMap.of()), "gravitino-host"); Response response = ops.getIcebergRestServiceUri(""); assertEquals("no-store", response.getHeaderString("Cache-Control")); } 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 e11c16b30b..ce9dc5b7da 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 @@ -73,6 +73,9 @@ public class CatalogConnectorManager { private String targetMetalake; private final Map<String, GravitinoMetalake> metalakes = new ConcurrentHashMap<>(); + // Tracks which metalakes' Iceberg REST discovery is currently failing, so a failure is logged + // at WARN only on the transition into/out of that state rather than on every poll. + private final Set<String> icebergRestDiscoveryFailing = ConcurrentHashMap.newKeySet(); private GravitinoAdminClient gravitinoClient; private GravitinoConfig config; @@ -213,15 +216,32 @@ public class CatalogConnectorManager { * caches the answer on the shared {@link GravitinoConfig} for {@code IcebergConnectorAdapter} to * read on the next catalog load. Failures — including talking to a Gravitino server older than * this endpoint — must not interrupt catalog loading, so they are swallowed here; Iceberg - * catalogs simply keep their last known routing decision until the next successful poll. + * catalogs simply keep their last known routing decision until the next successful poll. A + * failure is logged at WARN only on the transition into that state (not on every poll), so a + * persistently failing/misconfigured discovery is still visible at the default log level without + * spamming the log for the common case of an older Gravitino server that lacks this endpoint + * entirely. */ private void refreshIcebergRestUri(String metalakeName) { try { config.setDiscoveredIcebergRestUri( metalakeName, gravitinoClient.icebergRestServiceUri(metalakeName).orElse(null)); + if (icebergRestDiscoveryFailing.remove(metalakeName)) { + LOG.info("Iceberg REST service discovery for metalake {} recovered.", metalakeName); + } } catch (Exception e) { - LOG.debug( - "Failed to query the Iceberg REST service endpoint for metalake {}.", metalakeName, e); + if (icebergRestDiscoveryFailing.add(metalakeName)) { + LOG.warn( + "Failed to query the Iceberg REST service endpoint for metalake {}; keeping the " + + "last known routing decision. This is expected when talking to a Gravitino " + + "server that predates this endpoint, but is otherwise worth investigating. " + + "Further failures for this metalake are logged at DEBUG until it recovers.", + metalakeName, + e); + } else { + LOG.debug( + "Failed to query the Iceberg REST service endpoint for metalake {}.", metalakeName, e); + } } } 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 35f722bb1e..4189a251cc 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 @@ -35,6 +35,8 @@ import java.sql.Statement; import java.util.Map; import java.util.Properties; import java.util.Set; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import org.apache.commons.lang3.StringUtils; import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; @@ -60,6 +62,16 @@ public class CatalogRegister { private static final Set<String> SSL_VERIFICATION_MODES = ImmutableSet.of(SSL_VERIFICATION_FULL, SSL_VERIFICATION_CA, SSL_VERIFICATION_NONE); + // Matches "some.key.credential"='value' style property assignments (as produced by + // GravitinoConfig.toCatalogConfig / GravitinoCatalog.toJson-embedded WITH(...) clauses) whose + // key looks like it carries a secret, e.g. gravitino.iceberg.rest-catalog.oauth2.credential. + // Values matching this are redacted before the CREATE CATALOG statement is logged, since it can + // otherwise leak IRC OAuth2 client secrets and similar into the log and Trino's query history. + private static final Pattern SECRET_PROPERTY_PATTERN = + Pattern.compile( + "\"([^\"]*(?:credential|token|secret|password)[^\"]*)\"\\s*=\\s*'([^']*)'", + Pattern.CASE_INSENSITIVE); + private Connection connection; private boolean isStarted = false; private String catalogStoreDirectory; @@ -303,7 +315,15 @@ public class CatalogRegister { } } - private String generateCreateCatalogCommand(String name, GravitinoCatalog gravitinoCatalog) + // Package-private (rather than private) so tests can exercise the embed-then-serialize wiring + // without a live Trino JDBC connection, which init() otherwise requires. + @VisibleForTesting + void setConfigForTesting(GravitinoConfig config) { + this.config = config; + } + + @VisibleForTesting + String generateCreateCatalogCommand(String name, GravitinoCatalog gravitinoCatalog) throws Exception { // This statement is replicated by Trino to every node in the cluster, coordinator and workers // alike, so it is the only place a value the coordinator alone knows (like the Iceberg REST @@ -320,6 +340,18 @@ public class CatalogRegister { config.toCatalogConfig()); } + @VisibleForTesting + static String redactSecrets(String createCatalogCommand) { + Matcher matcher = SECRET_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); } @@ -356,7 +388,7 @@ public class CatalogRegister { } String createCatalogCommand = generateCreateCatalogCommand(name, catalog); executeSql(createCatalogCommand); - LOG.info("Register catalog {} successfully: {}", name, createCatalogCommand); + LOG.info("Register catalog {} successfully: {}", name, redactSecrets(createCatalogCommand)); } catch (SQLException e) { throw new TrinoException(GravitinoErrorCode.GRAVITINO_RUNTIME_ERROR, e.getMessage(), e); } catch (Exception e) { diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java index f03cf0a474..5ac9efb5e7 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/TestCatalogConnectorManager.java @@ -30,7 +30,10 @@ import static org.mockito.Mockito.when; import com.google.common.collect.ImmutableMap; import io.trino.spi.TrinoException; import io.trino.spi.connector.ConnectorContext; +import java.util.Optional; import org.apache.gravitino.client.GravitinoAdminClient; +import org.apache.gravitino.client.GravitinoMetalake; +import org.apache.gravitino.exceptions.RESTException; import org.apache.gravitino.trino.connector.GravitinoConfig; import org.apache.gravitino.trino.connector.GravitinoErrorCode; import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; @@ -140,6 +143,56 @@ public class TestCatalogConnectorManager { assertFalse(manager.skipCatalog("b2")); } + @Test + public void testRefreshIcebergRestUriCachesDiscoveredUri() throws Exception { + GravitinoAdminClient client = mock(GravitinoAdminClient.class); + CatalogRegister catalogRegister = mock(CatalogRegister.class); + when(catalogRegister.isTrinoStarted()).thenReturn(true); + when(client.loadMetalake("test")).thenReturn(mock(GravitinoMetalake.class)); + when(client.icebergRestServiceUri("test")) + .thenReturn(Optional.of("http://irc-host:9001/iceberg")); + + CatalogConnectorManager manager = + new CatalogConnectorManager(catalogRegister, createCatalogConnectorFactory(), null); + GravitinoConfig config = + new GravitinoConfig( + ImmutableMap.of( + "gravitino.uri", "http://127.0.0.1:8090", + "gravitino.metalake", "test", + "gravitino.use-single-metalake", "true")); + manager.config(config, client); + + manager.loadMetalakeSync(); + + assertEquals("http://irc-host:9001/iceberg", config.getDiscoveredIcebergRestUri("test")); + } + + @Test + public void testRefreshIcebergRestUriSwallowsFailureAndKeepsCatalogLoadingGoing() + throws Exception { + GravitinoAdminClient client = mock(GravitinoAdminClient.class); + CatalogRegister catalogRegister = mock(CatalogRegister.class); + when(catalogRegister.isTrinoStarted()).thenReturn(true); + when(client.loadMetalake("test")).thenReturn(mock(GravitinoMetalake.class)); + when(client.icebergRestServiceUri("test")) + .thenThrow(new RESTException("simulated: endpoint not found on an older server")); + + CatalogConnectorManager manager = + new CatalogConnectorManager(catalogRegister, createCatalogConnectorFactory(), null); + GravitinoConfig config = + new GravitinoConfig( + ImmutableMap.of( + "gravitino.uri", "http://127.0.0.1:8090", + "gravitino.metalake", "test", + "gravitino.use-single-metalake", "true")); + manager.config(config, client); + + // A discovery failure must not abort the metalake load (which loads catalogs), and must + // leave the discovered URI at its previous value rather than throwing out of loadMetalake. + assertDoesNotThrow(manager::loadMetalakeSync); + assertEquals("", config.getDiscoveredIcebergRestUri("test")); + } + private CatalogConnectorManager createManager(ImmutableMap<String, String> configMap) throws Exception { return createManager(createCatalogConnectorFactory(), configMap); 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 2effd98843..488cfc7ed6 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 @@ -26,14 +26,19 @@ import static org.junit.jupiter.api.Assertions.assertNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; +import com.google.common.collect.ImmutableMap; import io.trino.spi.TrinoException; import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; +import java.util.Collections; import java.util.HashMap; import java.util.Map; import java.util.Properties; +import org.apache.gravitino.Catalog; import org.apache.gravitino.trino.connector.GravitinoConfig; +import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog; +import org.apache.gravitino.trino.connector.metadata.TestGravitinoCatalog; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.io.TempDir; @@ -377,4 +382,87 @@ public class TestCatalogRegister { assertTrue(e.getMessage().contains("does not exist")); assertTrue(e.getMessage().contains("trino.jdbc.ssl.truststore.path")); } + + private static final String DISCOVERED_ICEBERG_REST_URI_PROPERTY = "__gravitino.iceberg.rest-uri"; + + @Test + public void testGenerateCreateCatalogCommandEmbedsDiscoveredUriForIcebergCatalog() + throws Exception { + GravitinoConfig config = + new GravitinoConfig( + ImmutableMap.of( + "gravitino.uri", "http://127.0.0.1:8090", "gravitino.metalake", "test")); + config.setDiscoveredIcebergRestUri("test", "http://irc-host:9001/iceberg"); + + CatalogRegister catalogRegister = new CatalogRegister(); + catalogRegister.setConfigForTesting(config); + + Catalog mockCatalog = + TestGravitinoCatalog.mockCatalog( + "iceberg_catalog", + "lakehouse-iceberg", + "test catalog", + Catalog.Type.RELATIONAL, + Collections.emptyMap()); + GravitinoCatalog catalog = new GravitinoCatalog("test", mockCatalog); + + String command = catalogRegister.generateCreateCatalogCommand("iceberg_catalog", catalog); + + assertTrue( + command.contains( + "\"" + DISCOVERED_ICEBERG_REST_URI_PROPERTY + "\":\"http://irc-host:9001/iceberg\""), + "Expected the discovered Iceberg REST URI to be embedded in: " + command); + } + + @Test + public void testGenerateCreateCatalogCommandDoesNotEmbedUriForNonIcebergCatalog() + throws Exception { + // The discovered URI is per-Iceberg-catalog routing state; embedding it into every catalog's + // properties (e.g. a Hive catalog) would be a leaky abstraction and is guarded against in + // IcebergConnectorAdapter.embedDiscoveredIcebergRestUri. This asserts that guard actually + // takes effect when reached through CatalogRegister, not just when called directly. + GravitinoConfig config = + new GravitinoConfig( + ImmutableMap.of( + "gravitino.uri", "http://127.0.0.1:8090", "gravitino.metalake", "test")); + config.setDiscoveredIcebergRestUri("test", "http://irc-host:9001/iceberg"); + + CatalogRegister catalogRegister = new CatalogRegister(); + catalogRegister.setConfigForTesting(config); + + Catalog mockCatalog = + TestGravitinoCatalog.mockCatalog( + "hive_catalog", + "hive", + "test catalog", + Catalog.Type.RELATIONAL, + Collections.emptyMap()); + GravitinoCatalog catalog = new GravitinoCatalog("test", mockCatalog); + + String command = catalogRegister.generateCreateCatalogCommand("hive_catalog", catalog); + + assertFalse(command.contains(DISCOVERED_ICEBERG_REST_URI_PROPERTY)); + } + + @Test + public void testRedactSecretsMasksSecretBearingProperties() { + String command = + "CREATE CATALOG c USING gravitino WITH ( " + + "\"gravitino.iceberg.rest-catalog.oauth2.credential\"='client:secretvalue', " + + "\"gravitino.iceberg.rest-catalog.uri\"='http://irc-host:9001/iceberg', " + + "\"some.token\"='abc123', " + + "\"trino.bypass.password\"='hunter2')"; + + String redacted = CatalogRegister.redactSecrets(command); + + assertFalse(redacted.contains("secretvalue")); + assertFalse(redacted.contains("abc123")); + assertFalse(redacted.contains("hunter2")); + assertTrue(redacted.contains("\"gravitino.iceberg.rest-catalog.oauth2.credential\"='***'")); + assertTrue(redacted.contains("\"some.token\"='***'")); + assertTrue(redacted.contains("\"trino.bypass.password\"='***'")); + // Non-secret properties must survive redaction unchanged. + assertTrue( + redacted.contains("\"gravitino.iceberg.rest-catalog.uri\"='http://irc-host:9001/iceberg'")); + } } diff --git a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java index ee3e448c29..1951003a19 100644 --- a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java +++ b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/catalog/iceberg/TestIcebergCatalogPropertyConverter.java @@ -426,6 +426,46 @@ public class TestIcebergCatalogPropertyConverter { Assertions.assertNull(config.get("iceberg.rest-catalog.uri")); } + @Test + public void testManualUriWinsOverDiscoveredUri() throws Exception { + // Both a manual gravitino.iceberg.rest-uri and a discovered per-metalake URI are set; the + // manual override must win, matching IcebergConnectorAdapter.buildInternalConnectorConfig's + // documented precedence. + Map<String, String> properties = + ImmutableMap.<String, String>builder() + .put("catalog-backend", "jdbc") + .put("uri", "jdbc:postgresql://localhost:5432/iceberg") + .put("jdbc-driver", "org.postgresql.Driver") + .build(); + + GravitinoConfig config = icebergRestConfiguredConfig(ImmutableMap.of()); + config.setDiscoveredIcebergRestUri("test", "http://discovered:9001/iceberg"); + + Map<String, String> connectorConfig = + buildConnectorConfig("catalog1", properties, config, /* embedDiscovery= */ true); + + Assertions.assertEquals("rest", connectorConfig.get("iceberg.catalog.type")); + Assertions.assertEquals( + "http://localhost:9001/iceberg", connectorConfig.get("iceberg.rest-catalog.uri")); + } + + @Test + public void testEmbedDiscoveredIcebergRestUriDoesNotAffectNonIcebergCatalog() { + // CatalogRegister calls embedDiscoveredIcebergRestUri unconditionally for every catalog being + // registered, not just Iceberg ones, so this guard is load-bearing: a Hive/MySQL/etc. catalog + // must never end up carrying the synthetic discovered-URI property. + GravitinoConfig config = icebergRestDiscoveredConfig("test", "http://discovered:9001/iceberg"); + Catalog mockHiveCatalog = + TestGravitinoCatalog.mockCatalog( + "hive_catalog", "hive", "test catalog", Catalog.Type.RELATIONAL, ImmutableMap.of()); + GravitinoCatalog hiveCatalog = new GravitinoCatalog("test", mockHiveCatalog); + + GravitinoCatalog result = + IcebergConnectorAdapter.embedDiscoveredIcebergRestUri(hiveCatalog, config); + + Assertions.assertEquals(hiveCatalog.getProperties(), result.getProperties()); + } + @Test public void testBuildConnectorPropertiesWithIcebergRestAuthentication() throws Exception { Map<String, String> properties =
