This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new 9e5698a1fa [Cherry-pick to branch-1.3] [#13194]
improvement(trino-connector): Defer Iceberg REST passthrough metadata
initialization (#13195) (#13205)
9e5698a1fa is described below
commit 9e5698a1fa6f30a86284fc70860b0970fdfc91c7
Author: github-actions[bot]
<41898282+github-actions[bot]@users.noreply.github.com>
AuthorDate: Wed Sep 16 16:23:27 2026 +0800
[Cherry-pick to branch-1.3] [#13194] improvement(trino-connector): Defer
Iceberg REST passthrough metadata initialization (#13195) (#13205)
**Cherry-pick Information:**
- Original commit: 422af4130bd482ea62ced40d3948ead302846cea
- Target branch: `branch-1.3`
- Status: ✅ Clean cherry-pick (no conflicts)
Co-authored-by: Yuhui <[email protected]>
---
.../trino/connector/DeferredConnectorMetadata.java | 111 ++++++++++
.../gravitino/trino/connector/GravitinoConfig.java | 7 +-
.../trino/connector/GravitinoConnector.java | 31 ++-
.../connector/catalog/CatalogConnectorContext.java | 18 +-
.../connector/TestDeferredConnectorMetadata.java | 229 +++++++++++++++++++++
.../trino/connector/TestGravitinoConfig.java | 25 +++
.../TestGravitinoConnectorPassthrough.java | 211 +++++++++++++++++++
7 files changed, 627 insertions(+), 5 deletions(-)
diff --git
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/DeferredConnectorMetadata.java
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/DeferredConnectorMetadata.java
new file mode 100644
index 0000000000..1b2919de51
--- /dev/null
+++
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/DeferredConnectorMetadata.java
@@ -0,0 +1,111 @@
+/*
+ * 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;
+
+import io.trino.spi.connector.ConnectorMetadata;
+import io.trino.spi.connector.ConnectorSession;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
+import java.lang.reflect.Proxy;
+import java.util.Objects;
+import java.util.function.Function;
+import javax.annotation.Nullable;
+
+/** Lazily obtains native metadata without authenticating metadata-only
management queries. */
+final class DeferredConnectorMetadata {
+ private final Function<ConnectorSession, ConnectorMetadata> factory;
+ private final ConnectorSession querySession;
+ @Nullable private ConnectorMetadata delegate;
+ @Nullable private ConnectorSession pendingBegin;
+ private boolean closed;
+
+ private DeferredConnectorMetadata(
+ ConnectorSession querySession, Function<ConnectorSession,
ConnectorMetadata> factory) {
+ this.querySession = Objects.requireNonNull(querySession, "querySession");
+ this.factory = Objects.requireNonNull(factory, "factory");
+ }
+
+ static ConnectorMetadata create(
+ ConnectorSession querySession, Function<ConnectorSession,
ConnectorMetadata> factory) {
+ DeferredConnectorMetadata handler = new
DeferredConnectorMetadata(querySession, factory);
+ return (ConnectorMetadata)
+ Proxy.newProxyInstance(
+ ConnectorMetadata.class.getClassLoader(),
+ new Class<?>[] {ConnectorMetadata.class},
+ handler::invoke);
+ }
+
+ @Nullable
+ private synchronized Object invoke(Object proxy, Method method, @Nullable
Object[] args)
+ throws Throwable {
+ if (method.getDeclaringClass() == Object.class) {
+ switch (method.getName()) {
+ case "toString":
+ return "DeferredConnectorMetadata";
+ case "hashCode":
+ return System.identityHashCode(proxy);
+ case "equals":
+ return proxy == args[0];
+ default:
+ throw new UnsupportedOperationException(method.getName());
+ }
+ }
+ if (method.getName().equals("cleanupQuery")) {
+ closed = true;
+ pendingBegin = null;
+ if (delegate == null) {
+ return null;
+ }
+ } else {
+ if (method.getName().equals("beginQuery")) {
+ closed = false;
+ }
+ if (closed) {
+ throw new IllegalStateException("Metadata query is already closed");
+ }
+ if (method.getName().equals("beginQuery") && delegate == null) {
+ pendingBegin = (ConnectorSession) args[0];
+ return null;
+ }
+ if (delegate == null) {
+ ConnectorSession currentSession = querySession;
+ if (args != null) {
+ for (Object arg : args) {
+ if (arg instanceof ConnectorSession) {
+ currentSession = (ConnectorSession) arg;
+ break;
+ }
+ }
+ }
+ ConnectorMetadata metadata =
+ Objects.requireNonNull(factory.apply(currentSession), "native
metadata");
+ delegate = metadata;
+ if (pendingBegin != null) {
+ pendingBegin = null;
+ metadata.beginQuery(currentSession);
+ }
+ }
+ }
+ try {
+ return method.invoke(delegate, args);
+ } catch (InvocationTargetException e) {
+ throw e.getCause();
+ }
+ }
+}
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 491ac2fa9e..107cde4b82 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
@@ -80,6 +80,8 @@ public class GravitinoConfig {
/** The Trino Iceberg REST catalog property prefix. */
private static final String TRINO_ICEBERG_REST_CATALOG_PREFIX =
"iceberg.rest-catalog.";
+ private static final String OAUTH2 = "OAUTH2";
+
/** Prefix for environment-variable references propagated to dynamic
catalogs. */
static final String GRAVITINO_DYNAMIC_CATALOG_ENV_PREFIX =
"gravitino.dynamic-catalog.environment-variable.";
@@ -830,8 +832,9 @@ public class GravitinoConfig {
String prefix = GRAVITINO_ICEBERG_REST_CATALOG_CONFIG_PREFIX.key;
Map<String, String> restCatalogConfig = new HashMap<>();
- if
("oauth2".equalsIgnoreCase(config.get(GravitinoAuthProvider.AUTH_TYPE_KEY))) {
- restCatalogConfig.put(TRINO_ICEBERG_REST_CATALOG_PREFIX + "security",
"OAUTH2");
+ if
(OAUTH2.equalsIgnoreCase(config.get(GravitinoAuthProvider.AUTH_TYPE_KEY))
+ && OAUTH2.equalsIgnoreCase(config.getOrDefault(prefix + "security",
OAUTH2))) {
+ restCatalogConfig.put(TRINO_ICEBERG_REST_CATALOG_PREFIX + "security",
OAUTH2);
putIfNotBlank(
restCatalogConfig,
TRINO_ICEBERG_REST_CATALOG_PREFIX + "oauth2.credential",
diff --git
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java
index 57dc9913c6..0eb06e1495 100644
---
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java
+++
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/GravitinoConnector.java
@@ -67,6 +67,9 @@ import
org.apache.gravitino.trino.connector.security.GravitinoAuthProvider;
public class GravitinoConnector implements Connector {
private static final Logger LOG = Logger.get(GravitinoConnector.class);
+ private static final String ICEBERG_PROVIDER = "lakehouse-iceberg";
+ private static final String TRINO_ICEBERG_REST_SECURITY =
"iceberg.rest-catalog.security";
+ private static final String OAUTH2_PASSTHROUGH = "OAUTH2_PASSTHROUGH";
private final NameIdentifier catalogIdentifier;
protected final CatalogConnectorContext catalogConnectorContext;
@@ -111,9 +114,8 @@ public class GravitinoConnector implements Connector {
GravitinoTransactionHandle gravitinoTransactionHandle =
(GravitinoTransactionHandle) transactionHandle;
- Connector internalConnector =
catalogConnectorContext.getInternalConnector();
ConnectorMetadata internalMetadata =
- internalConnector.getMetadata(session,
gravitinoTransactionHandle.getInternalHandle());
+ getInternalMetadata(session,
gravitinoTransactionHandle.getInternalHandle());
Preconditions.checkArgument(internalMetadata != null, "Internal metadata
must not be null");
CatalogConnectorMetadata metadata =
@@ -122,6 +124,31 @@ public class GravitinoConnector implements Connector {
metadata, catalogConnectorContext.getMetadataAdapter(),
internalMetadata);
}
+ /**
+ * Defers native REST authentication until a data operation needs it.
Catalog registration uses a
+ * password-authenticated management session without a delegated user token.
+ *
+ * @param session the authenticated query session
+ * @param transactionHandle the native transaction handle
+ * @return metadata that preserves user authentication at first data access
+ */
+ protected ConnectorMetadata getInternalMetadata(
+ ConnectorSession session, ConnectorTransactionHandle transactionHandle) {
+ if
(ICEBERG_PROVIDER.equals(catalogConnectorContext.getCatalog().getProvider())
+ && OAUTH2_PASSTHROUGH.equalsIgnoreCase(
+ catalogConnectorContext
+ .getInternalConnectorConfig()
+ .get(TRINO_ICEBERG_REST_SECURITY))) {
+ return DeferredConnectorMetadata.create(
+ session,
+ currentSession ->
+ catalogConnectorContext
+ .getInternalConnector()
+ .getMetadata(currentSession, transactionHandle));
+ }
+ return catalogConnectorContext.getInternalConnector().getMetadata(session,
transactionHandle);
+ }
+
protected GravitinoMetadata createGravitinoMetadata(
CatalogConnectorMetadata catalogConnectorMetadata,
CatalogConnectorMetadataAdapter metadataAdapter,
diff --git
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorContext.java
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorContext.java
index bbed544317..d529e07a61 100644
---
a/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorContext.java
+++
b/trino-connector/trino-connector/src/main/java/org/apache/gravitino/trino/connector/catalog/CatalogConnectorContext.java
@@ -48,6 +48,8 @@ public class CatalogConnectorContext {
// Internal connector communicates with data storage
private final Connector internalConnector;
+ private final Map<String, String> internalConnectorConfig;
+
private final CatalogConnectorAdapter adapter;
private final GravitinoConfig config;
@@ -58,6 +60,8 @@ public class CatalogConnectorContext {
* @param catalog the Gravitino catalog
* @param metalake the Gravitino metalake
* @param internalConnector the internal connector
+ * @param internalConnectorConfig the effective configuration used to create
the internal
+ * connector
* @param adapter the catalog connector adapter
* @param config the Gravitino connector configuration
*/
@@ -65,11 +69,13 @@ public class CatalogConnectorContext {
GravitinoCatalog catalog,
GravitinoMetalake metalake,
Connector internalConnector,
+ Map<String, String> internalConnectorConfig,
CatalogConnectorAdapter adapter,
GravitinoConfig config) {
this.catalog = catalog;
this.metalake = metalake;
this.internalConnector = internalConnector;
+ this.internalConnectorConfig = Map.copyOf(internalConnectorConfig);
this.adapter = adapter;
this.config = config;
}
@@ -119,6 +125,15 @@ public class CatalogConnectorContext {
return internalConnector;
}
+ /**
+ * Returns the effective configuration used to create the internal connector.
+ *
+ * @return the immutable internal connector configuration
+ */
+ public Map<String, String> getInternalConnectorConfig() {
+ return internalConnectorConfig;
+ }
+
/**
* Returns the table properties associated with this context.
*
@@ -271,7 +286,8 @@ public class CatalogConnectorContext {
Connector connector =
GravitinoConnectorPluginManager.instance(context.getClass().getClassLoader())
.createConnector(internalConnectorName, connectorConfig,
context);
- return new CatalogConnectorContext(catalog, metalake, connector,
connectorAdapter, config);
+ return new CatalogConnectorContext(
+ catalog, metalake, connector, connectorConfig, connectorAdapter,
config);
}
}
}
diff --git
a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestDeferredConnectorMetadata.java
b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestDeferredConnectorMetadata.java
new file mode 100644
index 0000000000..688c1010c0
--- /dev/null
+++
b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestDeferredConnectorMetadata.java
@@ -0,0 +1,229 @@
+/*
+ * 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;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+import io.trino.spi.connector.ConnectorMetadata;
+import io.trino.spi.connector.ConnectorSession;
+import java.util.List;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.junit.jupiter.api.Test;
+
+class TestDeferredConnectorMetadata {
+ @Test
+ void managementLifecycleDoesNotAuthenticate() {
+ AtomicInteger calls = new AtomicInteger();
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ mock(ConnectorSession.class),
+ currentSession -> {
+ calls.incrementAndGet();
+ throw new IllegalStateException("No user token in management
session");
+ });
+ ConnectorSession session = mock(ConnectorSession.class);
+ metadata.beginQuery(session);
+ assertEquals("DeferredConnectorMetadata", metadata.toString());
+ assertEquals(metadata, metadata);
+ metadata.hashCode();
+ metadata.cleanupQuery(session);
+ assertEquals(0, calls.get());
+ assertThrows(IllegalStateException.class, () ->
metadata.listSchemaNames(session));
+ assertEquals(0, calls.get());
+ }
+
+ @Test
+ void dataOperationInitializesOnceAndPreservesLifecycleOrder() {
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ ConnectorSession session = mock(ConnectorSession.class);
+ when(nativeMetadata.listSchemaNames(session)).thenReturn(List.of("demo"));
+ AtomicInteger calls = new AtomicInteger();
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ mock(ConnectorSession.class),
+ currentSession -> {
+ calls.incrementAndGet();
+ return nativeMetadata;
+ });
+ metadata.beginQuery(session);
+ verifyNoInteractions(nativeMetadata);
+ assertEquals(List.of("demo"), metadata.listSchemaNames(session));
+ assertEquals(List.of("demo"), metadata.listSchemaNames(session));
+ metadata.cleanupQuery(session);
+ assertEquals(1, calls.get());
+ var order = inOrder(nativeMetadata);
+ order.verify(nativeMetadata).beginQuery(session);
+ order.verify(nativeMetadata, times(2)).listSchemaNames(session);
+ order.verify(nativeMetadata).cleanupQuery(session);
+ }
+
+ @Test
+ void initializedMetadataSupportsTwoQueryLifecycles() {
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ ConnectorSession first = mock(ConnectorSession.class);
+ ConnectorSession second = mock(ConnectorSession.class);
+ AtomicInteger calls = new AtomicInteger();
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ first,
+ currentSession -> {
+ calls.incrementAndGet();
+ return nativeMetadata;
+ });
+ metadata.beginQuery(first);
+ metadata.listSchemaNames(first);
+ metadata.cleanupQuery(first);
+ assertThrows(IllegalStateException.class, () ->
metadata.listSchemaNames(second));
+ metadata.beginQuery(second);
+ metadata.listSchemaNames(second);
+ metadata.cleanupQuery(second);
+ assertEquals(1, calls.get());
+ var order = inOrder(nativeMetadata);
+ order.verify(nativeMetadata).beginQuery(first);
+ order.verify(nativeMetadata).listSchemaNames(first);
+ order.verify(nativeMetadata).cleanupQuery(first);
+ order.verify(nativeMetadata).beginQuery(second);
+ order.verify(nativeMetadata).listSchemaNames(second);
+ order.verify(nativeMetadata).cleanupQuery(second);
+ order.verifyNoMoreInteractions();
+ }
+
+ @Test
+ void managementQueryCanBeFollowedByDeferredDataQuery() {
+ ConnectorSession management = mock(ConnectorSession.class);
+ ConnectorSession user = mock(ConnectorSession.class);
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ AtomicInteger calls = new AtomicInteger();
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ management,
+ currentSession -> {
+ calls.incrementAndGet();
+ assertSame(user, currentSession);
+ return nativeMetadata;
+ });
+ metadata.beginQuery(management);
+ metadata.cleanupQuery(management);
+ assertThrows(IllegalStateException.class, () ->
metadata.listSchemaNames(user));
+ assertEquals(0, calls.get());
+ metadata.beginQuery(user);
+ assertEquals(0, calls.get());
+ metadata.listSchemaNames(user);
+ metadata.cleanupQuery(user);
+ assertEquals(1, calls.get());
+ var order = inOrder(nativeMetadata);
+ order.verify(nativeMetadata).beginQuery(user);
+ order.verify(nativeMetadata).listSchemaNames(user);
+ order.verify(nativeMetadata).cleanupQuery(user);
+ order.verifyNoMoreInteractions();
+ }
+
+ @Test
+ void failedBeginQueryStillCleansUpNativeMetadata() {
+ ConnectorSession session = mock(ConnectorSession.class);
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ IllegalStateException failure = new IllegalStateException("Query
initialization failed");
+ doThrow(failure).when(nativeMetadata).beginQuery(session);
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(session, currentSession ->
nativeMetadata);
+ metadata.beginQuery(session);
+ assertSame(
+ failure,
+ assertThrows(IllegalStateException.class, () ->
metadata.listSchemaNames(session)));
+ metadata.cleanupQuery(session);
+ var order = inOrder(nativeMetadata);
+ order.verify(nativeMetadata).beginQuery(session);
+ order.verify(nativeMetadata).cleanupQuery(session);
+ order.verifyNoMoreInteractions();
+ }
+
+ @Test
+ void authenticationFailurePropagatesWithoutFallback() {
+ IllegalArgumentException failure = new IllegalArgumentException("Token
rejected");
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ mock(ConnectorSession.class),
+ currentSession -> {
+ throw failure;
+ });
+ ConnectorSession session = mock(ConnectorSession.class);
+ metadata.beginQuery(session);
+ assertSame(
+ failure,
+ assertThrows(IllegalArgumentException.class, () ->
metadata.listSchemaNames(session)));
+ assertDoesNotThrow(() -> metadata.cleanupQuery(session));
+ }
+
+ @Test
+ void permissionFailureIsNotWrappedByReflection() {
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ ConnectorSession session = mock(ConnectorSession.class);
+ SecurityException failure = new SecurityException("Access denied");
+ when(nativeMetadata.listSchemaNames(session)).thenThrow(failure);
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ mock(ConnectorSession.class), currentSession -> nativeMetadata);
+ assertSame(
+ failure, assertThrows(SecurityException.class, () ->
metadata.listSchemaNames(session)));
+ }
+
+ @Test
+ void differentQueriesDoNotShareNativeMetadata() {
+ ConnectorMetadata alice = mock(ConnectorMetadata.class);
+ ConnectorMetadata bob = mock(ConnectorMetadata.class);
+ ConnectorSession session = mock(ConnectorSession.class);
+ ConnectorMetadata a =
+ DeferredConnectorMetadata.create(mock(ConnectorSession.class),
currentSession -> alice);
+ ConnectorMetadata b =
+ DeferredConnectorMetadata.create(mock(ConnectorSession.class),
currentSession -> bob);
+ a.listSchemaNames(session);
+ verifyNoInteractions(bob);
+ b.listSchemaNames(session);
+ verify(alice).listSchemaNames(session);
+ verify(bob).listSchemaNames(session);
+ }
+
+ @Test
+ void initializationUsesOperationSessionInsteadOfCreationSession() {
+ ConnectorSession management = mock(ConnectorSession.class);
+ ConnectorSession user = mock(ConnectorSession.class);
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ ConnectorMetadata metadata =
+ DeferredConnectorMetadata.create(
+ management,
+ currentSession -> {
+ assertSame(user, currentSession);
+ return nativeMetadata;
+ });
+ metadata.beginQuery(management);
+ metadata.listSchemaNames(user);
+ verify(nativeMetadata).beginQuery(user);
+ verify(nativeMetadata).listSchemaNames(user);
+ }
+}
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 73b1192169..0c55d5a26e 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
@@ -389,6 +389,31 @@ public class TestGravitinoConfig {
restCatalogConfig.get("iceberg.rest-catalog.oauth2.server-uri"));
}
+ /** Verifies non-OAuth2 REST modes do not inherit Gravitino service
credentials. */
+ @Test
+ public void testIcebergRestPassthroughDoesNotInheritServiceCredentials() {
+ for (String security : new String[] {"OAUTH2_PASSTHROUGH", "NONE"}) {
+ GravitinoConfig config =
+ new GravitinoConfig(
+ ImmutableMap.<String, String>builder()
+ .put("gravitino.metalake", "test")
+ .put("gravitino.client.authType", "oauth2")
+ .put("gravitino.client.oauth2.serverUri",
"https://idp.example.com")
+ .put("gravitino.client.oauth2.path", "token")
+ .put("gravitino.client.oauth2.credential", "service:secret")
+ .put("gravitino.client.oauth2.scope", "openid")
+ .put("gravitino.iceberg.rest-catalog.security", security)
+ .put("gravitino.iceberg.rest-catalog.session", "NONE")
+ .build());
+ assertEquals(
+ ImmutableMap.of(
+ "iceberg.rest-catalog.security", security,
"iceberg.rest-catalog.session", "NONE"),
+ config.getIcebergRestCatalogConfig());
+ assertEquals(
+ "service:secret",
config.getClientConfig().get("gravitino.client.oauth2.credential"));
+ }
+ }
+
@Test
public void testIcebergRestOAuthOverridesGravitinoClientOAuthByField() {
GravitinoConfig config =
diff --git
a/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnectorPassthrough.java
b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnectorPassthrough.java
new file mode 100644
index 0000000000..5091df8d3f
--- /dev/null
+++
b/trino-connector/trino-connector/src/test/java/org/apache/gravitino/trino/connector/TestGravitinoConnectorPassthrough.java
@@ -0,0 +1,211 @@
+/*
+ * 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;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
+
+import io.trino.spi.connector.Connector;
+import io.trino.spi.connector.ConnectorMetadata;
+import io.trino.spi.connector.ConnectorSession;
+import io.trino.spi.connector.ConnectorTransactionHandle;
+import io.trino.spi.security.ConnectorIdentity;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import javax.annotation.Nullable;
+import org.apache.gravitino.Catalog;
+import org.apache.gravitino.NameIdentifier;
+import org.apache.gravitino.SupportsSchemas;
+import org.apache.gravitino.client.GravitinoMetalake;
+import org.apache.gravitino.trino.connector.catalog.CatalogConnectorContext;
+import org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadata;
+import
org.apache.gravitino.trino.connector.catalog.CatalogConnectorMetadataAdapter;
+import
org.apache.gravitino.trino.connector.catalog.iceberg.IcebergCatalogPropertyConverter;
+import org.apache.gravitino.trino.connector.metadata.GravitinoCatalog;
+import org.junit.jupiter.api.Test;
+
+class TestGravitinoConnectorPassthrough {
+ @Test
+ void passwordCatalogLifecycleAndSchemaListingDoNotRequireNativeToken() {
+ Connector delegate = mock(Connector.class);
+ GravitinoConnector connector =
+ connector(context("OAUTH2_PASSTHROUGH", "lakehouse-iceberg",
delegate));
+ ConnectorSession session = mock(ConnectorSession.class);
+
when(session.getIdentity()).thenReturn(ConnectorIdentity.ofUser("catalog_manager"));
+ ConnectorTransactionHandle transaction =
mock(ConnectorTransactionHandle.class);
+ ConnectorMetadata metadata =
+ connector.getMetadata(session, new
GravitinoTransactionHandle(transaction));
+ metadata.beginQuery(session);
+ assertEquals(List.of("demo"), metadata.listSchemaNames(session));
+ metadata.cleanupQuery(session);
+ verifyNoInteractions(delegate);
+ }
+
+ @Test
+ void passthroughDataAccessUsesOriginalSessionAndPropagatesMissingToken() {
+ Connector delegate = mock(Connector.class);
+ GravitinoConnector connector =
+ connector(context("OAUTH2_PASSTHROUGH", "lakehouse-iceberg",
delegate));
+ ConnectorSession session = mock(ConnectorSession.class);
+ ConnectorTransactionHandle transaction =
mock(ConnectorTransactionHandle.class);
+ RuntimeException failure = new IllegalArgumentException("Missing delegated
token");
+ when(delegate.getMetadata(session, transaction)).thenThrow(failure);
+ ConnectorMetadata metadata = connector.getInternalMetadata(session,
transaction);
+ verifyNoInteractions(delegate);
+ assertSame(
+ failure,
+ assertThrows(IllegalArgumentException.class, () ->
metadata.listSchemaNames(session)));
+ verify(delegate).getMetadata(session, transaction);
+ }
+
+ @Test
+ void otherModesAndProvidersRetainEagerMetadata() {
+ for (String[] mode :
+ new String[][] {{"OAUTH2", "lakehouse-iceberg"},
{"OAUTH2_PASSTHROUGH", "hive"}}) {
+ Connector delegate = mock(Connector.class);
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ ConnectorSession session = mock(ConnectorSession.class);
+ ConnectorTransactionHandle transaction =
mock(ConnectorTransactionHandle.class);
+ when(delegate.getMetadata(session,
transaction)).thenReturn(nativeMetadata);
+ GravitinoConnector connector = connector(context(mode[0], mode[1],
delegate));
+ assertSame(nativeMetadata, connector.getInternalMetadata(session,
transaction));
+ verify(delegate).getMetadata(session, transaction);
+ }
+ }
+
+ @Test
+ void catalogPassthroughDefersMetadata() {
+ assertSecurityPrecedence(null, null, "OAUTH2_PASSTHROUGH", true);
+ }
+
+ @Test
+ void clusterNoneOverridesCatalogPassthrough() {
+ assertSecurityPrecedence(null, "NONE", "OAUTH2_PASSTHROUGH", false);
+ }
+
+ @Test
+ void clusterOAuth2OverridesCatalogPassthrough() {
+ assertSecurityPrecedence(null, "OAUTH2", "OAUTH2_PASSTHROUGH", false);
+ }
+
+ @Test
+ void clusterPassthroughOverridesCatalogNone() {
+ assertSecurityPrecedence(null, "OAUTH2_PASSTHROUGH", "NONE", true);
+ }
+
+ @Test
+ void inheritedClusterOAuth2OverridesCatalogPassthrough() {
+ assertSecurityPrecedence("oauth2", null, "OAUTH2_PASSTHROUGH", false);
+ }
+
+ private void assertSecurityPrecedence(
+ @Nullable String authType,
+ @Nullable String clusterSecurity,
+ String catalogSecurity,
+ boolean deferred) {
+ Map<String, String> config = new HashMap<>();
+ config.put("gravitino.metalake", "demo");
+ if (authType != null) {
+ config.put("gravitino.client.authType", authType);
+ }
+ if (clusterSecurity != null) {
+ config.put("gravitino.iceberg.rest-catalog.security", clusterSecurity);
+ }
+ Connector delegate = mock(Connector.class);
+ ConnectorMetadata nativeMetadata = mock(ConnectorMetadata.class);
+ ConnectorSession session = mock(ConnectorSession.class);
+ ConnectorTransactionHandle transaction =
mock(ConnectorTransactionHandle.class);
+ when(delegate.getMetadata(session,
transaction)).thenReturn(nativeMetadata);
+ GravitinoConnector connector =
+ connector(
+ context(
+ new GravitinoConfig(config),
+ "lakehouse-iceberg",
+ delegate,
+ Map.of("trino.bypass.iceberg.rest-catalog.security",
catalogSecurity)));
+ ConnectorMetadata metadata = connector.getInternalMetadata(session,
transaction);
+ if (deferred) {
+ verifyNoInteractions(delegate);
+ metadata.listSchemaNames(session);
+ verify(nativeMetadata).listSchemaNames(session);
+ } else {
+ assertSame(nativeMetadata, metadata);
+ }
+ verify(delegate).getMetadata(session, transaction);
+ }
+
+ private GravitinoConnector connector(CatalogConnectorContext context) {
+ return new GravitinoConnector(context) {
+ @Override
+ protected GravitinoMetadata createGravitinoMetadata(
+ CatalogConnectorMetadata metadata,
+ CatalogConnectorMetadataAdapter adapter,
+ ConnectorMetadata delegate) {
+ return new GravitinoMetadata(metadata, adapter, delegate) {};
+ }
+ };
+ }
+
+ private CatalogConnectorContext context(String security, String provider,
Connector delegate) {
+ return context(
+ new GravitinoConfig(
+ Map.of(
+ "gravitino.metalake", "demo",
+ "gravitino.client.authType", "oauth2",
+ "gravitino.client.session.forwardUser", "true",
+ "gravitino.iceberg.rest-catalog.security", security)),
+ provider,
+ delegate,
+ Map.of());
+ }
+
+ private CatalogConnectorContext context(
+ GravitinoConfig config,
+ String provider,
+ Connector delegate,
+ Map<String, String> catalogProperties) {
+ CatalogConnectorContext context = mock(CatalogConnectorContext.class);
+ GravitinoCatalog catalog = mock(GravitinoCatalog.class);
+ when(catalog.geNameIdentifier()).thenReturn(NameIdentifier.of("demo",
"iceberg_demo"));
+ when(catalog.getProvider()).thenReturn(provider);
+ when(catalog.getProperties()).thenReturn(catalogProperties);
+ GravitinoMetalake metalake = mock(GravitinoMetalake.class);
+ Catalog live = mock(Catalog.class);
+ SupportsSchemas schemas = mock(SupportsSchemas.class);
+ when(schemas.listSchemas()).thenReturn(new String[] {"demo"});
+ when(live.asSchemas()).thenReturn(schemas);
+ when(metalake.loadCatalog(any())).thenReturn(live);
+ when(context.getCatalog()).thenReturn(catalog);
+ when(context.getMetalake()).thenReturn(metalake);
+ when(context.getInternalConnector()).thenReturn(delegate);
+ when(context.getConfig()).thenReturn(config);
+ Map<String, String> internalConnectorConfig =
+ new IcebergCatalogPropertyConverter()
+ .buildIcebergRestProperties(catalog, config,
"http://localhost:9001/iceberg");
+
when(context.getInternalConnectorConfig()).thenReturn(internalConnectorConfig);
+ return context;
+ }
+}