This is an automated email from the ASF dual-hosted git repository.

jerryshao pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new 422af4130b [#13194] improvement(trino-connector): Defer Iceberg REST 
passthrough metadata initialization (#13195)
422af4130b is described below

commit 422af4130bd482ea62ced40d3948ead302846cea
Author: Yuhui <[email protected]>
AuthorDate: Wed Sep 16 15:06:12 2026 +0800

    [#13194] improvement(trino-connector): Defer Iceberg REST passthrough 
metadata initialization (#13195)
    
    ### What changes were proposed in this pull request?
    
    - Defer native metadata initialization for `lakehouse-iceberg` catalogs
    using `OAUTH2_PASSTHROUGH`.
    - Use the operation session when initializing native metadata, preserve
    query lifecycle ordering, and propagate native errors.
    - Inherit Gravitino service OAuth2 properties only when REST security is
    `OAUTH2`.
    - Add 18 regression tests covering metadata initialization, session
    handling, repeated query lifecycles, cleanup after initialization
    failure, and cluster-over-catalog configuration precedence.
    
    ### Why are the changes needed?
    
    Metadata-only management operations can be served through Gravitino
    without accessing native Iceberg REST metadata. Eager native metadata
    initialization can nevertheless trigger authentication in the native
    connector, requiring a delegated user token before these operations run.
    
    The configuration code also passes service OAuth2 defaults to explicitly
    configured non-OAuth2 REST security modes.
    
    Fix: #13194
    
    ### Does this PR introduce _any_ user-facing change?
    
    For managed Iceberg catalogs configured with `OAUTH2_PASSTHROUGH`,
    metadata-only management operations that are served through Gravitino no
    longer initialize native metadata. Operations requiring native metadata
    still use the operation session and propagate authentication failures.
    
    Explicit non-OAuth2 REST security modes no longer inherit Gravitino
    service OAuth2 properties. Explicitly configured REST properties remain
    effective.
    
    No configuration keys are added. The native connector must support the
    selected REST authentication mode; this change does not add
    `OAUTH2_PASSTHROUGH` support to upstream Trino.
    
    ### How was this patch tested?
    
    - `./gradlew spotlessApply :trino-connector:trino-connector:test
    -PskipITs`: passed; 338 tests, 0 failures, 1 skipped, including 18 new
    regression tests.
    - `./gradlew :trino-connector:trino-connector:javadoc -PskipITs`: passed
    with existing missing-comment warnings.
    - `git diff --check`: passed.
    - JaCoCo emitted JDK instrumentation compatibility warnings; coverage
    may be incomplete.
    - Authentication behavior was validated with mocked native connectors in
    unit tests; no end-to-end passthrough authentication tests were run.
---
 .../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 6bad499794..dcf663546d 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;
+  }
+}

Reply via email to