This is an automated email from the ASF dual-hosted git repository. diqiu50 pushed a commit to branch trino-irc-1.3 in repository https://gitbox.apache.org/repos/asf/gravitino.git
commit e554a3398dd3986734dfcf30236902c42781ed6c Author: yuhui <[email protected]> AuthorDate: Tue Aug 25 13:18:22 2026 +0000 improvement(spark-connector): Reuse OAuth for routed Iceberg REST catalogs --- .../spark/connector/GravitinoSparkConfig.java | 5 + .../connector/iceberg/GravitinoIcebergCatalog.java | 22 +- .../connector/iceberg/IcebergRestOAuthConfig.java | 89 ++++++++ .../iceberg/TestIcebergRestOAuthConfig.java | 123 +++++++++++ spark-connector/v3.5/spark/build.gradle.kts | 2 + ...kIcebergCatalogRestS3CredentialVendingIT35.java | 240 +++++++++++++++++++++ 6 files changed, 471 insertions(+), 10 deletions(-) diff --git a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/GravitinoSparkConfig.java b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/GravitinoSparkConfig.java index 78f8de661e..119c2f9501 100644 --- a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/GravitinoSparkConfig.java +++ b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/GravitinoSparkConfig.java @@ -36,6 +36,11 @@ public class GravitinoSparkConfig { // lakehouse-iceberg catalogs are routed through; takes precedence over auto-discovery. public static final String GRAVITINO_ICEBERG_REST_URI = GRAVITINO_PREFIX + "iceberg.rest-uri"; + // Reuses the Gravitino OAuth2 client configuration for automatically routed Iceberg REST + // catalogs. Iceberg obtains and refreshes its own access token with the same client identity. + public static final String GRAVITINO_ICEBERG_REUSE_OAUTH2 = + GRAVITINO_PREFIX + "iceberg.reuseOAuth2"; + // Pass-through prefix for the Iceberg REST client config (e.g. rest.auth.type, // rest.auth.basic.username), applied when a catalog is routed through the Iceberg REST server. public static final String GRAVITINO_ICEBERG_REST_CONFIG_PREFIX = diff --git a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java index cb0615dfad..1fc5b5b031 100644 --- a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java +++ b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/GravitinoIcebergCatalog.java @@ -43,6 +43,7 @@ import org.apache.iceberg.spark.SparkCatalog; import org.apache.iceberg.spark.procedures.SparkProcedures; import org.apache.iceberg.spark.source.HasIcebergCatalog; import org.apache.iceberg.spark.source.SparkTable; +import org.apache.spark.SparkConf; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.catalyst.analysis.NoSuchFunctionException; import org.apache.spark.sql.catalyst.analysis.NoSuchNamespaceException; @@ -83,7 +84,7 @@ public class GravitinoIcebergCatalog extends BaseCatalog Optional<String> icebergRestUri = resolveIcebergRestUri(properties); Map<String, String> all; if (icebergRestUri.isPresent()) { - all = buildIcebergRestSparkCatalogProperties(name, options, properties, icebergRestUri.get()); + all = buildAutoRoutedIcebergRestProperties(name, options, properties, icebergRestUri.get()); } else { all = getPropertiesConverter().toSparkCatalogProperties(options, properties); CredentialPropertyUtils.applyIcebergCredentials( @@ -121,7 +122,7 @@ public class GravitinoIcebergCatalog extends BaseCatalog return GravitinoCatalogManager.get().getIcebergRestUri(); } - private Map<String, String> buildIcebergRestSparkCatalogProperties( + private Map<String, String> buildAutoRoutedIcebergRestProperties( String gravitinoCatalogName, CaseInsensitiveStringMap options, Map<String, String> properties, @@ -130,20 +131,21 @@ public class GravitinoIcebergCatalog extends BaseCatalog Map<String, String> all = new HashMap<>( converter.buildIcebergRestProperties( - gravitinoCatalogName, restUri, properties, getIcebergRestClientConfig())); + gravitinoCatalogName, restUri, properties, getAutoRoutedIcebergRestClientConfig())); if (options != null) { all.putAll(options); } return all; } - private Map<String, String> getIcebergRestClientConfig() { - return Stream.of( - SparkSession.active() - .sparkContext() - .conf() - .getAllWithPrefix(GravitinoSparkConfig.GRAVITINO_ICEBERG_REST_CONFIG_PREFIX)) - .collect(Collectors.toMap(t -> t._1, t -> t._2, (oldVal, newVal) -> newVal)); + private Map<String, String> getAutoRoutedIcebergRestClientConfig() { + SparkConf sparkConf = SparkSession.active().sparkContext().conf(); + Map<String, String> explicitRestConfig = + Stream.of( + sparkConf.getAllWithPrefix( + GravitinoSparkConfig.GRAVITINO_ICEBERG_REST_CONFIG_PREFIX)) + .collect(Collectors.toMap(t -> t._1, t -> t._2, (oldVal, newVal) -> newVal)); + return IcebergRestOAuthConfig.resolve(sparkConf, explicitRestConfig); } @Override diff --git a/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/IcebergRestOAuthConfig.java b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/IcebergRestOAuthConfig.java new file mode 100644 index 0000000000..92b548ca00 --- /dev/null +++ b/spark-connector/spark-common/src/main/java/org/apache/gravitino/spark/connector/iceberg/IcebergRestOAuthConfig.java @@ -0,0 +1,89 @@ +/* + * 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.spark.connector.iceberg; + +import com.google.common.base.Preconditions; +import com.google.common.collect.ImmutableSet; +import java.util.HashMap; +import java.util.Map; +import java.util.Set; +import org.apache.commons.lang3.StringUtils; +import org.apache.gravitino.auth.AuthProperties; +import org.apache.gravitino.spark.connector.GravitinoSparkConfig; +import org.apache.spark.SparkConf; + +class IcebergRestOAuthConfig { + + static final String AUTH_TYPE = "rest.auth.type"; + static final String AUTH_TYPE_OAUTH2 = "oauth2"; + static final String CREDENTIAL = "credential"; + static final String OAUTH2_SERVER_URI = "oauth2-server-uri"; + static final String SCOPE = "scope"; + + private static final Set<String> LEGACY_OAUTH_PROPERTIES = + ImmutableSet.of( + "token", + CREDENTIAL, + SCOPE, + OAUTH2_SERVER_URI, + "audience", + "resource", + "token-refresh-enabled", + "token-exchange-enabled"); + + private IcebergRestOAuthConfig() {} + + static Map<String, String> resolve(SparkConf sparkConf, Map<String, String> explicitRestConfig) { + Map<String, String> result = new HashMap<>(explicitRestConfig); + if (hasExplicitAuthentication(result) + || !sparkConf.getBoolean(GravitinoSparkConfig.GRAVITINO_ICEBERG_REUSE_OAUTH2, true)) { + return result; + } + + String authType = + sparkConf.get(GravitinoSparkConfig.GRAVITINO_AUTH_TYPE, AuthProperties.SIMPLE_AUTH_TYPE); + if (!AuthProperties.isOAuth2(authType)) { + return result; + } + + String serverUri = required(sparkConf, GravitinoSparkConfig.GRAVITINO_OAUTH2_URI); + String tokenPath = required(sparkConf, GravitinoSparkConfig.GRAVITINO_OAUTH2_PATH); + result.put(AUTH_TYPE, AUTH_TYPE_OAUTH2); + result.put(CREDENTIAL, required(sparkConf, GravitinoSparkConfig.GRAVITINO_OAUTH2_CREDENTIAL)); + result.put(SCOPE, required(sparkConf, GravitinoSparkConfig.GRAVITINO_OAUTH2_SCOPE)); + result.put(OAUTH2_SERVER_URI, joinUri(serverUri, tokenPath)); + return result; + } + + private static boolean hasExplicitAuthentication(Map<String, String> restConfig) { + return restConfig.keySet().stream() + .anyMatch(key -> key.startsWith("rest.auth.") || LEGACY_OAUTH_PROPERTIES.contains(key)); + } + + private static String required(SparkConf sparkConf, String key) { + String value = sparkConf.get(key, null); + Preconditions.checkArgument(StringUtils.isNotBlank(value), key + " should not be empty"); + return value; + } + + private static String joinUri(String serverUri, String tokenPath) { + return StringUtils.removeEnd(serverUri, "/") + "/" + StringUtils.removeStart(tokenPath, "/"); + } +} diff --git a/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestIcebergRestOAuthConfig.java b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestIcebergRestOAuthConfig.java new file mode 100644 index 0000000000..eff21194ba --- /dev/null +++ b/spark-connector/spark-common/src/test/java/org/apache/gravitino/spark/connector/iceberg/TestIcebergRestOAuthConfig.java @@ -0,0 +1,123 @@ +/* + * 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.spark.connector.iceberg; + +import com.google.common.collect.ImmutableMap; +import com.google.common.collect.ImmutableSet; +import java.util.Collections; +import java.util.Map; +import org.apache.gravitino.spark.connector.GravitinoSparkConfig; +import org.apache.spark.SparkConf; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +/** Tests automatic reuse of Gravitino OAuth2 client settings for Iceberg REST catalogs. */ +public class TestIcebergRestOAuthConfig { + + @Test + void testDerivesIcebergOAuthConfig() { + SparkConf sparkConf = oauthSparkConf("https://identity.example.com/", "/oauth/token"); + + Map<String, String> result = IcebergRestOAuthConfig.resolve(sparkConf, Collections.emptyMap()); + + Assertions.assertEquals("oauth2", result.get(IcebergRestOAuthConfig.AUTH_TYPE)); + Assertions.assertEquals("alice:secret", result.get(IcebergRestOAuthConfig.CREDENTIAL)); + Assertions.assertEquals("openid", result.get(IcebergRestOAuthConfig.SCOPE)); + Assertions.assertEquals( + "https://identity.example.com/oauth/token", + result.get(IcebergRestOAuthConfig.OAUTH2_SERVER_URI)); + } + + @Test + void testExplicitRestAuthenticationTakesPrecedence() { + SparkConf sparkConf = oauthSparkConf("https://identity.example.com", "oauth/token"); + Map<String, String> explicit = + ImmutableMap.of("rest.auth.type", "basic", "rest.auth.basic.username", "admin"); + + Map<String, String> result = IcebergRestOAuthConfig.resolve(sparkConf, explicit); + + Assertions.assertEquals(explicit, result); + } + + @Test + void testLegacyExplicitOAuthPropertiesDisableAutomaticReuse() { + SparkConf sparkConf = oauthSparkConf("https://identity.example.com", "oauth/token"); + + for (String property : + ImmutableSet.of( + "token", + "credential", + "scope", + "oauth2-server-uri", + "audience", + "resource", + "token-refresh-enabled", + "token-exchange-enabled")) { + Map<String, String> explicit = ImmutableMap.of(property, "explicit-value"); + + Map<String, String> result = IcebergRestOAuthConfig.resolve(sparkConf, explicit); + + Assertions.assertEquals(explicit, result, property); + } + } + + @Test + void testNonAuthenticationRestPropertiesDoNotDisableAutomaticReuse() { + SparkConf sparkConf = oauthSparkConf("https://identity.example.com", "oauth/token"); + + Map<String, String> result = + IcebergRestOAuthConfig.resolve( + sparkConf, ImmutableMap.of("header.X-Iceberg-Custom", "value")); + + Assertions.assertEquals("value", result.get("header.X-Iceberg-Custom")); + Assertions.assertEquals("oauth2", result.get(IcebergRestOAuthConfig.AUTH_TYPE)); + Assertions.assertEquals("alice:secret", result.get(IcebergRestOAuthConfig.CREDENTIAL)); + } + + @Test + void testCanDisableOAuthReuse() { + SparkConf sparkConf = oauthSparkConf("https://identity.example.com", "oauth/token"); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_ICEBERG_REUSE_OAUTH2, "false"); + + Map<String, String> result = IcebergRestOAuthConfig.resolve(sparkConf, Collections.emptyMap()); + + Assertions.assertTrue(result.isEmpty()); + } + + @Test + void testDoesNotDeriveConfigForOtherAuthenticationTypes() { + SparkConf sparkConf = new SparkConf(false); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_AUTH_TYPE, "simple"); + + Map<String, String> result = IcebergRestOAuthConfig.resolve(sparkConf, Collections.emptyMap()); + + Assertions.assertTrue(result.isEmpty()); + } + + private SparkConf oauthSparkConf(String serverUri, String tokenPath) { + SparkConf sparkConf = new SparkConf(false); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_AUTH_TYPE, "oauth2"); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_OAUTH2_URI, serverUri); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_OAUTH2_PATH, tokenPath); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_OAUTH2_CREDENTIAL, "alice:secret"); + sparkConf.set(GravitinoSparkConfig.GRAVITINO_OAUTH2_SCOPE, "openid"); + return sparkConf; + } +} diff --git a/spark-connector/v3.5/spark/build.gradle.kts b/spark-connector/v3.5/spark/build.gradle.kts index 4b2a7b7402..d59088d5bf 100644 --- a/spark-connector/v3.5/spark/build.gradle.kts +++ b/spark-connector/v3.5/spark/build.gradle.kts @@ -151,6 +151,8 @@ dependencies { testImplementation(libs.mysql.driver) testImplementation(libs.postgresql.driver) testImplementation(libs.testcontainers) + testImplementation(libs.aws.policy) + testImplementation(project(":iceberg:iceberg-common")) // org.apache.iceberg.rest.RESTSerializers#registerAll(ObjectMapper) has different method signature for iceberg-core and iceberg-spark-runtime package, we must make sure iceberg-core is in front to start up MiniGravitino server. testImplementation("org.apache.iceberg:iceberg-core:$icebergVersion") diff --git a/spark-connector/v3.5/spark/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogRestS3CredentialVendingIT35.java b/spark-connector/v3.5/spark/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogRestS3CredentialVendingIT35.java new file mode 100644 index 0000000000..a118b43bf9 --- /dev/null +++ b/spark-connector/v3.5/spark/src/test/java/org/apache/gravitino/spark/connector/integration/test/iceberg/SparkIcebergCatalogRestS3CredentialVendingIT35.java @@ -0,0 +1,240 @@ +/* + * 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.spark.connector.integration.test.iceberg; + +import com.google.common.collect.ImmutableList; +import com.google.common.collect.Maps; +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.sql.SQLException; +import java.time.Instant; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import org.apache.gravitino.Catalog; +import org.apache.gravitino.Configs; +import org.apache.gravitino.MetadataObject; +import org.apache.gravitino.MetadataObjects; +import org.apache.gravitino.auth.AuthenticatorType; +import org.apache.gravitino.authorization.Owner; +import org.apache.gravitino.authorization.Privileges; +import org.apache.gravitino.authorization.SecurableObject; +import org.apache.gravitino.authorization.SecurableObjects; +import org.apache.gravitino.client.GravitinoMetalake; +import org.apache.gravitino.integration.test.container.ContainerSuite; +import org.apache.gravitino.integration.test.container.PostgreSQLContainer; +import org.apache.gravitino.integration.test.util.BaseIT; +import org.apache.gravitino.integration.test.util.JwksMockServerHelper; +import org.apache.gravitino.integration.test.util.OAuthMockDataProvider; +import org.apache.gravitino.integration.test.util.TestDatabaseName; +import org.apache.gravitino.server.authentication.OAuthConfig; +import org.apache.gravitino.spark.connector.GravitinoSparkConfig; +import org.apache.gravitino.spark.connector.iceberg.GravitinoIcebergCatalog; +import org.apache.gravitino.spark.connector.plugin.GravitinoSparkPlugin; +import org.apache.spark.SparkConf; +import org.apache.spark.sql.Row; +import org.apache.spark.sql.SparkSession; +import org.apache.spark.sql.connector.catalog.CatalogPlugin; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.condition.EnabledIfEnvironmentVariable; + +/** Spark 3.5 integration test for OAuth2 IRC routing with AWS credential vending. */ +@Tag("gravitino-docker-test") +@EnabledIfEnvironmentVariable(named = "GRAVITINO_TEST_CLOUD_IT", matches = "true") +public class SparkIcebergCatalogRestS3CredentialVendingIT35 extends BaseIT { + + private static final String ADMIN = "gravitino"; + private static final String ALICE = "alice"; + private static final String ALICE_CREDENTIAL = "alice:alice-secret"; + private static final String AUDIENCE = "service1"; + private static final String METALAKE = "spark35_irc_oauth"; + private static final String CATALOG = "iceberg_pg"; + private static final String SCHEMA = "aws_vending"; + private static final String TABLE = "spark35_test"; + private static final String ROLE = "spark35-irc-role"; + private static final String AWS_ROLE_ARN = + "arn:aws:iam::730335553010:role/gravitino-irc-s3-vending-demo"; + + private final ContainerSuite containerSuite = ContainerSuite.getInstance(); + + private JwksMockServerHelper mockServer; + private SparkSession spark; + + @BeforeAll + @Override + public void startIntegrationTest() throws Exception { + initializeOAuthServer(); + initializeGravitinoConfig(); + initializePostgreSQL(); + super.startIntegrationTest(); + initializeMetadata(); + spark = createSparkSession(); + } + + @AfterAll + @Override + public void stopIntegrationTest() throws IOException, InterruptedException { + if (spark != null) { + spark.stop(); + } + if (mockServer != null) { + mockServer.close(); + } + super.stopIntegrationTest(); + } + + @Test + void testPostgreSQLCatalogRoutesThroughOAuthIrcAndVendsS3Credentials() { + CatalogPlugin catalogPlugin = spark.sessionState().catalogManager().catalog(CATALOG); + Assertions.assertInstanceOf(GravitinoIcebergCatalog.class, catalogPlugin); + org.apache.iceberg.catalog.Catalog icebergCatalog = + ((GravitinoIcebergCatalog) catalogPlugin).icebergCatalog(); + Assertions.assertEquals( + "org.apache.iceberg.rest.RESTCatalog", icebergCatalog.getClass().getName()); + + spark.sql(String.format("DROP TABLE IF EXISTS %s.%s.%s", CATALOG, SCHEMA, TABLE)); + spark.sql( + String.format( + "CREATE TABLE %s.%s.%s (id BIGINT, data STRING) USING iceberg", + CATALOG, SCHEMA, TABLE)); + spark.sql( + String.format( + "INSERT INTO %s.%s.%s VALUES (1, 'one'), (2, 'two')", CATALOG, SCHEMA, TABLE)); + + List<Row> rows = + spark + .sql(String.format("SELECT id, data FROM %s.%s.%s ORDER BY id", CATALOG, SCHEMA, TABLE)) + .collectAsList(); + Assertions.assertEquals(2, rows.size()); + Assertions.assertEquals(1L, rows.get(0).getLong(0)); + Assertions.assertEquals("one", rows.get(0).getString(1)); + Assertions.assertEquals(2L, rows.get(1).getLong(0)); + Assertions.assertEquals("two", rows.get(1).getString(1)); + + Optional<Owner> owner = + client + .loadMetalake(METALAKE) + .getOwner( + MetadataObjects.of( + ImmutableList.of(CATALOG, SCHEMA, TABLE), MetadataObject.Type.TABLE)); + Assertions.assertTrue(owner.isPresent()); + Assertions.assertEquals(ALICE, owner.get().name()); + + spark.sql(String.format("DROP TABLE %s.%s.%s PURGE", CATALOG, SCHEMA, TABLE)); + } + + private void initializeOAuthServer() throws Exception { + mockServer = JwksMockServerHelper.create("spark35-irc-kid"); + Instant expiration = Instant.now().plusSeconds(3600); + String adminToken = mockServer.mintToken(ADMIN, AUDIENCE, expiration); + String aliceToken = mockServer.mintToken(ALICE, AUDIENCE, expiration); + mockServer.registerUserToken(ADMIN, adminToken); + mockServer.registerUserToken(ALICE, aliceToken); + mockServer.setFallbackToken(adminToken); + OAuthMockDataProvider.getInstance().setTokenData(adminToken.getBytes(StandardCharsets.UTF_8)); + } + + private void initializeGravitinoConfig() { + ignoreIcebergAuxRestService = false; + Map<String, String> configs = Maps.newHashMap(); + configs.put(Configs.AUTHENTICATORS.getKey(), AuthenticatorType.OAUTH.name().toLowerCase()); + configs.put(OAuthConfig.SERVICE_AUDIENCE.getKey(), AUDIENCE); + configs.put( + OAuthConfig.TOKEN_VALIDATOR_CLASS.getKey(), + "org.apache.gravitino.server.authentication.JwksTokenValidator"); + configs.put(OAuthConfig.JWKS_URI.getKey(), mockServer.jwksUri()); + configs.put(OAuthConfig.PRINCIPAL_FIELDS.getKey(), "sub"); + configs.put(Configs.ENABLE_AUTHORIZATION.getKey(), "true"); + configs.put(Configs.SERVICE_ADMINS.getKey(), ADMIN); + configs.put("gravitino.iceberg-rest.catalog-config-provider", "dynamic-config-provider"); + configs.put("gravitino.iceberg-rest.gravitino-metalake", METALAKE); + registerCustomConfigs(configs); + } + + private void initializePostgreSQL() { + containerSuite.startPostgreSQLContainer( + TestDatabaseName.PG_TEST_ICEBERG_CATALOG_MULTIPLE_JDBC_LOAD); + } + + private void initializeMetadata() throws SQLException { + PostgreSQLContainer postgres = containerSuite.getPostgreSQLContainer(); + TestDatabaseName database = TestDatabaseName.PG_TEST_ICEBERG_CATALOG_MULTIPLE_JDBC_LOAD; + client.createMetalake(METALAKE, "", new HashMap<>()); + GravitinoMetalake metalake = client.loadMetalake(METALAKE); + metalake.addUser(ALICE); + + Map<String, String> properties = Maps.newHashMap(); + properties.put("catalog-backend", "jdbc"); + properties.put("uri", postgres.getJdbcUrl(database)); + properties.put("jdbc-driver", postgres.getDriverClassName(database)); + properties.put("jdbc-user", postgres.getUsername()); + properties.put("jdbc-password", postgres.getPassword()); + properties.put( + "warehouse", + String.format( + "s3://%s/gravitino-irc-oauth-demo/spark-3.5-pg", System.getenv("AWS_S3_TEST_BUCKET"))); + properties.put("credential-providers", "s3-token"); + properties.put("s3-access-key-id", System.getenv("AWS_ACCESS_KEY_ID")); + properties.put("s3-secret-access-key", System.getenv("AWS_SECRET_ACCESS_KEY")); + properties.put("s3-region", System.getenv("AWS_DEFAULT_REGION")); + properties.put("s3-role-arn", AWS_ROLE_ARN); + properties.put("io-impl", "org.apache.iceberg.aws.s3.S3FileIO"); + Catalog catalog = + metalake.createCatalog( + CATALOG, Catalog.Type.RELATIONAL, "lakehouse-iceberg", "", properties); + catalog.asSchemas().createSchema(SCHEMA, "", new HashMap<>()); + + SecurableObject access = + SecurableObjects.ofCatalog( + CATALOG, + ImmutableList.of( + Privileges.UseCatalog.allow(), + Privileges.UseSchema.allow(), + Privileges.CreateTable.allow(), + Privileges.ModifyTable.allow(), + Privileges.SelectTable.allow())); + metalake.createRole(ROLE, new HashMap<>(), ImmutableList.of(access)); + metalake.grantRolesToUser(ImmutableList.of(ROLE), ALICE); + } + + private SparkSession createSparkSession() { + SparkConf conf = + new SparkConf() + .set("spark.plugins", GravitinoSparkPlugin.class.getName()) + .set(GravitinoSparkConfig.GRAVITINO_URI, serverUri) + .set(GravitinoSparkConfig.GRAVITINO_METALAKE, METALAKE) + .set(GravitinoSparkConfig.GRAVITINO_ENABLE_ICEBERG_SUPPORT, "true") + .set(GravitinoSparkConfig.GRAVITINO_AUTH_TYPE, "oauth2") + .set(GravitinoSparkConfig.GRAVITINO_OAUTH2_URI, mockServer.baseUri()) + .set(GravitinoSparkConfig.GRAVITINO_OAUTH2_PATH, "token") + .set(GravitinoSparkConfig.GRAVITINO_OAUTH2_CREDENTIAL, ALICE_CREDENTIAL) + .set(GravitinoSparkConfig.GRAVITINO_OAUTH2_SCOPE, "openid"); + return SparkSession.builder() + .master("local[1]") + .appName("SparkIcebergCatalogRestS3CredentialVendingIT35") + .config(conf) + .getOrCreate(); + } +}
