From 93f7fc7be73343d887a6f3fa0238ebda73186678 Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Fri, 25 Sep 2026 15:08:24 +0800 Subject: [PATCH 1/6] [oss] Support credentials provider so open streams pick up refreshed REST tokens RESTTokenFileIO bakes the vended token into the delegate FileIO's options, so a stream that is already open keeps signing with the token it was opened with and fails once that token expires. OSSFileIO now resolves credentials per request through a provider that reads the current token from RESTTokenFileIO via CredentialsSupplierRegistry. OSSLoader also accepts fs.oss.credentials.provider in place of the access keys. --- docs/docs/maintenance/filesystems.mdx | 6 + .../fs/CredentialsSupplierRegistry.java | 57 ++++++ .../apache/paimon/rest/RESTTokenFileIO.java | 48 ++++- .../paimon/rest/RESTTokenFileIOTest.java | 50 +++++ .../java/org/apache/paimon/oss/OSSFileIO.java | 16 +- .../oss/RegisteredCredentialsProvider.java | 129 ++++++++++++ .../RegisteredCredentialsProviderTest.java | 185 ++++++++++++++++++ .../java/org/apache/paimon/oss/OSSLoader.java | 4 +- .../org/apache/paimon/oss/OSSLoaderTest.java | 46 +++++ 9 files changed, 528 insertions(+), 13 deletions(-) create mode 100644 paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java create mode 100644 paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java create mode 100644 paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java create mode 100644 paimon-filesystems/paimon-oss/src/test/java/org/apache/paimon/oss/OSSLoaderTest.java diff --git a/docs/docs/maintenance/filesystems.mdx b/docs/docs/maintenance/filesystems.mdx index 54f0817f9991..f708a73df667 100644 --- a/docs/docs/maintenance/filesystems.mdx +++ b/docs/docs/maintenance/filesystems.mdx @@ -333,6 +333,12 @@ Download [paimon-jindo-@@VERSION@@.jar](https://repository.apache.org/snapshots/ +### Credentials Provider + +Instead of `fs.oss.accessKeyId` and `fs.oss.accessKeySecret`, you can set `fs.oss.credentials.provider` to a class implementing `com.aliyun.oss.common.auth.CredentialsProvider`, for example `com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider`. The provider is asked for credentials on every request. + +With a REST catalog that vends data tokens, the OSS FileIO signs every request with the catalog's current token, so streams that are already open keep working after the token is refreshed. + ### Server-Side Encryption Paimon can stamp OSS server-side-encryption headers on the writes it performs (`PutObject`, server-side `CopyObject` used by rename/commit, and multipart-upload initiation) via three keys that map 1:1 to the OSS request headers: diff --git a/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java b/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java new file mode 100644 index 000000000000..cd1dcb336a78 --- /dev/null +++ b/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java @@ -0,0 +1,57 @@ +/* + * 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.paimon.fs; + +import javax.annotation.Nullable; + +import java.util.Map; +import java.util.UUID; +import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Supplier; + +/** + * Credential suppliers that a plugin {@link FileIO} looks up by id, so it can resolve fresh + * credentials on every request instead of the ones it was created with. + */ +public final class CredentialsSupplierRegistry { + + /** Option that passes the id of a registered supplier to the {@link FileIO} being created. */ + public static final String SUPPLIER_ID = "fs.credentials-supplier.id"; + + private static final Map>> SUPPLIERS = + new ConcurrentHashMap<>(); + + private CredentialsSupplierRegistry() {} + + /** Registers a supplier of credential options and returns its id. */ + public static String register(Supplier> supplier) { + String id = UUID.randomUUID().toString(); + SUPPLIERS.put(id, supplier); + return id; + } + + @Nullable + public static Supplier> get(String id) { + return SUPPLIERS.get(id); + } + + public static void unregister(String id) { + SUPPLIERS.remove(id); + } +} diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java index 830fe058c43c..0e6ebe852ff6 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java @@ -21,6 +21,7 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.data.BlobDescriptor; +import org.apache.paimon.fs.CredentialsSupplierRegistry; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.FileStatus; import org.apache.paimon.fs.Path; @@ -67,12 +68,12 @@ public class RESTTokenFileIO implements FileIO { .defaultValue(false) .withDescription("Whether to support data token provided by the REST server."); - private static final Cache FILE_IO_CACHE = + private static final Cache FILE_IO_CACHE = Caffeine.newBuilder() .maximumSize(1000) .expireAfterAccess(10, TimeUnit.HOURS) .removalListener( - (ignored, value, cause) -> IOUtils.closeQuietly((FileIO) value)) + (ignored, value, cause) -> IOUtils.closeQuietly((CachedFileIO) value)) .scheduler( Scheduler.forScheduledExecutorService( Executors.newSingleThreadScheduledExecutor( @@ -245,28 +246,37 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc + "REST credential lifetime after refresh."); } - FileIO fileIO = FILE_IO_CACHE.getIfPresent(currentToken); - if (fileIO != null) { - return new FileIOWithToken(fileIO, currentToken); + CachedFileIO cached = FILE_IO_CACHE.getIfPresent(currentToken); + if (cached != null) { + return new FileIOWithToken(cached.fileIO, currentToken); } synchronized (FILE_IO_CACHE) { - fileIO = FILE_IO_CACHE.getIfPresent(currentToken); - if (fileIO != null) { - return new FileIOWithToken(fileIO, currentToken); + cached = FILE_IO_CACHE.getIfPresent(currentToken); + if (cached != null) { + return new FileIOWithToken(cached.fileIO, currentToken); } + // Lets a FileIO that supports it, such as OSS, sign each request with a fresh token. + String supplierId = CredentialsSupplierRegistry.register(() -> validToken().token()); Options options = catalogContext.options(); options = new Options(RESTUtil.merge(options.toMap(), currentToken.token())); options.set(FILE_IO_ALLOW_CACHE, false); + options.set(CredentialsSupplierRegistry.SUPPLIER_ID, supplierId); CatalogContext context = CatalogContext.create( options, catalogContext.hadoopConf(), catalogContext.preferIO(), catalogContext.fallbackIO()); - fileIO = FileIO.get(path, context); - FILE_IO_CACHE.put(currentToken, fileIO); + FileIO fileIO; + try { + fileIO = FileIO.get(path, context); + } catch (IOException | RuntimeException e) { + CredentialsSupplierRegistry.unregister(supplierId); + throw e; + } + FILE_IO_CACHE.put(currentToken, new CachedFileIO(fileIO, supplierId)); return new FileIOWithToken(fileIO, currentToken); } } @@ -300,6 +310,24 @@ private boolean shouldRefresh(long minimumValidityMillis) { < Math.max(TOKEN_EXPIRATION_SAFE_TIME_MILLIS, minimumValidityMillis); } + /** A delegate {@link FileIO} and the credentials supplier registered for it. */ + private static class CachedFileIO implements AutoCloseable { + + private final FileIO fileIO; + private final String supplierId; + + private CachedFileIO(FileIO fileIO, String supplierId) { + this.fileIO = fileIO; + this.supplierId = supplierId; + } + + @Override + public void close() throws Exception { + CredentialsSupplierRegistry.unregister(supplierId); + fileIO.close(); + } + } + private static class FileIOWithToken { private final FileIO fileIO; diff --git a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java index e80f097fcf2e..e577a5594883 100644 --- a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java @@ -21,6 +21,7 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.data.BlobDescriptor; +import org.apache.paimon.fs.CredentialsSupplierRegistry; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.FileIOLoader; import org.apache.paimon.fs.FileStatus; @@ -30,16 +31,20 @@ import org.apache.paimon.rest.responses.GetTableTokenResponse; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; import java.io.IOException; import java.time.Duration; import java.util.Collections; +import java.util.Map; import java.util.UUID; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Supplier; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.atLeastOnce; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; @@ -226,6 +231,51 @@ long currentTimeMillis() { verify(api, times(2)).loadTableToken(identifier); } + @Test + void testDelegateCredentialsSupplierFollowsTokenRefresh() throws IOException { + Path root = new Path("oss://bucket/table"); + AtomicLong now = new AtomicLong(1700000000000L); + FileIO delegate = mock(FileIO.class); + when(delegate.exists(any())).thenReturn(true); + FileIOLoader loader = mock(FileIOLoader.class); + when(loader.load(any())).thenReturn(delegate); + when(loader.getScheme()).thenReturn("oss"); + RESTApi api = mock(RESTApi.class); + Identifier identifier = Identifier.create("db", "table"); + String first = UUID.randomUUID().toString(); + String second = UUID.randomUUID().toString(); + when(api.loadTableToken(identifier)) + .thenReturn( + new GetTableTokenResponse( + Collections.singletonMap("test.token", first), + now.get() + Duration.ofHours(2).toMillis()), + new GetTableTokenResponse( + Collections.singletonMap("test.token", second), + now.get() + Duration.ofHours(4).toMillis())); + RESTTokenFileIO fileIO = + new RESTTokenFileIO( + CatalogContext.create(new Options(), loader, null), api, identifier, root) { + @Override + long currentTimeMillis() { + return now.get(); + } + }; + + fileIO.exists(root); + ArgumentCaptor context = ArgumentCaptor.forClass(CatalogContext.class); + verify(delegate, atLeastOnce()).configure(context.capture()); + String supplierId = + context.getValue().options().get(CredentialsSupplierRegistry.SUPPLIER_ID); + Supplier> supplier = CredentialsSupplierRegistry.get(supplierId); + assertThat(supplier).isNotNull(); + assertThat(supplier.get()).containsEntry("test.token", first); + + // 30 minutes left is inside the safe window, so the delegate is handed a new token + now.addAndGet(Duration.ofMinutes(90).toMillis()); + assertThat(supplier.get()).containsEntry("test.token", second); + verify(api, times(2)).loadTableToken(identifier); + } + @Test void testFileIOCreationFailureSurfacesAsCheckedIOException() throws IOException { Path tableRoot = new Path("resttoken-broken://bucket/table"); diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java index 3249cbfcef77..9f5be43fd3a8 100644 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java @@ -20,6 +20,7 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.data.BlobDescriptor; +import org.apache.paimon.fs.CredentialsSupplierRegistry; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.HadoopOptionsProvider; import org.apache.paimon.fs.Path; @@ -51,6 +52,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import javax.annotation.Nullable; + import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.UncheckedIOException; @@ -84,6 +87,7 @@ public class OSSFileIO extends HadoopCompliantFileIO implements HadoopOptionsPro private static final String OSS_ACCESS_KEY_SECRET = "fs.oss.accessKeySecret"; private static final String OSS_SECURITY_TOKEN = "fs.oss.securityToken"; private static final String OSS_SECOND_LEVEL_DOMAIN_ENABLED = "fs.oss.sld.enabled"; + private static final String OSS_CREDENTIALS_PROVIDER = "fs.oss.credentials.provider"; /** * Set to false for an OSS-compatible endpoint that is neither an official Aliyun domain nor a @@ -133,6 +137,7 @@ public class OSSFileIO extends HadoopCompliantFileIO implements HadoopOptionsPro private Options hadoopOptions; private boolean allowCache = true; + @Nullable private String credentialsSupplierId; @Override public boolean isObjectStore() { @@ -141,7 +146,9 @@ public boolean isObjectStore() { @Override public void configure(CatalogContext context) { - allowCache = context.options().get(FILE_IO_ALLOW_CACHE); + credentialsSupplierId = context.options().get(CredentialsSupplierRegistry.SUPPLIER_ID); + // The file system is bound to the supplier, so it must not be shared through the cache. + allowCache = context.options().get(FILE_IO_ALLOW_CACHE) && credentialsSupplierId == null; hadoopOptions = new Options(); // read all configuration with prefix 'CONFIG_PREFIXES' for (String key : context.options().keySet()) { @@ -195,6 +202,13 @@ protected AliyunOSSFileSystem createFileSystem(org.apache.hadoop.fs.Path path) { // retrieve props from the file, which comes at a high cost Configuration hadoopConf = new Configuration(SHARED_CONFIG); hadoopOptions.toMap().forEach(hadoopConf::set); + if (credentialsSupplierId != null) { + hadoopConf.set( + OSS_CREDENTIALS_PROVIDER, + RegisteredCredentialsProvider.class.getName()); + hadoopConf.set( + CredentialsSupplierRegistry.SUPPLIER_ID, credentialsSupplierId); + } URI fsUri = path.toUri(); if (scheme == null && authority == null) { fsUri = FileSystem.getDefaultUri(hadoopConf); diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java new file mode 100644 index 000000000000..f34f2498d7a3 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java @@ -0,0 +1,129 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.fs.CredentialsSupplierRegistry; + +import com.aliyun.oss.common.auth.Credentials; +import com.aliyun.oss.common.auth.CredentialsProvider; +import com.aliyun.oss.common.auth.DefaultCredentials; +import org.apache.hadoop.conf.Configuration; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.net.URI; +import java.util.Locale; +import java.util.Map; +import java.util.function.Supplier; + +/** + * An OSS {@link CredentialsProvider} that asks a supplier from {@link CredentialsSupplierRegistry} + * for credentials on every request, so streams that are already open also pick up refreshed ones. + */ +public class RegisteredCredentialsProvider implements CredentialsProvider { + + private static final Logger LOG = LoggerFactory.getLogger(RegisteredCredentialsProvider.class); + + static final String ACCESS_KEY_ID = "fs.oss.accessKeyId"; + static final String ACCESS_KEY_SECRET = "fs.oss.accessKeySecret"; + static final String SECURITY_TOKEN = "fs.oss.securityToken"; + + private final String supplierId; + + // The last options from the supplier and the credentials built from them, swapped together. + private volatile Resolved last; + + public RegisteredCredentialsProvider(URI uri, Configuration conf) { + this.supplierId = conf.get(CredentialsSupplierRegistry.SUPPLIER_ID); + String accessKeyId = conf.get(ACCESS_KEY_ID); + String accessKeySecret = conf.get(ACCESS_KEY_SECRET); + if (accessKeyId != null && accessKeySecret != null) { + this.last = + new Resolved( + null, + new DefaultCredentials( + accessKeyId, accessKeySecret, conf.get(SECURITY_TOKEN))); + } + } + + @Override + public void setCredentials(Credentials credentials) { + this.last = new Resolved(null, credentials); + } + + @Override + public Credentials getCredentials() { + Resolved resolved = last; + Supplier> supplier = + supplierId == null ? null : CredentialsSupplierRegistry.get(supplierId); + if (supplier != null) { + try { + return resolve(supplier.get(), resolved); + } catch (RuntimeException e) { + if (resolved == null) { + throw e; + } + LOG.warn("Failed to refresh OSS credentials, reusing the last ones.", e); + } + } + if (resolved == null) { + throw new IllegalStateException( + "No OSS credentials available for credentials supplier " + supplierId); + } + return resolved.credentials; + } + + private Credentials resolve(Map options, Resolved resolved) { + if (resolved != null && resolved.options == options) { + return resolved.credentials; + } + String accessKeyId = null; + String accessKeySecret = null; + String securityToken = null; + for (Map.Entry entry : options.entrySet()) { + String key = entry.getKey().toLowerCase(Locale.ROOT); + if (key.equals(ACCESS_KEY_ID.toLowerCase(Locale.ROOT))) { + accessKeyId = entry.getValue(); + } else if (key.equals(ACCESS_KEY_SECRET.toLowerCase(Locale.ROOT))) { + accessKeySecret = entry.getValue(); + } else if (key.equals(SECURITY_TOKEN.toLowerCase(Locale.ROOT))) { + securityToken = entry.getValue(); + } + } + if (accessKeyId == null || accessKeySecret == null) { + throw new IllegalStateException( + "Credentials supplier " + supplierId + " returned no OSS access key."); + } + Credentials credentials = + new DefaultCredentials(accessKeyId, accessKeySecret, securityToken); + last = new Resolved(options, credentials); + return credentials; + } + + private static final class Resolved { + + private final Map options; + private final Credentials credentials; + + private Resolved(Map options, Credentials credentials) { + this.options = options; + this.credentials = credentials; + } + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java new file mode 100644 index 000000000000..7a38c860c033 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java @@ -0,0 +1,185 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.fs.CredentialsSupplierRegistry; +import org.apache.paimon.fs.Path; +import org.apache.paimon.options.Options; + +import com.aliyun.oss.ClientConfiguration; +import com.aliyun.oss.ClientException; +import com.aliyun.oss.OSSClient; +import com.aliyun.oss.common.comm.ExecutionContext; +import com.aliyun.oss.common.comm.RequestMessage; +import com.aliyun.oss.common.comm.ResponseMessage; +import com.aliyun.oss.common.comm.RetryStrategy; +import com.aliyun.oss.common.comm.ServiceClient; +import com.aliyun.oss.internal.OSSObjectOperation; +import com.aliyun.oss.model.GenericRequest; +import org.apache.hadoop.conf.Configuration; +import org.junit.jupiter.api.Test; + +import java.net.URI; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests for {@link RegisteredCredentialsProvider}. */ +public class RegisteredCredentialsProviderTest { + + @Test + public void testEveryRequestIsSignedWithTheCurrentCredentials() throws Exception { + AtomicReference accessKeyId = new AtomicReference<>("ak-1"); + String id = CredentialsSupplierRegistry.register(() -> credentials(accessKeyId.get())); + try { + List authorizations = new ArrayList<>(); + OSSObjectOperation operation = + new OSSObjectOperation(capturingClient(authorizations), provider(id, null)); + operation.setEndpoint(new URI("http://oss-cn-hangzhou.aliyuncs.com")); + + GenericRequest request = new GenericRequest("bucket", "key"); + assertThatThrownBy(() -> operation.getObjectMetadata(request)) + .isInstanceOf(ClientException.class); + accessKeyId.set("ak-2"); + assertThatThrownBy(() -> operation.getObjectMetadata(request)) + .isInstanceOf(ClientException.class); + + assertThat(authorizations).hasSize(2); + assertThat(authorizations.get(0)).contains("ak-1").doesNotContain("ak-2"); + assertThat(authorizations.get(1)).contains("ak-2").doesNotContain("ak-1"); + } finally { + CredentialsSupplierRegistry.unregister(id); + } + } + + @Test + public void testKeepsTheLastCredentialsOnceTheSupplierIsGone() { + String id = CredentialsSupplierRegistry.register(() -> credentials("ak-2")); + RegisteredCredentialsProvider provider = provider(id, "ak-1"); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + + CredentialsSupplierRegistry.unregister(id); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + assertThat(provider(id, "ak-1").getCredentials().getAccessKeyId()).isEqualTo("ak-1"); + } + + @Test + public void testKeepsTheLastCredentialsWhenTheSupplierFails() { + AtomicBoolean failing = new AtomicBoolean(false); + String id = + CredentialsSupplierRegistry.register( + () -> { + if (failing.get()) { + throw new IllegalStateException("REST server unavailable"); + } + return credentials("ak-2"); + }); + try { + RegisteredCredentialsProvider provider = provider(id, null); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + + failing.set(true); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + assertThatThrownBy(() -> provider(id, null).getCredentials()) + .hasMessageContaining("REST server unavailable"); + } finally { + CredentialsSupplierRegistry.unregister(id); + } + } + + @Test + public void testOSSFileIOResolvesCredentialsFromTheRegisteredSupplier() throws Exception { + String id = CredentialsSupplierRegistry.register(() -> credentials("ak-2")); + OSSFileIO fileIO = new OSSFileIO(); + try { + Options options = new Options(); + options.set("fs.oss.endpoint", "http://oss-cn-hangzhou.aliyuncs.com"); + options.set("fs.oss.accessKeyId", "ak-1"); + options.set("fs.oss.accessKeySecret", "sk-1"); + options.set(CredentialsSupplierRegistry.SUPPLIER_ID, id); + fileIO.configure(CatalogContext.create(options)); + + OSSClient client = fileIO.ossClient(new Path("oss://bucket/key")); + assertThat(client.getCredentialsProvider()) + .isInstanceOf(RegisteredCredentialsProvider.class); + assertThat(client.getCredentialsProvider().getCredentials().getAccessKeyId()) + .isEqualTo("ak-2"); + // readers that build their own clients from these options keep the static keys + assertThat(fileIO.hadoopOptions().keySet()) + .doesNotContain("fs.oss.credentials.provider"); + } finally { + fileIO.close(); + CredentialsSupplierRegistry.unregister(id); + } + } + + private static Map credentials(String accessKeyId) { + Map credentials = new HashMap<>(); + credentials.put("fs.oss.accessKeyId", accessKeyId); + credentials.put("fs.oss.accessKeySecret", "secret-" + accessKeyId); + credentials.put("fs.oss.securityToken", "token-" + accessKeyId); + return credentials; + } + + private static RegisteredCredentialsProvider provider(String id, String accessKeyId) { + Configuration conf = new Configuration(false); + conf.set(CredentialsSupplierRegistry.SUPPLIER_ID, id); + if (accessKeyId != null) { + conf.set("fs.oss.accessKeyId", accessKeyId); + conf.set("fs.oss.accessKeySecret", "secret-" + accessKeyId); + } + return new RegisteredCredentialsProvider(null, conf); + } + + /** Records the Authorization header of each request instead of sending it. */ + private static ServiceClient capturingClient(List authorizations) { + return new ServiceClient(new ClientConfiguration()) { + @Override + protected ResponseMessage sendRequestCore( + ServiceClient.Request request, ExecutionContext context) { + authorizations.add(request.getHeaders().get("Authorization")); + throw new ClientException("captured"); + } + + @Override + protected RetryStrategy getDefaultRetryStrategy() { + return new RetryStrategy() { + @Override + public boolean shouldRetry( + Exception e, + RequestMessage request, + ResponseMessage response, + int retries) { + return false; + } + }; + } + + @Override + public void shutdown() {} + }; + } +} diff --git a/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java b/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java index bef7c2a5218a..c78fcbc7cfe8 100644 --- a/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java +++ b/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java @@ -58,8 +58,8 @@ public String getScheme() { public List requiredOptions() { List options = new ArrayList<>(); options.add(new String[] {"fs.oss.endpoint"}); - options.add(new String[] {"fs.oss.accessKeyId"}); - options.add(new String[] {"fs.oss.accessKeySecret"}); + options.add(new String[] {"fs.oss.accessKeyId", "fs.oss.credentials.provider"}); + options.add(new String[] {"fs.oss.accessKeySecret", "fs.oss.credentials.provider"}); return options; } diff --git a/paimon-filesystems/paimon-oss/src/test/java/org/apache/paimon/oss/OSSLoaderTest.java b/paimon-filesystems/paimon-oss/src/test/java/org/apache/paimon/oss/OSSLoaderTest.java new file mode 100644 index 000000000000..650db27e5eea --- /dev/null +++ b/paimon-filesystems/paimon-oss/src/test/java/org/apache/paimon/oss/OSSLoaderTest.java @@ -0,0 +1,46 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.Path; +import org.apache.paimon.options.Options; + +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link OSSLoader}. */ +public class OSSLoaderTest { + + /** A credentials provider replaces the access keys, so the OSS plugin is still chosen. */ + @Test + public void testCredentialsProviderSatisfiesRequiredOptions() throws Exception { + Options options = new Options(); + options.set("fs.oss.endpoint", "oss-cn-hangzhou.aliyuncs.com"); + options.set( + "fs.oss.credentials.provider", + "com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider"); + + FileIO fileIO = FileIO.get(new Path("oss://bucket/key"), CatalogContext.create(options)); + + assertThat(fileIO.getClass().getEnclosingClass()).isEqualTo(OSSLoader.class); + } +} From b4ef06b172b20a994a8f74e07f30e8c33080a30f Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Fri, 25 Sep 2026 20:57:16 +0800 Subject: [PATCH 2/6] [oss] Keep the credentials supplier alive while an open stream still uses it Evicting a delegate from RESTTokenFileIO's cache unregistered its supplier even though the OSS client and its open streams were still live, so they fell back to their last credentials and failed at the next expiry. The registry now holds suppliers weakly. OSSFileIO and RegisteredCredentialsProvider keep a strong reference once they resolve one, so a supplier lives exactly as long as a client or stream that can still sign with it. --- .../fs/CredentialsSupplierRegistry.java | 46 ++++- .../apache/paimon/rest/RESTTokenFileIO.java | 33 ++-- .../fs/CredentialsSupplierRegistryTest.java | 52 +++++ .../java/org/apache/paimon/oss/OSSFileIO.java | 8 +- .../oss/RegisteredCredentialsProvider.java | 7 +- .../RegisteredCredentialsProviderTest.java | 86 ++++---- .../paimon/rest/RESTTokenFileIOOnOSSTest.java | 187 ++++++++++++++++++ 7 files changed, 348 insertions(+), 71 deletions(-) create mode 100644 paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java create mode 100644 paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java diff --git a/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java b/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java index cd1dcb336a78..6643344324f8 100644 --- a/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java +++ b/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java @@ -20,6 +20,9 @@ import javax.annotation.Nullable; +import java.lang.ref.Reference; +import java.lang.ref.ReferenceQueue; +import java.lang.ref.WeakReference; import java.util.Map; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -28,30 +31,61 @@ /** * Credential suppliers that a plugin {@link FileIO} looks up by id, so it can resolve fresh * credentials on every request instead of the ones it was created with. + * + *

Suppliers are held weakly: whoever uses one keeps it alive, and it goes away with its users. */ public final class CredentialsSupplierRegistry { /** Option that passes the id of a registered supplier to the {@link FileIO} being created. */ public static final String SUPPLIER_ID = "fs.credentials-supplier.id"; - private static final Map>> SUPPLIERS = - new ConcurrentHashMap<>(); + private static final Map SUPPLIERS = new ConcurrentHashMap<>(); + + private static final ReferenceQueue>> COLLECTED = + new ReferenceQueue<>(); private CredentialsSupplierRegistry() {} /** Registers a supplier of credential options and returns its id. */ public static String register(Supplier> supplier) { + expungeCollected(); String id = UUID.randomUUID().toString(); - SUPPLIERS.put(id, supplier); + SUPPLIERS.put(id, new SupplierReference(id, supplier, COLLECTED)); return id; } + /** Returns the supplier, or null if it was never registered or nothing uses it anymore. */ @Nullable public static Supplier> get(String id) { - return SUPPLIERS.get(id); + expungeCollected(); + SupplierReference reference = SUPPLIERS.get(id); + return reference == null ? null : reference.get(); + } + + static int size() { + expungeCollected(); + return SUPPLIERS.size(); } - public static void unregister(String id) { - SUPPLIERS.remove(id); + private static void expungeCollected() { + Reference>> reference; + while ((reference = COLLECTED.poll()) != null) { + SupplierReference collected = (SupplierReference) reference; + SUPPLIERS.remove(collected.id, collected); + } + } + + private static final class SupplierReference + extends WeakReference>> { + + private final String id; + + private SupplierReference( + String id, + Supplier> supplier, + ReferenceQueue>> queue) { + super(supplier, queue); + this.id = id; + } } } diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java index 0e6ebe852ff6..eba1edf36815 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java @@ -18,6 +18,7 @@ package org.apache.paimon.rest; +import org.apache.paimon.annotation.VisibleForTesting; import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.data.BlobDescriptor; @@ -50,6 +51,7 @@ import java.util.Map; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import static org.apache.paimon.options.CatalogOptions.FILE_IO_ALLOW_CACHE; import static org.apache.paimon.rest.RESTApi.TOKEN_EXPIRATION_SAFE_TIME_MILLIS; @@ -93,6 +95,12 @@ public static void setFileIOCacheMaximumSize(long maximumSize) { .setMaximum(maximumSize); } + @VisibleForTesting + static void invalidateFileIOCache() { + FILE_IO_CACHE.invalidateAll(); + FILE_IO_CACHE.cleanUp(); + } + static long fileIOCacheMaximumSize() { return FILE_IO_CACHE .policy() @@ -258,7 +266,8 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc } // Lets a FileIO that supports it, such as OSS, sign each request with a fresh token. - String supplierId = CredentialsSupplierRegistry.register(() -> validToken().token()); + Supplier> supplier = () -> validToken().token(); + String supplierId = CredentialsSupplierRegistry.register(supplier); Options options = catalogContext.options(); options = new Options(RESTUtil.merge(options.toMap(), currentToken.token())); options.set(FILE_IO_ALLOW_CACHE, false); @@ -269,14 +278,8 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc catalogContext.hadoopConf(), catalogContext.preferIO(), catalogContext.fallbackIO()); - FileIO fileIO; - try { - fileIO = FileIO.get(path, context); - } catch (IOException | RuntimeException e) { - CredentialsSupplierRegistry.unregister(supplierId); - throw e; - } - FILE_IO_CACHE.put(currentToken, new CachedFileIO(fileIO, supplierId)); + FileIO fileIO = FileIO.get(path, context); + FILE_IO_CACHE.put(currentToken, new CachedFileIO(fileIO, supplier)); return new FileIOWithToken(fileIO, currentToken); } } @@ -310,20 +313,22 @@ private boolean shouldRefresh(long minimumValidityMillis) { < Math.max(TOKEN_EXPIRATION_SAFE_TIME_MILLIS, minimumValidityMillis); } - /** A delegate {@link FileIO} and the credentials supplier registered for it. */ + /** + * A delegate {@link FileIO} and its credentials supplier, kept alive here until the delegate + * picks it up, since the registry only holds it weakly. + */ private static class CachedFileIO implements AutoCloseable { private final FileIO fileIO; - private final String supplierId; + private final Supplier> credentialsSupplier; - private CachedFileIO(FileIO fileIO, String supplierId) { + private CachedFileIO(FileIO fileIO, Supplier> credentialsSupplier) { this.fileIO = fileIO; - this.supplierId = supplierId; + this.credentialsSupplier = credentialsSupplier; } @Override public void close() throws Exception { - CredentialsSupplierRegistry.unregister(supplierId); fileIO.close(); } } diff --git a/paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java b/paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java new file mode 100644 index 000000000000..c7f67ccf3fd0 --- /dev/null +++ b/paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java @@ -0,0 +1,52 @@ +/* + * 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.paimon.fs; + +import org.junit.jupiter.api.Test; + +import java.util.Collections; +import java.util.Map; +import java.util.function.Supplier; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link CredentialsSupplierRegistry}. */ +class CredentialsSupplierRegistryTest { + + @Test + void testSupplierStaysWhileReferencedAndGoesAwayAfterwards() throws InterruptedException { + Supplier> held = () -> Collections.singletonMap("k", "v"); + String heldId = CredentialsSupplierRegistry.register(held); + String droppedId = registerUnreferenced(); + + for (int i = 0; i < 100 && CredentialsSupplierRegistry.get(droppedId) != null; i++) { + System.gc(); + Thread.sleep(10); + } + + assertThat(CredentialsSupplierRegistry.get(droppedId)).isNull(); + assertThat(CredentialsSupplierRegistry.get(heldId)).isSameAs(held); + } + + private static String registerUnreferenced() { + // a capturing lambda, since a non-capturing one is a cached singleton that is never freed + String value = String.valueOf(System.nanoTime()); + return CredentialsSupplierRegistry.register(() -> Collections.singletonMap("k", value)); + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java index 9f5be43fd3a8..9a7ec35fb16d 100644 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java @@ -139,6 +139,9 @@ public class OSSFileIO extends HadoopCompliantFileIO implements HadoopOptionsPro private boolean allowCache = true; @Nullable private String credentialsSupplierId; + // Keeps the supplier alive until the file systems created here have picked it up. + @Nullable private transient Supplier> credentialsSupplier; + @Override public boolean isObjectStore() { return true; @@ -146,7 +149,10 @@ public boolean isObjectStore() { @Override public void configure(CatalogContext context) { - credentialsSupplierId = context.options().get(CredentialsSupplierRegistry.SUPPLIER_ID); + String supplierId = context.options().get(CredentialsSupplierRegistry.SUPPLIER_ID); + credentialsSupplier = + supplierId == null ? null : CredentialsSupplierRegistry.get(supplierId); + credentialsSupplierId = credentialsSupplier == null ? null : supplierId; // The file system is bound to the supplier, so it must not be shared through the cache. allowCache = context.options().get(FILE_IO_ALLOW_CACHE) && credentialsSupplierId == null; hadoopOptions = new Options(); diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java index f34f2498d7a3..0db1010b525c 100644 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java @@ -27,6 +27,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import javax.annotation.Nullable; + import java.net.URI; import java.util.Locale; import java.util.Map; @@ -35,6 +37,7 @@ /** * An OSS {@link CredentialsProvider} that asks a supplier from {@link CredentialsSupplierRegistry} * for credentials on every request, so streams that are already open also pick up refreshed ones. + * It holds the supplier for as long as the OSS client that uses it is alive. */ public class RegisteredCredentialsProvider implements CredentialsProvider { @@ -45,12 +48,14 @@ public class RegisteredCredentialsProvider implements CredentialsProvider { static final String SECURITY_TOKEN = "fs.oss.securityToken"; private final String supplierId; + @Nullable private final Supplier> supplier; // The last options from the supplier and the credentials built from them, swapped together. private volatile Resolved last; public RegisteredCredentialsProvider(URI uri, Configuration conf) { this.supplierId = conf.get(CredentialsSupplierRegistry.SUPPLIER_ID); + this.supplier = supplierId == null ? null : CredentialsSupplierRegistry.get(supplierId); String accessKeyId = conf.get(ACCESS_KEY_ID); String accessKeySecret = conf.get(ACCESS_KEY_SECRET); if (accessKeyId != null && accessKeySecret != null) { @@ -70,8 +75,6 @@ public void setCredentials(Credentials credentials) { @Override public Credentials getCredentials() { Resolved resolved = last; - Supplier> supplier = - supplierId == null ? null : CredentialsSupplierRegistry.get(supplierId); if (supplier != null) { try { return resolve(supplier.get(), resolved); diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java index 7a38c860c033..96d9a7608f67 100644 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java @@ -43,6 +43,7 @@ import java.util.Map; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Supplier; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -53,66 +54,56 @@ public class RegisteredCredentialsProviderTest { @Test public void testEveryRequestIsSignedWithTheCurrentCredentials() throws Exception { AtomicReference accessKeyId = new AtomicReference<>("ak-1"); - String id = CredentialsSupplierRegistry.register(() -> credentials(accessKeyId.get())); - try { - List authorizations = new ArrayList<>(); - OSSObjectOperation operation = - new OSSObjectOperation(capturingClient(authorizations), provider(id, null)); - operation.setEndpoint(new URI("http://oss-cn-hangzhou.aliyuncs.com")); - - GenericRequest request = new GenericRequest("bucket", "key"); - assertThatThrownBy(() -> operation.getObjectMetadata(request)) - .isInstanceOf(ClientException.class); - accessKeyId.set("ak-2"); - assertThatThrownBy(() -> operation.getObjectMetadata(request)) - .isInstanceOf(ClientException.class); - - assertThat(authorizations).hasSize(2); - assertThat(authorizations.get(0)).contains("ak-1").doesNotContain("ak-2"); - assertThat(authorizations.get(1)).contains("ak-2").doesNotContain("ak-1"); - } finally { - CredentialsSupplierRegistry.unregister(id); - } + Supplier> supplier = () -> credentials(accessKeyId.get()); + String id = CredentialsSupplierRegistry.register(supplier); + List authorizations = new ArrayList<>(); + OSSObjectOperation operation = + new OSSObjectOperation(capturingClient(authorizations), provider(id, null)); + operation.setEndpoint(new URI("http://oss-cn-hangzhou.aliyuncs.com")); + + GenericRequest request = new GenericRequest("bucket", "key"); + assertThatThrownBy(() -> operation.getObjectMetadata(request)) + .isInstanceOf(ClientException.class); + accessKeyId.set("ak-2"); + assertThatThrownBy(() -> operation.getObjectMetadata(request)) + .isInstanceOf(ClientException.class); + + assertThat(authorizations).hasSize(2); + assertThat(authorizations.get(0)).contains("ak-1").doesNotContain("ak-2"); + assertThat(authorizations.get(1)).contains("ak-2").doesNotContain("ak-1"); } @Test - public void testKeepsTheLastCredentialsOnceTheSupplierIsGone() { - String id = CredentialsSupplierRegistry.register(() -> credentials("ak-2")); - RegisteredCredentialsProvider provider = provider(id, "ak-1"); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); - - CredentialsSupplierRegistry.unregister(id); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); - assertThat(provider(id, "ak-1").getCredentials().getAccessKeyId()).isEqualTo("ak-1"); + public void testUsesConfiguredCredentialsWithoutARegisteredSupplier() { + RegisteredCredentialsProvider provider = provider("unknown-id", "ak-1"); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-1"); } @Test public void testKeepsTheLastCredentialsWhenTheSupplierFails() { AtomicBoolean failing = new AtomicBoolean(false); - String id = - CredentialsSupplierRegistry.register( - () -> { - if (failing.get()) { - throw new IllegalStateException("REST server unavailable"); - } - return credentials("ak-2"); - }); - try { - RegisteredCredentialsProvider provider = provider(id, null); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + Supplier> supplier = + () -> { + if (failing.get()) { + throw new IllegalStateException("REST server unavailable"); + } + return credentials("ak-2"); + }; + String id = CredentialsSupplierRegistry.register(supplier); + RegisteredCredentialsProvider provider = provider(id, null); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); - failing.set(true); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); - assertThatThrownBy(() -> provider(id, null).getCredentials()) - .hasMessageContaining("REST server unavailable"); - } finally { - CredentialsSupplierRegistry.unregister(id); - } + failing.set(true); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + assertThatThrownBy(() -> provider(id, null).getCredentials()) + .hasMessageContaining("REST server unavailable"); + assertThat(supplier).isNotNull(); } @Test public void testOSSFileIOResolvesCredentialsFromTheRegisteredSupplier() throws Exception { - String id = CredentialsSupplierRegistry.register(() -> credentials("ak-2")); + Supplier> supplier = () -> credentials("ak-2"); + String id = CredentialsSupplierRegistry.register(supplier); OSSFileIO fileIO = new OSSFileIO(); try { Options options = new Options(); @@ -132,7 +123,6 @@ public void testOSSFileIOResolvesCredentialsFromTheRegisteredSupplier() throws E .doesNotContain("fs.oss.credentials.provider"); } finally { fileIO.close(); - CredentialsSupplierRegistry.unregister(id); } } diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java new file mode 100644 index 000000000000..295b4207f223 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java @@ -0,0 +1,187 @@ +/* + * 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.paimon.rest; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.fs.FileIO; +import org.apache.paimon.fs.FileIOLoader; +import org.apache.paimon.fs.Path; +import org.apache.paimon.fs.PluginFileIO; +import org.apache.paimon.fs.SeekableInputStream; +import org.apache.paimon.options.Options; +import org.apache.paimon.oss.OSSFileIO; +import org.apache.paimon.rest.responses.GetTableTokenResponse; + +import com.sun.net.httpserver.HttpServer; +import org.apache.hadoop.conf.Configuration; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicLong; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** Tests {@link RESTTokenFileIO} over a real {@link OSSFileIO} against a local OSS endpoint. */ +public class RESTTokenFileIOOnOSSTest { + + private static final byte[] CONTENT = new byte[4096]; + + private final List getAuthorizations = Collections.synchronizedList(new ArrayList<>()); + private HttpServer server; + + @BeforeEach + public void startServer() throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext( + "/", + exchange -> { + exchange.getResponseHeaders() + .add("Last-Modified", "Thu, 24 Sep 2026 00:00:00 GMT"); + if ("HEAD".equals(exchange.getRequestMethod())) { + exchange.getResponseHeaders() + .add("Content-Length", String.valueOf(CONTENT.length)); + exchange.sendResponseHeaders(200, -1); + } else { + getAuthorizations.add( + exchange.getRequestHeaders().getFirst("Authorization")); + String[] range = + exchange.getRequestHeaders() + .getFirst("Range") + .substring("bytes=".length()) + .split("-"); + int start = Integer.parseInt(range[0]); + int end = Math.min(Integer.parseInt(range[1]), CONTENT.length - 1); + exchange.getResponseHeaders() + .add( + "Content-Range", + "bytes " + start + "-" + end + "/" + CONTENT.length); + exchange.sendResponseHeaders(206, end - start + 1); + try (OutputStream out = exchange.getResponseBody()) { + out.write(CONTENT, start, end - start + 1); + } + } + exchange.close(); + }); + server.start(); + } + + @AfterEach + public void stopServer() { + server.stop(0); + } + + /** An open stream keeps refreshing even after its delegate FileIO left the cache. */ + @Test + public void testOpenStreamPicksUpRefreshedTokenAfterCacheEviction() throws Exception { + Identifier identifier = Identifier.create("db", "table"); + AtomicLong now = new AtomicLong(1700000000000L); + RESTApi api = mock(RESTApi.class); + when(api.loadTableToken(identifier)) + .thenReturn( + new GetTableTokenResponse( + token("ak-1"), now.get() + Duration.ofHours(2).toMillis()), + new GetTableTokenResponse( + token("ak-2"), now.get() + Duration.ofHours(4).toMillis())); + Options options = new Options(); + options.set("fs.oss.multipart.download.size", "1024"); + options.set("fs.oss.multipart.download.threads", "1"); + RESTTokenFileIO fileIO = + new RESTTokenFileIO( + CatalogContext.create( + options, new Configuration(false), new PluginOSSLoader(), null), + api, + identifier, + new Path("oss://bucket/table")) { + @Override + long currentTimeMillis() { + return now.get(); + } + }; + + try (SeekableInputStream in = fileIO.newInputStream(new Path("oss://bucket/table/f"))) { + in.read(new byte[1024]); + + RESTTokenFileIO.invalidateFileIOCache(); + System.gc(); + // 30 minutes left is inside the refresh window + now.addAndGet(Duration.ofMinutes(90).toMillis()); + in.seek(3000); + in.read(new byte[1024]); + } + + assertThat(getAuthorizations.get(0)).contains("ak-1"); + assertThat(getAuthorizations.get(getAuthorizations.size() - 1)).contains("ak-2"); + verify(api, times(2)).loadTableToken(identifier); + } + + private Map token(String accessKeyId) { + Map token = new HashMap<>(); + token.put("fs.oss.endpoint", "http://127.0.0.1:" + server.getAddress().getPort()); + token.put("fs.oss.accessKeyId", accessKeyId); + token.put("fs.oss.accessKeySecret", "secret-" + accessKeyId); + token.put("fs.oss.securityToken", "token-" + accessKeyId); + return token; + } + + /** Wraps {@link OSSFileIO} the way the OSS plugin does, including its no-op close. */ + private static class PluginOSSLoader implements FileIOLoader { + + @Override + public String getScheme() { + return "oss"; + } + + @Override + public FileIO load(Path path) { + return new PluginFileIO() { + @Override + public boolean isObjectStore() { + return true; + } + + @Override + protected FileIO createFileIO(Path path) { + FileIO fileIO = new OSSFileIO(); + fileIO.configure(CatalogContext.create(options)); + return fileIO; + } + + @Override + protected ClassLoader pluginClassLoader() { + return OSSFileIO.class.getClassLoader(); + } + }; + } + } +} From 21a64e9f27b8a93f3d2184617b2e0389f45d896c Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Mon, 28 Sep 2026 12:30:14 +0800 Subject: [PATCH 3/6] [oss] Refresh the REST data token inside the OSS credentials provider Passing a live supplier from RESTTokenFileIO to the OSS client through a static registry tied the supplier to objects the client does not own, and got the lifetime wrong on cache eviction. Like Iceberg's VendedCredentialsProvider and Hadoop's credential providers, the OSS provider is now built from configuration alone: RESTTokenFileIO names the table and the token expiry in the delegate options, OSSFileIO hands the catalog options to RESTTokenCredentialsProvider, and the provider loads the table's token through RESTTokenRefresher when it is about to expire. The registry is removed. --- docs/docs/maintenance/filesystems.mdx | 2 +- .../fs/CredentialsSupplierRegistry.java | 91 --------- .../apache/paimon/rest/RESTTokenFileIO.java | 68 +++---- .../paimon/rest/RESTTokenRefresher.java | 145 +++++++++++++++ .../fs/CredentialsSupplierRegistryTest.java | 52 ------ .../paimon/rest/RESTTokenFileIOTest.java | 47 ++--- .../paimon/rest/RESTTokenRefresherTest.java | 127 +++++++++++++ .../java/org/apache/paimon/oss/OSSFileIO.java | 31 ++-- .../oss/RESTTokenCredentialsProvider.java | 112 +++++++++++ .../oss/RegisteredCredentialsProvider.java | 132 ------------- .../paimon/oss/FakeRESTAndOSSServer.java | 133 +++++++++++++ .../oss/RESTTokenCredentialsProviderTest.java | 124 +++++++++++++ .../RegisteredCredentialsProviderTest.java | 175 ------------------ .../paimon/rest/RESTTokenFileIOOnOSSTest.java | 146 ++++----------- 14 files changed, 735 insertions(+), 650 deletions(-) delete mode 100644 paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java create mode 100644 paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java delete mode 100644 paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java create mode 100644 paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java create mode 100644 paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RESTTokenCredentialsProvider.java delete mode 100644 paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java create mode 100644 paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java create mode 100644 paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java delete mode 100644 paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java diff --git a/docs/docs/maintenance/filesystems.mdx b/docs/docs/maintenance/filesystems.mdx index f708a73df667..f69aaf3ebbbf 100644 --- a/docs/docs/maintenance/filesystems.mdx +++ b/docs/docs/maintenance/filesystems.mdx @@ -337,7 +337,7 @@ Download [paimon-jindo-@@VERSION@@.jar](https://repository.apache.org/snapshots/ Instead of `fs.oss.accessKeyId` and `fs.oss.accessKeySecret`, you can set `fs.oss.credentials.provider` to a class implementing `com.aliyun.oss.common.auth.CredentialsProvider`, for example `com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider`. The provider is asked for credentials on every request. -With a REST catalog that vends data tokens, the OSS FileIO signs every request with the catalog's current token, so streams that are already open keep working after the token is refreshed. +With a REST catalog that vends data tokens, the OSS FileIO refreshes the table's token from the catalog by itself before it expires, so streams that are already open keep working after the token is refreshed. ### Server-Side Encryption diff --git a/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java b/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java deleted file mode 100644 index 6643344324f8..000000000000 --- a/paimon-common/src/main/java/org/apache/paimon/fs/CredentialsSupplierRegistry.java +++ /dev/null @@ -1,91 +0,0 @@ -/* - * 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.paimon.fs; - -import javax.annotation.Nullable; - -import java.lang.ref.Reference; -import java.lang.ref.ReferenceQueue; -import java.lang.ref.WeakReference; -import java.util.Map; -import java.util.UUID; -import java.util.concurrent.ConcurrentHashMap; -import java.util.function.Supplier; - -/** - * Credential suppliers that a plugin {@link FileIO} looks up by id, so it can resolve fresh - * credentials on every request instead of the ones it was created with. - * - *

Suppliers are held weakly: whoever uses one keeps it alive, and it goes away with its users. - */ -public final class CredentialsSupplierRegistry { - - /** Option that passes the id of a registered supplier to the {@link FileIO} being created. */ - public static final String SUPPLIER_ID = "fs.credentials-supplier.id"; - - private static final Map SUPPLIERS = new ConcurrentHashMap<>(); - - private static final ReferenceQueue>> COLLECTED = - new ReferenceQueue<>(); - - private CredentialsSupplierRegistry() {} - - /** Registers a supplier of credential options and returns its id. */ - public static String register(Supplier> supplier) { - expungeCollected(); - String id = UUID.randomUUID().toString(); - SUPPLIERS.put(id, new SupplierReference(id, supplier, COLLECTED)); - return id; - } - - /** Returns the supplier, or null if it was never registered or nothing uses it anymore. */ - @Nullable - public static Supplier> get(String id) { - expungeCollected(); - SupplierReference reference = SUPPLIERS.get(id); - return reference == null ? null : reference.get(); - } - - static int size() { - expungeCollected(); - return SUPPLIERS.size(); - } - - private static void expungeCollected() { - Reference>> reference; - while ((reference = COLLECTED.poll()) != null) { - SupplierReference collected = (SupplierReference) reference; - SUPPLIERS.remove(collected.id, collected); - } - } - - private static final class SupplierReference - extends WeakReference>> { - - private final String id; - - private SupplierReference( - String id, - Supplier> supplier, - ReferenceQueue>> queue) { - super(supplier, queue); - this.id = id; - } - } -} diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java index eba1edf36815..25a2dd9924be 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java @@ -22,7 +22,6 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.data.BlobDescriptor; -import org.apache.paimon.fs.CredentialsSupplierRegistry; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.FileStatus; import org.apache.paimon.fs.Path; @@ -51,7 +50,6 @@ import java.util.Map; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; -import java.util.function.Supplier; import static org.apache.paimon.options.CatalogOptions.FILE_IO_ALLOW_CACHE; import static org.apache.paimon.rest.RESTApi.TOKEN_EXPIRATION_SAFE_TIME_MILLIS; @@ -70,12 +68,12 @@ public class RESTTokenFileIO implements FileIO { .defaultValue(false) .withDescription("Whether to support data token provided by the REST server."); - private static final Cache FILE_IO_CACHE = + private static final Cache FILE_IO_CACHE = Caffeine.newBuilder() .maximumSize(1000) .expireAfterAccess(10, TimeUnit.HOURS) .removalListener( - (ignored, value, cause) -> IOUtils.closeQuietly((CachedFileIO) value)) + (ignored, value, cause) -> IOUtils.closeQuietly((FileIO) value)) .scheduler( Scheduler.forScheduledExecutorService( Executors.newSingleThreadScheduledExecutor( @@ -254,32 +252,30 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc + "REST credential lifetime after refresh."); } - CachedFileIO cached = FILE_IO_CACHE.getIfPresent(currentToken); - if (cached != null) { - return new FileIOWithToken(cached.fileIO, currentToken); + FileIO fileIO = FILE_IO_CACHE.getIfPresent(currentToken); + if (fileIO != null) { + return new FileIOWithToken(fileIO, currentToken); } synchronized (FILE_IO_CACHE) { - cached = FILE_IO_CACHE.getIfPresent(currentToken); - if (cached != null) { - return new FileIOWithToken(cached.fileIO, currentToken); + fileIO = FILE_IO_CACHE.getIfPresent(currentToken); + if (fileIO != null) { + return new FileIOWithToken(fileIO, currentToken); } - // Lets a FileIO that supports it, such as OSS, sign each request with a fresh token. - Supplier> supplier = () -> validToken().token(); - String supplierId = CredentialsSupplierRegistry.register(supplier); Options options = catalogContext.options(); options = new Options(RESTUtil.merge(options.toMap(), currentToken.token())); options.set(FILE_IO_ALLOW_CACHE, false); - options.set(CredentialsSupplierRegistry.SUPPLIER_ID, supplierId); + // Lets a FileIO that supports it, such as OSS, refresh this token by itself. + RESTTokenRefresher.configure(options, tokenIdentifier(), currentToken.expireAtMillis()); CatalogContext context = CatalogContext.create( options, catalogContext.hadoopConf(), catalogContext.preferIO(), catalogContext.fallbackIO()); - FileIO fileIO = FileIO.get(path, context); - FILE_IO_CACHE.put(currentToken, new CachedFileIO(fileIO, supplier)); + fileIO = FileIO.get(path, context); + FILE_IO_CACHE.put(currentToken, fileIO); return new FileIOWithToken(fileIO, currentToken); } } @@ -313,26 +309,6 @@ private boolean shouldRefresh(long minimumValidityMillis) { < Math.max(TOKEN_EXPIRATION_SAFE_TIME_MILLIS, minimumValidityMillis); } - /** - * A delegate {@link FileIO} and its credentials supplier, kept alive here until the delegate - * picks it up, since the registry only holds it weakly. - */ - private static class CachedFileIO implements AutoCloseable { - - private final FileIO fileIO; - private final Supplier> credentialsSupplier; - - private CachedFileIO(FileIO fileIO, Supplier> credentialsSupplier) { - this.fileIO = fileIO; - this.credentialsSupplier = credentialsSupplier; - } - - @Override - public void close() throws Exception { - fileIO.close(); - } - } - private static class FileIOWithToken { private final FileIO fileIO; @@ -349,15 +325,7 @@ private void refreshToken() { if (apiInstance == null) { apiInstance = new RESTApi(catalogContext.options(), false); } - Identifier tableIdentifier = identifier; - if (identifier.isSystemTable()) { - tableIdentifier = - new Identifier( - identifier.getDatabaseName(), - identifier.getTableName(), - identifier.getBranchName()); - } - GetTableTokenResponse response = apiInstance.loadTableToken(tableIdentifier); + GetTableTokenResponse response = apiInstance.loadTableToken(tokenIdentifier()); LOG.info( "end refresh data token for identifier [{}] expiresAtMillis [{}]", identifier, @@ -369,6 +337,16 @@ private void refreshToken() { response.getExpiresAtMillis()); } + private Identifier tokenIdentifier() { + if (identifier.isSystemTable()) { + return new Identifier( + identifier.getDatabaseName(), + identifier.getTableName(), + identifier.getBranchName()); + } + return identifier; + } + private Map mergeTokenWithCatalogOptions(Map token) { Map newToken = Maps.newLinkedHashMap(token); Options catalogOptions = catalogContext.options(); diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java new file mode 100644 index 000000000000..ae481ff5ed32 --- /dev/null +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java @@ -0,0 +1,145 @@ +/* + * 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.paimon.rest; + +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.options.Options; +import org.apache.paimon.rest.responses.GetTableTokenResponse; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.annotation.Nullable; + +import static org.apache.paimon.rest.RESTApi.TOKEN_EXPIRATION_SAFE_TIME_MILLIS; + +/** + * Loads the data token of a table from the REST server and refreshes it before it expires, so a + * storage credentials provider can refresh by itself from the options it was created with. + */ +public class RESTTokenRefresher { + + private static final Logger LOG = LoggerFactory.getLogger(RESTTokenRefresher.class); + + /** Database of the table whose token is refreshed, set by {@link RESTTokenFileIO}. */ + public static final String DATABASE = "data-token.database"; + + /** Object name of the table whose token is refreshed, including any branch. */ + public static final String OBJECT = "data-token.object"; + + /** Expiration of the token already present in the options. */ + public static final String EXPIRES_AT_MILLIS = "data-token.expires-at-millis"; + + // After a failed refresh, keep using the current token this long before trying again. + static final long RETRY_INTERVAL_MILLIS = 10_000L; + + private final Options catalogOptions; + private final Identifier identifier; + + @Nullable private RESTApi api; + @Nullable private volatile RESTToken token; + private long nextAttemptMillis; + + RESTTokenRefresher( + Options catalogOptions, + Identifier identifier, + @Nullable RESTApi api, + @Nullable RESTToken token) { + this.catalogOptions = catalogOptions; + this.identifier = identifier; + this.api = api; + this.token = token; + } + + /** Names the table and the expiration of the merged token, so a refresher can be created. */ + public static void configure(Options options, Identifier identifier, long expiresAtMillis) { + options.set(DATABASE, identifier.getDatabaseName()); + options.set(OBJECT, identifier.getObjectName()); + options.set(EXPIRES_AT_MILLIS, String.valueOf(expiresAtMillis)); + } + + /** Whether the options name a table, see {@link #configure}. */ + public static boolean isConfigured(Options options) { + return options.containsKey(DATABASE) && options.containsKey(OBJECT); + } + + /** + * Creates a refresher from catalog options that also carry the current token and the keys set + * by {@link #configure}. + */ + public static RESTTokenRefresher fromOptions(Options options) { + Identifier identifier = new Identifier(options.get(DATABASE), options.get(OBJECT)); + String expiresAt = options.get(EXPIRES_AT_MILLIS); + RESTToken token = + expiresAt == null + ? null + : new RESTToken(options.toMap(), Long.parseLong(expiresAt)); + return new RESTTokenRefresher(options, identifier, null, token); + } + + /** Returns the current token, loading a new one when it is about to expire. */ + public RESTToken token() { + RESTToken current = token; + if (current != null && !expiresSoon(current)) { + return current; + } + synchronized (this) { + current = token; + if (current != null && !expiresSoon(current)) { + return current; + } + long now = currentTimeMillis(); + boolean usable = current != null && now < current.expireAtMillis(); + if (usable && now < nextAttemptMillis) { + return current; + } + try { + RESTToken loaded = load(); + token = loaded; + return loaded; + } catch (RuntimeException e) { + if (!usable) { + throw e; + } + nextAttemptMillis = now + RETRY_INTERVAL_MILLIS; + LOG.warn( + "Failed to refresh the data token of {}, keeping the current one.", + identifier, + e); + return current; + } + } + } + + private boolean expiresSoon(RESTToken token) { + return token.expireAtMillis() - currentTimeMillis() < TOKEN_EXPIRATION_SAFE_TIME_MILLIS; + } + + private RESTToken load() { + if (api == null) { + api = new RESTApi(catalogOptions, false); + } + GetTableTokenResponse response = api.loadTableToken(identifier); + return new RESTToken(response.getToken(), response.getExpiresAtMillis()); + } + + long currentTimeMillis() { + return System.currentTimeMillis(); + } +} diff --git a/paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java b/paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java deleted file mode 100644 index c7f67ccf3fd0..000000000000 --- a/paimon-common/src/test/java/org/apache/paimon/fs/CredentialsSupplierRegistryTest.java +++ /dev/null @@ -1,52 +0,0 @@ -/* - * 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.paimon.fs; - -import org.junit.jupiter.api.Test; - -import java.util.Collections; -import java.util.Map; -import java.util.function.Supplier; - -import static org.assertj.core.api.Assertions.assertThat; - -/** Tests for {@link CredentialsSupplierRegistry}. */ -class CredentialsSupplierRegistryTest { - - @Test - void testSupplierStaysWhileReferencedAndGoesAwayAfterwards() throws InterruptedException { - Supplier> held = () -> Collections.singletonMap("k", "v"); - String heldId = CredentialsSupplierRegistry.register(held); - String droppedId = registerUnreferenced(); - - for (int i = 0; i < 100 && CredentialsSupplierRegistry.get(droppedId) != null; i++) { - System.gc(); - Thread.sleep(10); - } - - assertThat(CredentialsSupplierRegistry.get(droppedId)).isNull(); - assertThat(CredentialsSupplierRegistry.get(heldId)).isSameAs(held); - } - - private static String registerUnreferenced() { - // a capturing lambda, since a non-capturing one is a cached singleton that is never freed - String value = String.valueOf(System.nanoTime()); - return CredentialsSupplierRegistry.register(() -> Collections.singletonMap("k", value)); - } -} diff --git a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java index e577a5594883..1f63b10e9009 100644 --- a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java @@ -21,7 +21,6 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.data.BlobDescriptor; -import org.apache.paimon.fs.CredentialsSupplierRegistry; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.FileIOLoader; import org.apache.paimon.fs.FileStatus; @@ -36,10 +35,8 @@ import java.io.IOException; import java.time.Duration; import java.util.Collections; -import java.util.Map; import java.util.UUID; import java.util.concurrent.atomic.AtomicLong; -import java.util.function.Supplier; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; @@ -232,48 +229,38 @@ long currentTimeMillis() { } @Test - void testDelegateCredentialsSupplierFollowsTokenRefresh() throws IOException { + void testDelegateOptionsNameTheTableAndTokenExpiry() throws IOException { Path root = new Path("oss://bucket/table"); - AtomicLong now = new AtomicLong(1700000000000L); FileIO delegate = mock(FileIO.class); when(delegate.exists(any())).thenReturn(true); FileIOLoader loader = mock(FileIOLoader.class); when(loader.load(any())).thenReturn(delegate); when(loader.getScheme()).thenReturn("oss"); RESTApi api = mock(RESTApi.class); - Identifier identifier = Identifier.create("db", "table"); - String first = UUID.randomUUID().toString(); - String second = UUID.randomUUID().toString(); - when(api.loadTableToken(identifier)) + Identifier table = new Identifier("db", "table", "b1"); + long expiresAt = System.currentTimeMillis() + Duration.ofHours(2).toMillis(); + when(api.loadTableToken(table)) .thenReturn( new GetTableTokenResponse( - Collections.singletonMap("test.token", first), - now.get() + Duration.ofHours(2).toMillis()), - new GetTableTokenResponse( - Collections.singletonMap("test.token", second), - now.get() + Duration.ofHours(4).toMillis())); + Collections.singletonMap("token", UUID.randomUUID().toString()), + expiresAt)); RESTTokenFileIO fileIO = new RESTTokenFileIO( - CatalogContext.create(new Options(), loader, null), api, identifier, root) { - @Override - long currentTimeMillis() { - return now.get(); - } - }; + CatalogContext.create(new Options(), loader, null), + api, + new Identifier("db", "table", "b1", "files"), + root); fileIO.exists(root); + ArgumentCaptor context = ArgumentCaptor.forClass(CatalogContext.class); verify(delegate, atLeastOnce()).configure(context.capture()); - String supplierId = - context.getValue().options().get(CredentialsSupplierRegistry.SUPPLIER_ID); - Supplier> supplier = CredentialsSupplierRegistry.get(supplierId); - assertThat(supplier).isNotNull(); - assertThat(supplier.get()).containsEntry("test.token", first); - - // 30 minutes left is inside the safe window, so the delegate is handed a new token - now.addAndGet(Duration.ofMinutes(90).toMillis()); - assertThat(supplier.get()).containsEntry("test.token", second); - verify(api, times(2)).loadTableToken(identifier); + Options options = context.getValue().options(); + // a system table refreshes the token of its table + assertThat(options.get(RESTTokenRefresher.DATABASE)).isEqualTo("db"); + assertThat(options.get(RESTTokenRefresher.OBJECT)).isEqualTo("table$branch_b1"); + assertThat(options.get(RESTTokenRefresher.EXPIRES_AT_MILLIS)) + .isEqualTo(String.valueOf(expiresAt)); } @Test diff --git a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java new file mode 100644 index 000000000000..3fcf1e4b5d81 --- /dev/null +++ b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java @@ -0,0 +1,127 @@ +/* + * 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.paimon.rest; + +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.options.Options; +import org.apache.paimon.rest.responses.GetTableTokenResponse; + +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.Collections; +import java.util.concurrent.atomic.AtomicLong; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** Tests for {@link RESTTokenRefresher}. */ +class RESTTokenRefresherTest { + + private static final long NOW = 1700000000000L; + + private final Identifier identifier = new Identifier("db", "table", "b1"); + private final AtomicLong now = new AtomicLong(NOW); + private final RESTApi api = mock(RESTApi.class); + + @Test + void testKeepsTheTokenUntilItIsAboutToExpire() { + when(api.loadTableToken(identifier)).thenReturn(response("second", hours(4))); + RESTTokenRefresher refresher = refresher(token("first", hours(2))); + + assertThat(refresher.token().token()).containsEntry("k", "first"); + verify(api, never()).loadTableToken(identifier); + + // 30 minutes left is inside the refresh window + now.addAndGet(Duration.ofMinutes(90).toMillis()); + assertThat(refresher.token().token()).containsEntry("k", "second"); + assertThat(refresher.token().token()).containsEntry("k", "second"); + verify(api, times(1)).loadTableToken(identifier); + } + + @Test + void testKeepsTheCurrentTokenWhileRefreshFailsAndRetriesLater() { + when(api.loadTableToken(identifier)) + .thenThrow(new IllegalStateException("REST server unavailable")) + .thenReturn(response("second", hours(4))); + RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(30))); + + assertThat(refresher.token().token()).containsEntry("k", "first"); + assertThat(refresher.token().token()).containsEntry("k", "first"); + verify(api, times(1)).loadTableToken(identifier); + + now.addAndGet(RESTTokenRefresher.RETRY_INTERVAL_MILLIS); + assertThat(refresher.token().token()).containsEntry("k", "second"); + verify(api, times(2)).loadTableToken(identifier); + } + + @Test + void testFailsWhenTheTokenExpiredAndRefreshFails() { + when(api.loadTableToken(identifier)) + .thenThrow(new IllegalStateException("REST server unavailable")); + RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(-1))); + + assertThatThrownBy(refresher::token).hasMessageContaining("REST server unavailable"); + } + + @Test + void testConfigureRoundTripsTheTable() { + Options options = new Options(); + options.set("k", "first"); + Identifier dotted = new Identifier("my.db", "table", "b1"); + long expiresAt = System.currentTimeMillis() + hours(2).toMillis(); + RESTTokenRefresher.configure(options, dotted, expiresAt); + + assertThat(RESTTokenRefresher.isConfigured(options)).isTrue(); + assertThat(RESTTokenRefresher.isConfigured(new Options())).isFalse(); + RESTTokenRefresher refresher = RESTTokenRefresher.fromOptions(options); + // the token already in the options is used without a request + assertThat(refresher.token().token()).containsEntry("k", "first"); + assertThat(refresher.token().expireAtMillis()).isEqualTo(expiresAt); + assertThat(options.get(RESTTokenRefresher.DATABASE)).isEqualTo("my.db"); + assertThat(options.get(RESTTokenRefresher.OBJECT)).isEqualTo("table$branch_b1"); + } + + private RESTTokenRefresher refresher(RESTToken token) { + return new RESTTokenRefresher(new Options(), identifier, api, token) { + @Override + long currentTimeMillis() { + return now.get(); + } + }; + } + + private RESTToken token(String value, Duration lifetime) { + return new RESTToken(Collections.singletonMap("k", value), now.get() + lifetime.toMillis()); + } + + private GetTableTokenResponse response(String value, Duration lifetime) { + return new GetTableTokenResponse( + Collections.singletonMap("k", value), NOW + lifetime.toMillis()); + } + + private static Duration hours(int hours) { + return Duration.ofHours(hours); + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java index 9a7ec35fb16d..21667818b41c 100644 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java @@ -20,12 +20,12 @@ import org.apache.paimon.catalog.CatalogContext; import org.apache.paimon.data.BlobDescriptor; -import org.apache.paimon.fs.CredentialsSupplierRegistry; import org.apache.paimon.fs.FileIO; import org.apache.paimon.fs.HadoopOptionsProvider; import org.apache.paimon.fs.Path; import org.apache.paimon.fs.TwoPhaseOutputStream; import org.apache.paimon.options.Options; +import org.apache.paimon.rest.RESTTokenRefresher; import org.apache.paimon.utils.IOUtils; import org.apache.paimon.utils.ReflectionUtils; import org.apache.paimon.utils.SensitiveConfigUtils; @@ -137,10 +137,9 @@ public class OSSFileIO extends HadoopCompliantFileIO implements HadoopOptionsPro private Options hadoopOptions; private boolean allowCache = true; - @Nullable private String credentialsSupplierId; - // Keeps the supplier alive until the file systems created here have picked it up. - @Nullable private transient Supplier> credentialsSupplier; + // Catalog options for RESTTokenCredentialsProvider when the options name a REST catalog table. + @Nullable private Map restTokenOptions; @Override public boolean isObjectStore() { @@ -149,12 +148,12 @@ public boolean isObjectStore() { @Override public void configure(CatalogContext context) { - String supplierId = context.options().get(CredentialsSupplierRegistry.SUPPLIER_ID); - credentialsSupplier = - supplierId == null ? null : CredentialsSupplierRegistry.get(supplierId); - credentialsSupplierId = credentialsSupplier == null ? null : supplierId; - // The file system is bound to the supplier, so it must not be shared through the cache. - allowCache = context.options().get(FILE_IO_ALLOW_CACHE) && credentialsSupplierId == null; + restTokenOptions = + RESTTokenRefresher.isConfigured(context.options()) + ? context.options().toMap() + : null; + // The file system refreshes a table's token, so it must not be shared through the cache. + allowCache = context.options().get(FILE_IO_ALLOW_CACHE) && restTokenOptions == null; hadoopOptions = new Options(); // read all configuration with prefix 'CONFIG_PREFIXES' for (String key : context.options().keySet()) { @@ -208,12 +207,16 @@ protected AliyunOSSFileSystem createFileSystem(org.apache.hadoop.fs.Path path) { // retrieve props from the file, which comes at a high cost Configuration hadoopConf = new Configuration(SHARED_CONFIG); hadoopOptions.toMap().forEach(hadoopConf::set); - if (credentialsSupplierId != null) { + if (restTokenOptions != null) { hadoopConf.set( OSS_CREDENTIALS_PROVIDER, - RegisteredCredentialsProvider.class.getName()); - hadoopConf.set( - CredentialsSupplierRegistry.SUPPLIER_ID, credentialsSupplierId); + RESTTokenCredentialsProvider.class.getName()); + restTokenOptions.forEach( + (key, value) -> + hadoopConf.set( + RESTTokenCredentialsProvider.CATALOG_OPTIONS_PREFIX + + key, + value)); } URI fsUri = path.toUri(); if (scheme == null && authority == null) { diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RESTTokenCredentialsProvider.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RESTTokenCredentialsProvider.java new file mode 100644 index 000000000000..8531115363b1 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RESTTokenCredentialsProvider.java @@ -0,0 +1,112 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.options.Options; +import org.apache.paimon.rest.RESTToken; +import org.apache.paimon.rest.RESTTokenRefresher; + +import com.aliyun.oss.common.auth.Credentials; +import com.aliyun.oss.common.auth.CredentialsProvider; +import com.aliyun.oss.common.auth.DefaultCredentials; +import org.apache.hadoop.conf.Configuration; + +import java.net.URI; +import java.util.Locale; +import java.util.Map; + +import static org.apache.paimon.utils.Preconditions.checkArgument; + +/** + * An OSS {@link CredentialsProvider} that refreshes the data token of a REST catalog table by + * itself and is asked on every request, so streams that are already open also pick up refreshed + * tokens. + */ +public class RESTTokenCredentialsProvider implements CredentialsProvider { + + /** Prefix under which {@link OSSFileIO} passes the catalog options to this provider. */ + static final String CATALOG_OPTIONS_PREFIX = "fs.oss.rest-token.catalog."; + + private static final String ACCESS_KEY_ID = "fs.oss.accessKeyId"; + private static final String ACCESS_KEY_SECRET = "fs.oss.accessKeySecret"; + private static final String SECURITY_TOKEN = "fs.oss.securityToken"; + + private final RESTTokenRefresher refresher; + + // The last token and the credentials built from it, swapped together. + private volatile Resolved last; + + public RESTTokenCredentialsProvider(URI uri, Configuration conf) { + Options options = new Options(conf.getPropsWithPrefix(CATALOG_OPTIONS_PREFIX)); + checkArgument( + RESTTokenRefresher.isConfigured(options), + "No REST catalog table configured under '%s'.", + CATALOG_OPTIONS_PREFIX); + this.refresher = RESTTokenRefresher.fromOptions(options); + } + + @Override + public void setCredentials(Credentials credentials) { + throw new UnsupportedOperationException( + "Credentials come from the REST catalog and cannot be set."); + } + + @Override + public Credentials getCredentials() { + RESTToken token = refresher.token(); + Resolved resolved = last; + if (resolved != null && resolved.token == token) { + return resolved.credentials; + } + Credentials credentials = toCredentials(token.token()); + last = new Resolved(token, credentials); + return credentials; + } + + private static Credentials toCredentials(Map token) { + String accessKeyId = null; + String accessKeySecret = null; + String securityToken = null; + for (Map.Entry entry : token.entrySet()) { + String key = entry.getKey().toLowerCase(Locale.ROOT); + if (key.equals(ACCESS_KEY_ID.toLowerCase(Locale.ROOT))) { + accessKeyId = entry.getValue(); + } else if (key.equals(ACCESS_KEY_SECRET.toLowerCase(Locale.ROOT))) { + accessKeySecret = entry.getValue(); + } else if (key.equals(SECURITY_TOKEN.toLowerCase(Locale.ROOT))) { + securityToken = entry.getValue(); + } + } + checkArgument( + accessKeyId != null && accessKeySecret != null, + "The data token from the REST catalog has no OSS access key."); + return new DefaultCredentials(accessKeyId, accessKeySecret, securityToken); + } + + private static final class Resolved { + + private final RESTToken token; + private final Credentials credentials; + + private Resolved(RESTToken token, Credentials credentials) { + this.token = token; + this.credentials = credentials; + } + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java deleted file mode 100644 index 0db1010b525c..000000000000 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/RegisteredCredentialsProvider.java +++ /dev/null @@ -1,132 +0,0 @@ -/* - * 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.paimon.oss; - -import org.apache.paimon.fs.CredentialsSupplierRegistry; - -import com.aliyun.oss.common.auth.Credentials; -import com.aliyun.oss.common.auth.CredentialsProvider; -import com.aliyun.oss.common.auth.DefaultCredentials; -import org.apache.hadoop.conf.Configuration; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import javax.annotation.Nullable; - -import java.net.URI; -import java.util.Locale; -import java.util.Map; -import java.util.function.Supplier; - -/** - * An OSS {@link CredentialsProvider} that asks a supplier from {@link CredentialsSupplierRegistry} - * for credentials on every request, so streams that are already open also pick up refreshed ones. - * It holds the supplier for as long as the OSS client that uses it is alive. - */ -public class RegisteredCredentialsProvider implements CredentialsProvider { - - private static final Logger LOG = LoggerFactory.getLogger(RegisteredCredentialsProvider.class); - - static final String ACCESS_KEY_ID = "fs.oss.accessKeyId"; - static final String ACCESS_KEY_SECRET = "fs.oss.accessKeySecret"; - static final String SECURITY_TOKEN = "fs.oss.securityToken"; - - private final String supplierId; - @Nullable private final Supplier> supplier; - - // The last options from the supplier and the credentials built from them, swapped together. - private volatile Resolved last; - - public RegisteredCredentialsProvider(URI uri, Configuration conf) { - this.supplierId = conf.get(CredentialsSupplierRegistry.SUPPLIER_ID); - this.supplier = supplierId == null ? null : CredentialsSupplierRegistry.get(supplierId); - String accessKeyId = conf.get(ACCESS_KEY_ID); - String accessKeySecret = conf.get(ACCESS_KEY_SECRET); - if (accessKeyId != null && accessKeySecret != null) { - this.last = - new Resolved( - null, - new DefaultCredentials( - accessKeyId, accessKeySecret, conf.get(SECURITY_TOKEN))); - } - } - - @Override - public void setCredentials(Credentials credentials) { - this.last = new Resolved(null, credentials); - } - - @Override - public Credentials getCredentials() { - Resolved resolved = last; - if (supplier != null) { - try { - return resolve(supplier.get(), resolved); - } catch (RuntimeException e) { - if (resolved == null) { - throw e; - } - LOG.warn("Failed to refresh OSS credentials, reusing the last ones.", e); - } - } - if (resolved == null) { - throw new IllegalStateException( - "No OSS credentials available for credentials supplier " + supplierId); - } - return resolved.credentials; - } - - private Credentials resolve(Map options, Resolved resolved) { - if (resolved != null && resolved.options == options) { - return resolved.credentials; - } - String accessKeyId = null; - String accessKeySecret = null; - String securityToken = null; - for (Map.Entry entry : options.entrySet()) { - String key = entry.getKey().toLowerCase(Locale.ROOT); - if (key.equals(ACCESS_KEY_ID.toLowerCase(Locale.ROOT))) { - accessKeyId = entry.getValue(); - } else if (key.equals(ACCESS_KEY_SECRET.toLowerCase(Locale.ROOT))) { - accessKeySecret = entry.getValue(); - } else if (key.equals(SECURITY_TOKEN.toLowerCase(Locale.ROOT))) { - securityToken = entry.getValue(); - } - } - if (accessKeyId == null || accessKeySecret == null) { - throw new IllegalStateException( - "Credentials supplier " + supplierId + " returned no OSS access key."); - } - Credentials credentials = - new DefaultCredentials(accessKeyId, accessKeySecret, securityToken); - last = new Resolved(options, credentials); - return credentials; - } - - private static final class Resolved { - - private final Map options; - private final Credentials credentials; - - private Resolved(Map options, Credentials credentials) { - this.options = options; - this.credentials = credentials; - } - } -} diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java new file mode 100644 index 000000000000..7812a12b4e1f --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java @@ -0,0 +1,133 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.options.Options; +import org.apache.paimon.rest.RESTApi; +import org.apache.paimon.rest.responses.GetTableTokenResponse; + +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; + +import java.io.IOException; +import java.io.OutputStream; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +/** A local REST catalog token endpoint and OSS endpoint that record what they are asked. */ +public class FakeRESTAndOSSServer implements AutoCloseable { + + public static final int OBJECT_SIZE = 4096; + + private final HttpServer server; + private final List tokens = new ArrayList<>(); + private final AtomicInteger tokenRequests = new AtomicInteger(); + private final List ossGetAuthorizations = + Collections.synchronizedList(new ArrayList<>()); + + public FakeRESTAndOSSServer() throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/v1/", this::handleToken); + server.createContext("/bucket/", this::handleObject); + server.start(); + } + + /** Tokens returned by successive token requests; the last one repeats. */ + public void addToken(String accessKeyId, long expiresAtMillis) { + tokens.add(new GetTableTokenResponse(ossToken(accessKeyId), expiresAtMillis)); + } + + public Map ossToken(String accessKeyId) { + Map token = new HashMap<>(); + token.put("fs.oss.endpoint", endpoint()); + token.put("fs.oss.accessKeyId", accessKeyId); + token.put("fs.oss.accessKeySecret", "secret-" + accessKeyId); + token.put("fs.oss.securityToken", "token-" + accessKeyId); + return token; + } + + /** Options a REST catalog client needs to reach this server. */ + public Options catalogOptions() { + Options options = new Options(); + options.set("uri", endpoint()); + options.set("token.provider", "bear"); + options.set("token", "catalog-token"); + options.set("prefix", "catalog"); + return options; + } + + public int tokenRequests() { + return tokenRequests.get(); + } + + public List ossGetAuthorizations() { + return ossGetAuthorizations; + } + + public String endpoint() { + return "http://127.0.0.1:" + server.getAddress().getPort(); + } + + @Override + public void close() { + server.stop(0); + } + + private void handleToken(HttpExchange exchange) throws IOException { + int index = tokenRequests.getAndIncrement(); + GetTableTokenResponse token = tokens.get(Math.min(index, tokens.size() - 1)); + byte[] body = RESTApi.toJson(token).getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().add("Content-Type", "application/json"); + exchange.sendResponseHeaders(200, body.length); + try (OutputStream out = exchange.getResponseBody()) { + out.write(body); + } + exchange.close(); + } + + private void handleObject(HttpExchange exchange) throws IOException { + exchange.getResponseHeaders().add("Last-Modified", "Thu, 24 Sep 2026 00:00:00 GMT"); + if ("HEAD".equals(exchange.getRequestMethod())) { + exchange.getResponseHeaders().add("Content-Length", String.valueOf(OBJECT_SIZE)); + exchange.sendResponseHeaders(200, -1); + } else { + ossGetAuthorizations.add(exchange.getRequestHeaders().getFirst("Authorization")); + String[] range = + exchange.getRequestHeaders() + .getFirst("Range") + .substring("bytes=".length()) + .split("-"); + int start = Integer.parseInt(range[0]); + int end = Math.min(Integer.parseInt(range[1]), OBJECT_SIZE - 1); + exchange.getResponseHeaders() + .add("Content-Range", "bytes " + start + "-" + end + "/" + OBJECT_SIZE); + exchange.sendResponseHeaders(206, end - start + 1); + try (OutputStream out = exchange.getResponseBody()) { + out.write(new byte[end - start + 1]); + } + } + exchange.close(); + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java new file mode 100644 index 000000000000..32c0ec1ad256 --- /dev/null +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java @@ -0,0 +1,124 @@ +/* + * 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.paimon.oss; + +import org.apache.paimon.catalog.CatalogContext; +import org.apache.paimon.catalog.Identifier; +import org.apache.paimon.fs.Path; +import org.apache.paimon.options.Options; +import org.apache.paimon.rest.RESTTokenRefresher; + +import com.aliyun.oss.OSSClient; +import org.apache.hadoop.conf.Configuration; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +import java.time.Duration; + +import static org.assertj.core.api.Assertions.assertThat; + +/** Tests for {@link RESTTokenCredentialsProvider}. */ +public class RESTTokenCredentialsProviderTest { + + private FakeRESTAndOSSServer server; + + @BeforeEach + public void startServer() throws Exception { + server = new FakeRESTAndOSSServer(); + } + + @AfterEach + public void stopServer() { + server.close(); + } + + @Test + public void testUsesTheTokenInTheOptionsWhileItHasTimeLeft() { + RESTTokenCredentialsProvider provider = provider(tableOptions("ak-1", Duration.ofHours(2))); + + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-1"); + assertThat(provider.getCredentials().getSecurityToken()).isEqualTo("token-ak-1"); + assertThat(server.tokenRequests()).isZero(); + } + + @Test + public void testLoadsANewTokenFromTheCatalogWhenTheCurrentOneIsAboutToExpire() { + server.addToken("ak-2", System.currentTimeMillis() + Duration.ofHours(4).toMillis()); + RESTTokenCredentialsProvider provider = + provider(tableOptions("ak-1", Duration.ofMinutes(30))); + + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + int requests = server.tokenRequests(); + assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); + assertThat(server.tokenRequests()).isEqualTo(requests); + } + + @Test + public void testOSSFileIOInstallsTheProviderOnlyForARESTTable() throws Exception { + OSSFileIO restFileIO = new OSSFileIO(); + OSSFileIO plainFileIO = new OSSFileIO(); + try { + restFileIO.configure(CatalogContext.create(tableOptions("ak-1", Duration.ofHours(2)))); + OSSClient client = restFileIO.ossClient(new Path("oss://bucket/key")); + assertThat(client.getCredentialsProvider()) + .isInstanceOf(RESTTokenCredentialsProvider.class); + assertThat(client.getCredentialsProvider().getCredentials().getAccessKeyId()) + .isEqualTo("ak-1"); + // readers that build their own clients from these options keep the static keys + assertThat(restFileIO.hadoopOptions().keySet()) + .doesNotContain("fs.oss.credentials.provider") + .noneMatch(key -> key.startsWith("fs.oss.rest-token.")); + + Options plain = new Options(); + plain.set("fs.oss.endpoint", server.endpoint()); + plain.set("fs.oss.accessKeyId", "ak-1"); + plain.set("fs.oss.accessKeySecret", "secret-ak-1"); + plainFileIO.configure(CatalogContext.create(plain)); + assertThat(plainFileIO.ossClient(new Path("oss://bucket/key")).getCredentialsProvider()) + .isNotInstanceOf(RESTTokenCredentialsProvider.class); + } finally { + restFileIO.close(); + plainFileIO.close(); + } + } + + /** Catalog options with a merged token, as RESTTokenFileIO hands them to its delegate. */ + private Options tableOptions(String accessKeyId, Duration lifetime) { + Options options = server.catalogOptions(); + server.ossToken(accessKeyId).forEach(options::set); + options.set("file-io.allow-cache", "false"); + RESTTokenRefresher.configure( + options, + Identifier.create("db", "table"), + System.currentTimeMillis() + lifetime.toMillis()); + return options; + } + + private static RESTTokenCredentialsProvider provider(Options options) { + Configuration conf = new Configuration(false); + options.toMap() + .forEach( + (key, value) -> + conf.set( + RESTTokenCredentialsProvider.CATALOG_OPTIONS_PREFIX + key, + value)); + return new RESTTokenCredentialsProvider(null, conf); + } +} diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java deleted file mode 100644 index 96d9a7608f67..000000000000 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RegisteredCredentialsProviderTest.java +++ /dev/null @@ -1,175 +0,0 @@ -/* - * 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.paimon.oss; - -import org.apache.paimon.catalog.CatalogContext; -import org.apache.paimon.fs.CredentialsSupplierRegistry; -import org.apache.paimon.fs.Path; -import org.apache.paimon.options.Options; - -import com.aliyun.oss.ClientConfiguration; -import com.aliyun.oss.ClientException; -import com.aliyun.oss.OSSClient; -import com.aliyun.oss.common.comm.ExecutionContext; -import com.aliyun.oss.common.comm.RequestMessage; -import com.aliyun.oss.common.comm.ResponseMessage; -import com.aliyun.oss.common.comm.RetryStrategy; -import com.aliyun.oss.common.comm.ServiceClient; -import com.aliyun.oss.internal.OSSObjectOperation; -import com.aliyun.oss.model.GenericRequest; -import org.apache.hadoop.conf.Configuration; -import org.junit.jupiter.api.Test; - -import java.net.URI; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicReference; -import java.util.function.Supplier; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.assertj.core.api.Assertions.assertThatThrownBy; - -/** Tests for {@link RegisteredCredentialsProvider}. */ -public class RegisteredCredentialsProviderTest { - - @Test - public void testEveryRequestIsSignedWithTheCurrentCredentials() throws Exception { - AtomicReference accessKeyId = new AtomicReference<>("ak-1"); - Supplier> supplier = () -> credentials(accessKeyId.get()); - String id = CredentialsSupplierRegistry.register(supplier); - List authorizations = new ArrayList<>(); - OSSObjectOperation operation = - new OSSObjectOperation(capturingClient(authorizations), provider(id, null)); - operation.setEndpoint(new URI("http://oss-cn-hangzhou.aliyuncs.com")); - - GenericRequest request = new GenericRequest("bucket", "key"); - assertThatThrownBy(() -> operation.getObjectMetadata(request)) - .isInstanceOf(ClientException.class); - accessKeyId.set("ak-2"); - assertThatThrownBy(() -> operation.getObjectMetadata(request)) - .isInstanceOf(ClientException.class); - - assertThat(authorizations).hasSize(2); - assertThat(authorizations.get(0)).contains("ak-1").doesNotContain("ak-2"); - assertThat(authorizations.get(1)).contains("ak-2").doesNotContain("ak-1"); - } - - @Test - public void testUsesConfiguredCredentialsWithoutARegisteredSupplier() { - RegisteredCredentialsProvider provider = provider("unknown-id", "ak-1"); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-1"); - } - - @Test - public void testKeepsTheLastCredentialsWhenTheSupplierFails() { - AtomicBoolean failing = new AtomicBoolean(false); - Supplier> supplier = - () -> { - if (failing.get()) { - throw new IllegalStateException("REST server unavailable"); - } - return credentials("ak-2"); - }; - String id = CredentialsSupplierRegistry.register(supplier); - RegisteredCredentialsProvider provider = provider(id, null); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); - - failing.set(true); - assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); - assertThatThrownBy(() -> provider(id, null).getCredentials()) - .hasMessageContaining("REST server unavailable"); - assertThat(supplier).isNotNull(); - } - - @Test - public void testOSSFileIOResolvesCredentialsFromTheRegisteredSupplier() throws Exception { - Supplier> supplier = () -> credentials("ak-2"); - String id = CredentialsSupplierRegistry.register(supplier); - OSSFileIO fileIO = new OSSFileIO(); - try { - Options options = new Options(); - options.set("fs.oss.endpoint", "http://oss-cn-hangzhou.aliyuncs.com"); - options.set("fs.oss.accessKeyId", "ak-1"); - options.set("fs.oss.accessKeySecret", "sk-1"); - options.set(CredentialsSupplierRegistry.SUPPLIER_ID, id); - fileIO.configure(CatalogContext.create(options)); - - OSSClient client = fileIO.ossClient(new Path("oss://bucket/key")); - assertThat(client.getCredentialsProvider()) - .isInstanceOf(RegisteredCredentialsProvider.class); - assertThat(client.getCredentialsProvider().getCredentials().getAccessKeyId()) - .isEqualTo("ak-2"); - // readers that build their own clients from these options keep the static keys - assertThat(fileIO.hadoopOptions().keySet()) - .doesNotContain("fs.oss.credentials.provider"); - } finally { - fileIO.close(); - } - } - - private static Map credentials(String accessKeyId) { - Map credentials = new HashMap<>(); - credentials.put("fs.oss.accessKeyId", accessKeyId); - credentials.put("fs.oss.accessKeySecret", "secret-" + accessKeyId); - credentials.put("fs.oss.securityToken", "token-" + accessKeyId); - return credentials; - } - - private static RegisteredCredentialsProvider provider(String id, String accessKeyId) { - Configuration conf = new Configuration(false); - conf.set(CredentialsSupplierRegistry.SUPPLIER_ID, id); - if (accessKeyId != null) { - conf.set("fs.oss.accessKeyId", accessKeyId); - conf.set("fs.oss.accessKeySecret", "secret-" + accessKeyId); - } - return new RegisteredCredentialsProvider(null, conf); - } - - /** Records the Authorization header of each request instead of sending it. */ - private static ServiceClient capturingClient(List authorizations) { - return new ServiceClient(new ClientConfiguration()) { - @Override - protected ResponseMessage sendRequestCore( - ServiceClient.Request request, ExecutionContext context) { - authorizations.add(request.getHeaders().get("Authorization")); - throw new ClientException("captured"); - } - - @Override - protected RetryStrategy getDefaultRetryStrategy() { - return new RetryStrategy() { - @Override - public boolean shouldRetry( - Exception e, - RequestMessage request, - ResponseMessage response, - int retries) { - return false; - } - }; - } - - @Override - public void shutdown() {} - }; - } -} diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java index 295b4207f223..2e0f555caf5c 100644 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java @@ -26,132 +26,58 @@ import org.apache.paimon.fs.PluginFileIO; import org.apache.paimon.fs.SeekableInputStream; import org.apache.paimon.options.Options; +import org.apache.paimon.oss.FakeRESTAndOSSServer; import org.apache.paimon.oss.OSSFileIO; -import org.apache.paimon.rest.responses.GetTableTokenResponse; -import com.sun.net.httpserver.HttpServer; import org.apache.hadoop.conf.Configuration; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; -import java.io.IOException; -import java.io.OutputStream; -import java.net.InetSocketAddress; -import java.time.Duration; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; import java.util.List; -import java.util.Map; -import java.util.concurrent.atomic.AtomicLong; +import static org.apache.paimon.rest.RESTApi.TOKEN_EXPIRATION_SAFE_TIME_MILLIS; import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.times; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; -/** Tests {@link RESTTokenFileIO} over a real {@link OSSFileIO} against a local OSS endpoint. */ +/** Tests {@link RESTTokenFileIO} over a real {@link OSSFileIO} against a local catalog and OSS. */ public class RESTTokenFileIOOnOSSTest { - private static final byte[] CONTENT = new byte[4096]; - - private final List getAuthorizations = Collections.synchronizedList(new ArrayList<>()); - private HttpServer server; - - @BeforeEach - public void startServer() throws IOException { - server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); - server.createContext( - "/", - exchange -> { - exchange.getResponseHeaders() - .add("Last-Modified", "Thu, 24 Sep 2026 00:00:00 GMT"); - if ("HEAD".equals(exchange.getRequestMethod())) { - exchange.getResponseHeaders() - .add("Content-Length", String.valueOf(CONTENT.length)); - exchange.sendResponseHeaders(200, -1); - } else { - getAuthorizations.add( - exchange.getRequestHeaders().getFirst("Authorization")); - String[] range = - exchange.getRequestHeaders() - .getFirst("Range") - .substring("bytes=".length()) - .split("-"); - int start = Integer.parseInt(range[0]); - int end = Math.min(Integer.parseInt(range[1]), CONTENT.length - 1); - exchange.getResponseHeaders() - .add( - "Content-Range", - "bytes " + start + "-" + end + "/" + CONTENT.length); - exchange.sendResponseHeaders(206, end - start + 1); - try (OutputStream out = exchange.getResponseBody()) { - out.write(CONTENT, start, end - start + 1); - } - } - exchange.close(); - }); - server.start(); - } - - @AfterEach - public void stopServer() { - server.stop(0); - } - /** An open stream keeps refreshing even after its delegate FileIO left the cache. */ @Test public void testOpenStreamPicksUpRefreshedTokenAfterCacheEviction() throws Exception { - Identifier identifier = Identifier.create("db", "table"); - AtomicLong now = new AtomicLong(1700000000000L); - RESTApi api = mock(RESTApi.class); - when(api.loadTableToken(identifier)) - .thenReturn( - new GetTableTokenResponse( - token("ak-1"), now.get() + Duration.ofHours(2).toMillis()), - new GetTableTokenResponse( - token("ak-2"), now.get() + Duration.ofHours(4).toMillis())); - Options options = new Options(); - options.set("fs.oss.multipart.download.size", "1024"); - options.set("fs.oss.multipart.download.threads", "1"); - RESTTokenFileIO fileIO = - new RESTTokenFileIO( - CatalogContext.create( - options, new Configuration(false), new PluginOSSLoader(), null), - api, - identifier, - new Path("oss://bucket/table")) { - @Override - long currentTimeMillis() { - return now.get(); - } - }; - - try (SeekableInputStream in = fileIO.newInputStream(new Path("oss://bucket/table/f"))) { - in.read(new byte[1024]); + try (FakeRESTAndOSSServer server = new FakeRESTAndOSSServer()) { + // the first token enters the refresh window five seconds from now + long firstExpiresAt = + System.currentTimeMillis() + TOKEN_EXPIRATION_SAFE_TIME_MILLIS + 5000; + server.addToken("ak-1", firstExpiresAt); + server.addToken("ak-2", firstExpiresAt + TOKEN_EXPIRATION_SAFE_TIME_MILLIS * 4); + Options options = server.catalogOptions(); + options.set("fs.oss.multipart.download.size", "1024"); + options.set("fs.oss.multipart.download.threads", "1"); + RESTTokenFileIO fileIO = + new RESTTokenFileIO( + CatalogContext.create( + options, new Configuration(false), new PluginOSSLoader(), null), + null, + Identifier.create("db", "table"), + new Path("oss://bucket/table")); + + try (SeekableInputStream in = + fileIO.newInputStream(new Path("oss://bucket/table/data"))) { + in.read(new byte[1024]); + + RESTTokenFileIO.invalidateFileIOCache(); + System.gc(); + while (System.currentTimeMillis() + < firstExpiresAt - TOKEN_EXPIRATION_SAFE_TIME_MILLIS + 100) { + Thread.sleep(100); + } + in.seek(3000); + in.read(new byte[1024]); + } - RESTTokenFileIO.invalidateFileIOCache(); - System.gc(); - // 30 minutes left is inside the refresh window - now.addAndGet(Duration.ofMinutes(90).toMillis()); - in.seek(3000); - in.read(new byte[1024]); + List authorizations = server.ossGetAuthorizations(); + assertThat(authorizations.get(0)).contains("ak-1"); + assertThat(authorizations.get(authorizations.size() - 1)).contains("ak-2"); } - - assertThat(getAuthorizations.get(0)).contains("ak-1"); - assertThat(getAuthorizations.get(getAuthorizations.size() - 1)).contains("ak-2"); - verify(api, times(2)).loadTableToken(identifier); - } - - private Map token(String accessKeyId) { - Map token = new HashMap<>(); - token.put("fs.oss.endpoint", "http://127.0.0.1:" + server.getAddress().getPort()); - token.put("fs.oss.accessKeyId", accessKeyId); - token.put("fs.oss.accessKeySecret", "secret-" + accessKeyId); - token.put("fs.oss.securityToken", "token-" + accessKeyId); - return token; } /** Wraps {@link OSSFileIO} the way the OSS plugin does, including its no-op close. */ From 114d251666db565f35f58818d0de792debd8ed69 Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Mon, 28 Sep 2026 13:34:46 +0800 Subject: [PATCH 4/6] [oss] Drive the token refresh in RESTTokenFileIOOnOSSTest with a test clock The test waited on the wall clock for the first token to enter the refresh window, so on a slow CI runner the refresh happened before the first read. It now moves the refreshers' clock forward instead. --- .../apache/paimon/rest/RESTTokenRefresher.java | 6 +++++- .../paimon/rest/RESTTokenFileIOOnOSSTest.java | 18 ++++++++---------- 2 files changed, 13 insertions(+), 11 deletions(-) diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java index ae481ff5ed32..09251f6bfc91 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java @@ -18,6 +18,7 @@ package org.apache.paimon.rest; +import org.apache.paimon.annotation.VisibleForTesting; import org.apache.paimon.catalog.Identifier; import org.apache.paimon.options.Options; import org.apache.paimon.rest.responses.GetTableTokenResponse; @@ -49,6 +50,9 @@ public class RESTTokenRefresher { // After a failed refresh, keep using the current token this long before trying again. static final long RETRY_INTERVAL_MILLIS = 10_000L; + // Shifts the clock of refreshers that providers create from options, for tests only. + @VisibleForTesting static volatile long clockOffsetMillis; + private final Options catalogOptions; private final Identifier identifier; @@ -140,6 +144,6 @@ private RESTToken load() { } long currentTimeMillis() { - return System.currentTimeMillis(); + return System.currentTimeMillis() + clockOffsetMillis; } } diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java index 2e0f555caf5c..62867268c9a7 100644 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java @@ -32,9 +32,9 @@ import org.apache.hadoop.conf.Configuration; import org.junit.jupiter.api.Test; +import java.time.Duration; import java.util.List; -import static org.apache.paimon.rest.RESTApi.TOKEN_EXPIRATION_SAFE_TIME_MILLIS; import static org.assertj.core.api.Assertions.assertThat; /** Tests {@link RESTTokenFileIO} over a real {@link OSSFileIO} against a local catalog and OSS. */ @@ -44,11 +44,9 @@ public class RESTTokenFileIOOnOSSTest { @Test public void testOpenStreamPicksUpRefreshedTokenAfterCacheEviction() throws Exception { try (FakeRESTAndOSSServer server = new FakeRESTAndOSSServer()) { - // the first token enters the refresh window five seconds from now - long firstExpiresAt = - System.currentTimeMillis() + TOKEN_EXPIRATION_SAFE_TIME_MILLIS + 5000; - server.addToken("ak-1", firstExpiresAt); - server.addToken("ak-2", firstExpiresAt + TOKEN_EXPIRATION_SAFE_TIME_MILLIS * 4); + long now = System.currentTimeMillis(); + server.addToken("ak-1", now + Duration.ofHours(2).toMillis()); + server.addToken("ak-2", now + Duration.ofHours(4).toMillis()); Options options = server.catalogOptions(); options.set("fs.oss.multipart.download.size", "1024"); options.set("fs.oss.multipart.download.threads", "1"); @@ -66,12 +64,12 @@ options, new Configuration(false), new PluginOSSLoader(), null), RESTTokenFileIO.invalidateFileIOCache(); System.gc(); - while (System.currentTimeMillis() - < firstExpiresAt - TOKEN_EXPIRATION_SAFE_TIME_MILLIS + 100) { - Thread.sleep(100); - } + // 30 minutes left is inside the refresh window + RESTTokenRefresher.clockOffsetMillis = Duration.ofMinutes(90).toMillis(); in.seek(3000); in.read(new byte[1024]); + } finally { + RESTTokenRefresher.clockOffsetMillis = 0; } List authorizations = server.ossGetAuthorizations(); From 002bdccd1ec2f750f123f35eb3d79b27cede0dc1 Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Mon, 28 Sep 2026 15:22:51 +0800 Subject: [PATCH 5/6] [oss] Stop OSS requests from waiting on a data token reload RESTTokenRefresher held its lock across the catalog request, whose client may retry for minutes, so every OSS request on the client waited while the current token was still valid. Now one caller reloads and the others keep the valid token; with an expired token, callers fail fast within the retry interval. A reloaded token is layered over the catalog options as RESTTokenFileIO does, OSSFileIO reads only a constant from RESTTokenRefresher so an older paimon-common still loads it, and OSSLoader names its option keys. --- docs/docs/maintenance/filesystems.mdx | 4 +- .../apache/paimon/rest/RESTTokenFileIO.java | 2 +- .../paimon/rest/RESTTokenRefresher.java | 82 +++++++++++++------ .../paimon/rest/RESTTokenRefresherTest.java | 37 ++++++++- .../java/org/apache/paimon/oss/OSSFileIO.java | 3 +- .../java/org/apache/paimon/oss/OSSLoader.java | 12 ++- 6 files changed, 108 insertions(+), 32 deletions(-) diff --git a/docs/docs/maintenance/filesystems.mdx b/docs/docs/maintenance/filesystems.mdx index f69aaf3ebbbf..f1b837cd4f03 100644 --- a/docs/docs/maintenance/filesystems.mdx +++ b/docs/docs/maintenance/filesystems.mdx @@ -335,9 +335,9 @@ Download [paimon-jindo-@@VERSION@@.jar](https://repository.apache.org/snapshots/ ### Credentials Provider -Instead of `fs.oss.accessKeyId` and `fs.oss.accessKeySecret`, you can set `fs.oss.credentials.provider` to a class implementing `com.aliyun.oss.common.auth.CredentialsProvider`, for example `com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider`. The provider is asked for credentials on every request. +`fs.oss.credentials.provider` can replace `fs.oss.accessKeyId` and `fs.oss.accessKeySecret` with a `com.aliyun.oss.common.auth.CredentialsProvider` class, such as `com.aliyun.oss.common.auth.EnvironmentVariableCredentialsProvider`. -With a REST catalog that vends data tokens, the OSS FileIO refreshes the table's token from the catalog by itself before it expires, so streams that are already open keep working after the token is refreshed. +With a REST catalog, the data token is reloaded before it expires, also for streams that are already open. ### Server-Side Encryption diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java index 25a2dd9924be..50fe2cd52cc7 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java @@ -266,7 +266,7 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc Options options = catalogContext.options(); options = new Options(RESTUtil.merge(options.toMap(), currentToken.token())); options.set(FILE_IO_ALLOW_CACHE, false); - // Lets a FileIO that supports it, such as OSS, refresh this token by itself. + // Record whose data token the delegate holds, so it can reload it from the catalog. RESTTokenRefresher.configure(options, tokenIdentifier(), currentToken.expireAtMillis()); CatalogContext context = CatalogContext.create( diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java index 09251f6bfc91..5a0e70c67209 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java @@ -28,6 +28,8 @@ import javax.annotation.Nullable; +import java.util.concurrent.locks.ReentrantLock; + import static org.apache.paimon.rest.RESTApi.TOKEN_EXPIRATION_SAFE_TIME_MILLIS; /** @@ -47,7 +49,7 @@ public class RESTTokenRefresher { /** Expiration of the token already present in the options. */ public static final String EXPIRES_AT_MILLIS = "data-token.expires-at-millis"; - // After a failed refresh, keep using the current token this long before trying again. + // After a failed reload, wait this long before asking the catalog again. static final long RETRY_INTERVAL_MILLIS = 10_000L; // Shifts the clock of refreshers that providers create from options, for tests only. @@ -55,10 +57,14 @@ public class RESTTokenRefresher { private final Options catalogOptions; private final Identifier identifier; + private final ReentrantLock lock = new ReentrantLock(); - @Nullable private RESTApi api; @Nullable private volatile RESTToken token; + + // Guarded by lock. + @Nullable private RESTApi api; private long nextAttemptMillis; + @Nullable private RuntimeException lastFailure; RESTTokenRefresher( Options catalogOptions, @@ -97,40 +103,65 @@ public static RESTTokenRefresher fromOptions(Options options) { return new RESTTokenRefresher(options, identifier, null, token); } - /** Returns the current token, loading a new one when it is about to expire. */ + /** Returns the current token, reloading it from the catalog when it is about to expire. */ public RESTToken token() { RESTToken current = token; if (current != null && !expiresSoon(current)) { return current; } - synchronized (this) { - current = token; - if (current != null && !expiresSoon(current)) { + // While the token is still valid, one caller reloads it and the others keep using it. + if (isValid(current)) { + if (!lock.tryLock()) { return current; } - long now = currentTimeMillis(); - boolean usable = current != null && now < current.expireAtMillis(); - if (usable && now < nextAttemptMillis) { + } else { + lock.lock(); + } + try { + return reload(); + } finally { + lock.unlock(); + } + } + + private RESTToken reload() { + RESTToken current = token; + if (current != null && !expiresSoon(current)) { + return current; + } + boolean valid = isValid(current); + long now = currentTimeMillis(); + if (now < nextAttemptMillis) { + if (valid) { return current; } - try { - RESTToken loaded = load(); - token = loaded; - return loaded; - } catch (RuntimeException e) { - if (!usable) { - throw e; - } - nextAttemptMillis = now + RETRY_INTERVAL_MILLIS; - LOG.warn( - "Failed to refresh the data token of {}, keeping the current one.", - identifier, - e); - return current; + throw new IllegalStateException( + "The data token of " + identifier + " expired and reloading it failed.", + lastFailure); + } + try { + RESTToken loaded = load(); + token = loaded; + lastFailure = null; + return loaded; + } catch (RuntimeException e) { + lastFailure = e; + nextAttemptMillis = now + RETRY_INTERVAL_MILLIS; + if (!valid) { + throw e; } + LOG.warn( + "Failed to reload the data token of {}, keeping the current one.", + identifier, + e); + return current; } } + private boolean isValid(@Nullable RESTToken token) { + return token != null && currentTimeMillis() < token.expireAtMillis(); + } + private boolean expiresSoon(RESTToken token) { return token.expireAtMillis() - currentTimeMillis() < TOKEN_EXPIRATION_SAFE_TIME_MILLIS; } @@ -140,7 +171,10 @@ private RESTToken load() { api = new RESTApi(catalogOptions, false); } GetTableTokenResponse response = api.loadTableToken(identifier); - return new RESTToken(response.getToken(), response.getExpiresAtMillis()); + // Layered over the catalog options, the same way RESTTokenFileIO builds its delegate. + return new RESTToken( + RESTUtil.merge(catalogOptions.toMap(), response.getToken()), + response.getExpiresAtMillis()); } long currentTimeMillis() { diff --git a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java index 3fcf1e4b5d81..2a10eb22f541 100644 --- a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java @@ -26,6 +26,10 @@ import java.time.Duration; import java.util.Collections; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicLong; import static org.assertj.core.api.Assertions.assertThat; @@ -77,12 +81,43 @@ void testKeepsTheCurrentTokenWhileRefreshFailsAndRetriesLater() { } @Test - void testFailsWhenTheTokenExpiredAndRefreshFails() { + void testOtherCallersKeepTheValidTokenWhileOneReloads() throws Exception { + CountDownLatch loading = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + when(api.loadTableToken(identifier)) + .thenAnswer( + invocation -> { + loading.countDown(); + release.await(); + return response("second", hours(4)); + }); + RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(30))); + ExecutorService executor = Executors.newSingleThreadExecutor(); + try { + Future reloading = executor.submit(refresher::token); + loading.await(); + + // does not wait for the reload in progress + assertThat(refresher.token().token()).containsEntry("k", "first"); + release.countDown(); + assertThat(reloading.get().token()).containsEntry("k", "second"); + } finally { + release.countDown(); + executor.shutdownNow(); + } + } + + @Test + void testExpiredTokenFailsFastWithinTheRetryInterval() { when(api.loadTableToken(identifier)) .thenThrow(new IllegalStateException("REST server unavailable")); RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(-1))); assertThatThrownBy(refresher::token).hasMessageContaining("REST server unavailable"); + assertThatThrownBy(refresher::token) + .hasMessageContaining("reloading it failed") + .hasRootCauseMessage("REST server unavailable"); + verify(api, times(1)).loadTableToken(identifier); } @Test diff --git a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java index 21667818b41c..b726120e7713 100644 --- a/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java +++ b/paimon-filesystems/paimon-oss-impl/src/main/java/org/apache/paimon/oss/OSSFileIO.java @@ -148,8 +148,9 @@ public boolean isObjectStore() { @Override public void configure(CatalogContext context) { + // Only a constant is read here, so an older paimon-common without the class still works. restTokenOptions = - RESTTokenRefresher.isConfigured(context.options()) + context.options().containsKey(RESTTokenRefresher.DATABASE) ? context.options().toMap() : null; // The file system refreshes a table's token, so it must not be shared through the cache. diff --git a/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java b/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java index c78fcbc7cfe8..e4086b3a0de6 100644 --- a/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java +++ b/paimon-filesystems/paimon-oss/src/main/java/org/apache/paimon/oss/OSSLoader.java @@ -37,6 +37,11 @@ public class OSSLoader implements FileIOLoader { private static final String OSS_CLASS = "org.apache.paimon.oss.OSSFileIO"; + private static final String ENDPOINT = "fs.oss.endpoint"; + private static final String ACCESS_KEY_ID = "fs.oss.accessKeyId"; + private static final String ACCESS_KEY_SECRET = "fs.oss.accessKeySecret"; + private static final String CREDENTIALS_PROVIDER = "fs.oss.credentials.provider"; + // Singleton lazy initialization private static PluginLoader loader; @@ -57,9 +62,10 @@ public String getScheme() { @Override public List requiredOptions() { List options = new ArrayList<>(); - options.add(new String[] {"fs.oss.endpoint"}); - options.add(new String[] {"fs.oss.accessKeyId", "fs.oss.credentials.provider"}); - options.add(new String[] {"fs.oss.accessKeySecret", "fs.oss.credentials.provider"}); + options.add(new String[] {ENDPOINT}); + // Each entry lists alternatives: a credentials provider can replace the access keys. + options.add(new String[] {ACCESS_KEY_ID, CREDENTIALS_PROVIDER}); + options.add(new String[] {ACCESS_KEY_SECRET, CREDENTIALS_PROVIDER}); return options; } From 5d6796f9f847533a4e3022517d3272b602dddfdb Mon Sep 17 00:00:00 2001 From: Sun Dapeng Date: Mon, 28 Sep 2026 18:03:58 +0800 Subject: [PATCH 6/6] [oss] Share a refreshing delegate only within one table and catalog user RESTTokenFileIO cached delegates by token alone. Now that a delegate reloads the token as a specific table and catalog user, another table or user holding the same token reused it and reloaded as the wrong one. The cache key now also covers the token's table and the catalog options. RESTTokenRefresher reloaded whenever less than an hour was left, so a token that lives shorter was reloaded on every OSS request. Each token now carries its own reload time: an hour before expiry, or halfway through when it lives shorter, and at least ten seconds after it arrived. An expired token is never returned. --- .../apache/paimon/rest/RESTTokenFileIO.java | 86 ++++++++++++++++++- .../paimon/rest/RESTTokenRefresher.java | 61 +++++++------ .../paimon/rest/RESTTokenFileIOTest.java | 33 +++++++ .../paimon/rest/RESTTokenRefresherTest.java | 24 +++++- .../paimon/oss/FakeRESTAndOSSServer.java | 6 ++ .../oss/RESTTokenCredentialsProviderTest.java | 4 +- .../paimon/rest/RESTTokenFileIOOnOSSTest.java | 47 +++++++--- 7 files changed, 217 insertions(+), 44 deletions(-) diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java index 50fe2cd52cc7..3ed7657facce 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenFileIO.java @@ -48,6 +48,7 @@ import java.io.IOException; import java.time.Duration; import java.util.Map; +import java.util.Objects; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; @@ -68,7 +69,7 @@ public class RESTTokenFileIO implements FileIO { .defaultValue(false) .withDescription("Whether to support data token provided by the REST server."); - private static final Cache FILE_IO_CACHE = + private static final Cache FILE_IO_CACHE = Caffeine.newBuilder() .maximumSize(1000) .expireAfterAccess(10, TimeUnit.HOURS) @@ -119,6 +120,9 @@ static long fileIOCacheMaximumSize() { // Server again after serialization private volatile RESTToken token; + // The table and catalog options a delegate refreshes its token with, computed once. + private transient volatile DelegateContext delegateContext; + public RESTTokenFileIO( CatalogContext catalogContext, RESTApi apiInstance, Identifier identifier, Path path) { this.catalogContext = catalogContext; @@ -252,13 +256,16 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc + "REST credential lifetime after refresh."); } - FileIO fileIO = FILE_IO_CACHE.getIfPresent(currentToken); + // A delegate reloads its token as a table and a catalog user, so it is shared only + // with callers that have the same ones. + DelegateKey key = new DelegateKey(currentToken, delegateContext()); + FileIO fileIO = FILE_IO_CACHE.getIfPresent(key); if (fileIO != null) { return new FileIOWithToken(fileIO, currentToken); } synchronized (FILE_IO_CACHE) { - fileIO = FILE_IO_CACHE.getIfPresent(currentToken); + fileIO = FILE_IO_CACHE.getIfPresent(key); if (fileIO != null) { return new FileIOWithToken(fileIO, currentToken); } @@ -275,7 +282,7 @@ private FileIOWithToken fileIOWithToken(long minimumValidityMillis) throws IOExc catalogContext.preferIO(), catalogContext.fallbackIO()); fileIO = FileIO.get(path, context); - FILE_IO_CACHE.put(currentToken, fileIO); + FILE_IO_CACHE.put(key, fileIO); return new FileIOWithToken(fileIO, currentToken); } } @@ -309,6 +316,68 @@ private boolean shouldRefresh(long minimumValidityMillis) { < Math.max(TOKEN_EXPIRATION_SAFE_TIME_MILLIS, minimumValidityMillis); } + /** The table and catalog options that a delegate reloads its token with. */ + private static final class DelegateContext { + + private final Identifier table; + private final Map catalogOptions; + private final int hash; + + private DelegateContext(Identifier table, Map catalogOptions) { + this.table = table; + this.catalogOptions = catalogOptions; + this.hash = Objects.hash(table, catalogOptions); + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (!(o instanceof DelegateContext)) { + return false; + } + DelegateContext that = (DelegateContext) o; + return hash == that.hash + && table.equals(that.table) + && catalogOptions.equals(that.catalogOptions); + } + + @Override + public int hashCode() { + return hash; + } + } + + /** Cache key of a delegate: its token and the context it reloads the token with. */ + private static final class DelegateKey { + + private final RESTToken token; + private final DelegateContext context; + + private DelegateKey(RESTToken token, DelegateContext context) { + this.token = token; + this.context = context; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (!(o instanceof DelegateKey)) { + return false; + } + DelegateKey that = (DelegateKey) o; + return token.equals(that.token) && context.equals(that.context); + } + + @Override + public int hashCode() { + return 31 * token.hashCode() + context.hashCode(); + } + } + private static class FileIOWithToken { private final FileIO fileIO; @@ -337,6 +406,15 @@ private void refreshToken() { response.getExpiresAtMillis()); } + private DelegateContext delegateContext() { + DelegateContext context = delegateContext; + if (context == null) { + context = new DelegateContext(tokenIdentifier(), catalogContext.options().toMap()); + delegateContext = context; + } + return context; + } + private Identifier tokenIdentifier() { if (identifier.isSystemTable()) { return new Identifier( diff --git a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java index 5a0e70c67209..f047f640e03b 100644 --- a/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java +++ b/paimon-common/src/main/java/org/apache/paimon/rest/RESTTokenRefresher.java @@ -49,8 +49,8 @@ public class RESTTokenRefresher { /** Expiration of the token already present in the options. */ public static final String EXPIRES_AT_MILLIS = "data-token.expires-at-millis"; - // After a failed reload, wait this long before asking the catalog again. - static final long RETRY_INTERVAL_MILLIS = 10_000L; + // A reloaded token is kept at least this long, and a failed reload is retried after it. + static final long RELOAD_INTERVAL_MILLIS = 10_000L; // Shifts the clock of refreshers that providers create from options, for tests only. @VisibleForTesting static volatile long clockOffsetMillis; @@ -59,7 +59,7 @@ public class RESTTokenRefresher { private final Identifier identifier; private final ReentrantLock lock = new ReentrantLock(); - @Nullable private volatile RESTToken token; + @Nullable private volatile CachedToken cached; // Guarded by lock. @Nullable private RESTApi api; @@ -74,7 +74,7 @@ public class RESTTokenRefresher { this.catalogOptions = catalogOptions; this.identifier = identifier; this.api = api; - this.token = token; + this.cached = token == null ? null : new CachedToken(token, currentTimeMillis()); } /** Names the table and the expiration of the merged token, so a refresher can be created. */ @@ -105,14 +105,15 @@ public static RESTTokenRefresher fromOptions(Options options) { /** Returns the current token, reloading it from the catalog when it is about to expire. */ public RESTToken token() { - RESTToken current = token; - if (current != null && !expiresSoon(current)) { - return current; + CachedToken current = cached; + long now = currentTimeMillis(); + if (current != null && now < current.reloadAtMillis) { + return current.token; } // While the token is still valid, one caller reloads it and the others keep using it. - if (isValid(current)) { + if (current != null && now < current.token.expireAtMillis()) { if (!lock.tryLock()) { - return current; + return current.token; } } else { lock.lock(); @@ -125,15 +126,15 @@ public RESTToken token() { } private RESTToken reload() { - RESTToken current = token; - if (current != null && !expiresSoon(current)) { - return current; - } - boolean valid = isValid(current); + CachedToken current = cached; long now = currentTimeMillis(); + if (current != null && now < current.reloadAtMillis) { + return current.token; + } + boolean valid = current != null && now < current.token.expireAtMillis(); if (now < nextAttemptMillis) { if (valid) { - return current; + return current.token; } throw new IllegalStateException( "The data token of " + identifier + " expired and reloading it failed.", @@ -141,12 +142,12 @@ private RESTToken reload() { } try { RESTToken loaded = load(); - token = loaded; + cached = new CachedToken(loaded, now); lastFailure = null; return loaded; } catch (RuntimeException e) { lastFailure = e; - nextAttemptMillis = now + RETRY_INTERVAL_MILLIS; + nextAttemptMillis = now + RELOAD_INTERVAL_MILLIS; if (!valid) { throw e; } @@ -154,18 +155,10 @@ private RESTToken reload() { "Failed to reload the data token of {}, keeping the current one.", identifier, e); - return current; + return current.token; } } - private boolean isValid(@Nullable RESTToken token) { - return token != null && currentTimeMillis() < token.expireAtMillis(); - } - - private boolean expiresSoon(RESTToken token) { - return token.expireAtMillis() - currentTimeMillis() < TOKEN_EXPIRATION_SAFE_TIME_MILLIS; - } - private RESTToken load() { if (api == null) { api = new RESTApi(catalogOptions, false); @@ -180,4 +173,20 @@ private RESTToken load() { long currentTimeMillis() { return System.currentTimeMillis() + clockOffsetMillis; } + + /** A token and the time to reload it, before it expires but not right after it arrived. */ + private static final class CachedToken { + + private final RESTToken token; + private final long reloadAtMillis; + + private CachedToken(RESTToken token, long now) { + long expiresAt = token.expireAtMillis(); + // Ahead of expiry by the safe time, or by half the time left when that is shorter. + long ahead = Math.min(TOKEN_EXPIRATION_SAFE_TIME_MILLIS, (expiresAt - now) / 2); + this.token = token; + this.reloadAtMillis = + Math.min(expiresAt, Math.max(expiresAt - ahead, now + RELOAD_INTERVAL_MILLIS)); + } + } } diff --git a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java index 1f63b10e9009..b4ed19d7b4c2 100644 --- a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenFileIOTest.java @@ -263,6 +263,39 @@ void testDelegateOptionsNameTheTableAndTokenExpiry() throws IOException { .isEqualTo(String.valueOf(expiresAt)); } + @Test + void testDelegateIsSharedOnlyBySameTableAndCatalogOptions() throws IOException { + FileIOLoader loader = mock(FileIOLoader.class); + when(loader.load(any())).thenAnswer(invocation -> mock(FileIO.class)); + when(loader.getScheme()).thenReturn("oss"); + RESTApi api = mock(RESTApi.class); + // every table and user gets the same token + when(api.loadTableToken(any())) + .thenReturn( + new GetTableTokenResponse( + Collections.singletonMap("token", UUID.randomUUID().toString()), + System.currentTimeMillis() + Duration.ofHours(2).toMillis())); + Options userA = new Options(); + userA.set("token", "user-a"); + Options userB = new Options(); + userB.set("token", "user-b"); + + FileIO delegate = restTokenFileIO(userA, loader, api, "table_a").fileIO(); + + assertThat(restTokenFileIO(userA, loader, api, "table_a").fileIO()).isSameAs(delegate); + assertThat(restTokenFileIO(userA, loader, api, "table_b").fileIO()).isNotSameAs(delegate); + assertThat(restTokenFileIO(userB, loader, api, "table_a").fileIO()).isNotSameAs(delegate); + } + + private static RESTTokenFileIO restTokenFileIO( + Options options, FileIOLoader loader, RESTApi api, String table) { + return new RESTTokenFileIO( + CatalogContext.create(options, loader, null), + api, + Identifier.create("db", table), + new Path("oss://bucket/" + table)); + } + @Test void testFileIOCreationFailureSurfacesAsCheckedIOException() throws IOException { Path tableRoot = new Path("resttoken-broken://bucket/table"); diff --git a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java index 2a10eb22f541..21dd44a1e89c 100644 --- a/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java +++ b/paimon-common/src/test/java/org/apache/paimon/rest/RESTTokenRefresherTest.java @@ -70,12 +70,14 @@ void testKeepsTheCurrentTokenWhileRefreshFailsAndRetriesLater() { .thenThrow(new IllegalStateException("REST server unavailable")) .thenReturn(response("second", hours(4))); RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(30))); + // past half of its lifetime, so it is due for a reload + now.addAndGet(Duration.ofMinutes(20).toMillis()); assertThat(refresher.token().token()).containsEntry("k", "first"); assertThat(refresher.token().token()).containsEntry("k", "first"); verify(api, times(1)).loadTableToken(identifier); - now.addAndGet(RESTTokenRefresher.RETRY_INTERVAL_MILLIS); + now.addAndGet(RESTTokenRefresher.RELOAD_INTERVAL_MILLIS); assertThat(refresher.token().token()).containsEntry("k", "second"); verify(api, times(2)).loadTableToken(identifier); } @@ -92,6 +94,7 @@ void testOtherCallersKeepTheValidTokenWhileOneReloads() throws Exception { return response("second", hours(4)); }); RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(30))); + now.addAndGet(Duration.ofMinutes(20).toMillis()); ExecutorService executor = Executors.newSingleThreadExecutor(); try { Future reloading = executor.submit(refresher::token); @@ -107,6 +110,23 @@ void testOtherCallersKeepTheValidTokenWhileOneReloads() throws Exception { } } + @Test + void testShortLivedTokenIsReusedUntilHalfOfItsLifetime() { + when(api.loadTableToken(identifier)) + .thenAnswer(invocation -> response("next", Duration.ofMinutes(30))); + RESTTokenRefresher refresher = refresher(token("first", Duration.ofMinutes(30))); + + now.addAndGet(Duration.ofMinutes(16).toMillis()); + for (int i = 0; i < 20; i++) { + assertThat(refresher.token().token()).containsEntry("k", "next"); + } + verify(api, times(1)).loadTableToken(identifier); + + now.addAndGet(Duration.ofMinutes(16).toMillis()); + refresher.token(); + verify(api, times(2)).loadTableToken(identifier); + } + @Test void testExpiredTokenFailsFastWithinTheRetryInterval() { when(api.loadTableToken(identifier)) @@ -153,7 +173,7 @@ private RESTToken token(String value, Duration lifetime) { private GetTableTokenResponse response(String value, Duration lifetime) { return new GetTableTokenResponse( - Collections.singletonMap("k", value), NOW + lifetime.toMillis()); + Collections.singletonMap("k", value), now.get() + lifetime.toMillis()); } private static Duration hours(int hours) { diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java index 7812a12b4e1f..94f05c7f9a0f 100644 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/FakeRESTAndOSSServer.java @@ -44,6 +44,7 @@ public class FakeRESTAndOSSServer implements AutoCloseable { private final HttpServer server; private final List tokens = new ArrayList<>(); private final AtomicInteger tokenRequests = new AtomicInteger(); + private final List tokenRequestPaths = Collections.synchronizedList(new ArrayList<>()); private final List ossGetAuthorizations = Collections.synchronizedList(new ArrayList<>()); @@ -82,6 +83,10 @@ public int tokenRequests() { return tokenRequests.get(); } + public List tokenRequestPaths() { + return tokenRequestPaths; + } + public List ossGetAuthorizations() { return ossGetAuthorizations; } @@ -96,6 +101,7 @@ public void close() { } private void handleToken(HttpExchange exchange) throws IOException { + tokenRequestPaths.add(exchange.getRequestURI().getPath()); int index = tokenRequests.getAndIncrement(); GetTableTokenResponse token = tokens.get(Math.min(index, tokens.size() - 1)); byte[] body = RESTApi.toJson(token).getBytes(StandardCharsets.UTF_8); diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java index 32c0ec1ad256..47e0b646f944 100644 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/oss/RESTTokenCredentialsProviderTest.java @@ -59,10 +59,10 @@ public void testUsesTheTokenInTheOptionsWhileItHasTimeLeft() { } @Test - public void testLoadsANewTokenFromTheCatalogWhenTheCurrentOneIsAboutToExpire() { + public void testLoadsANewTokenFromTheCatalogOnceTheCurrentOneExpired() { server.addToken("ak-2", System.currentTimeMillis() + Duration.ofHours(4).toMillis()); RESTTokenCredentialsProvider provider = - provider(tableOptions("ak-1", Duration.ofMinutes(30))); + provider(tableOptions("ak-1", Duration.ofMinutes(-1))); assertThat(provider.getCredentials().getAccessKeyId()).isEqualTo("ak-2"); int requests = server.tokenRequests(); diff --git a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java index 62867268c9a7..4df50f4ae028 100644 --- a/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java +++ b/paimon-filesystems/paimon-oss-impl/src/test/java/org/apache/paimon/rest/RESTTokenFileIOOnOSSTest.java @@ -47,16 +47,7 @@ public void testOpenStreamPicksUpRefreshedTokenAfterCacheEviction() throws Excep long now = System.currentTimeMillis(); server.addToken("ak-1", now + Duration.ofHours(2).toMillis()); server.addToken("ak-2", now + Duration.ofHours(4).toMillis()); - Options options = server.catalogOptions(); - options.set("fs.oss.multipart.download.size", "1024"); - options.set("fs.oss.multipart.download.threads", "1"); - RESTTokenFileIO fileIO = - new RESTTokenFileIO( - CatalogContext.create( - options, new Configuration(false), new PluginOSSLoader(), null), - null, - Identifier.create("db", "table"), - new Path("oss://bucket/table")); + RESTTokenFileIO fileIO = restTokenFileIO(server, "table"); try (SeekableInputStream in = fileIO.newInputStream(new Path("oss://bucket/table/data"))) { @@ -78,6 +69,42 @@ options, new Configuration(false), new PluginOSSLoader(), null), } } + /** A stream reloads the token of its own table, even when another table got the same one. */ + @Test + public void testStreamReloadsTheTokenOfItsOwnTable() throws Exception { + try (FakeRESTAndOSSServer server = new FakeRESTAndOSSServer()) { + server.addToken("ak-1", System.currentTimeMillis() + Duration.ofHours(2).toMillis()); + RESTTokenFileIO tableA = restTokenFileIO(server, "table_a"); + RESTTokenFileIO tableB = restTokenFileIO(server, "table_b"); + + tableA.exists(new Path("oss://bucket/table_a/data")); + try (SeekableInputStream in = + tableB.newInputStream(new Path("oss://bucket/table_b/data"))) { + in.read(new byte[1024]); + RESTTokenRefresher.clockOffsetMillis = Duration.ofMinutes(90).toMillis(); + in.seek(3000); + in.read(new byte[1024]); + } finally { + RESTTokenRefresher.clockOffsetMillis = 0; + } + + List paths = server.tokenRequestPaths(); + assertThat(paths.get(paths.size() - 1)).endsWith("/tables/table_b/token"); + } + } + + private static RESTTokenFileIO restTokenFileIO(FakeRESTAndOSSServer server, String table) { + Options options = server.catalogOptions(); + options.set("fs.oss.multipart.download.size", "1024"); + options.set("fs.oss.multipart.download.threads", "1"); + return new RESTTokenFileIO( + CatalogContext.create( + options, new Configuration(false), new PluginOSSLoader(), null), + null, + Identifier.create("db", table), + new Path("oss://bucket/" + table)); + } + /** Wraps {@link OSSFileIO} the way the OSS plugin does, including its no-op close. */ private static class PluginOSSLoader implements FileIOLoader {