From 7556f5b14848932e54f5ae6e078b77e9b48aa539 Mon Sep 17 00:00:00 2001 From: Dongie Agnir Date: Wed, 15 Jul 2026 09:54:09 -0700 Subject: [PATCH 1/2] Wire in acquireAsync in adaptive strat Implement `DefaultAdaptiveRetryStrategy`'s `acquireInitialTokenAsync` and `refreshRetryTokenAsync` using the new `RateLimiterTokenBucket#acquireAsync`. The existing `acquireInitialToken` and `refreshRetryToken` delegate to these new methods. --- .../amazon/awssdk/spotbugs-suppressions.xml | 10 + .../retries/internal/BaseRetryStrategy.java | 28 ++- .../DefaultAdaptiveRetryStrategy.java | 50 +++- .../RateLimiterTokenBucketStore.java | 46 +++- .../RateLimiterTokenBucketStoreTest.java | 22 +- .../codegen-resources/json/service-2.json | 38 +++ .../retry/AdaptiveRetryRateLimitingTest.java | 236 ++++++++++++++++++ 7 files changed, 406 insertions(+), 24 deletions(-) create mode 100644 test/codegen-generated-classes-test/src/main/resources/codegen-resources/json/service-2.json create mode 100644 test/codegen-generated-classes-test/src/test/java/software/amazon/awssdk/services/retry/AdaptiveRetryRateLimitingTest.java diff --git a/build-tools/src/main/resources/software/amazon/awssdk/spotbugs-suppressions.xml b/build-tools/src/main/resources/software/amazon/awssdk/spotbugs-suppressions.xml index 05606dc6d574..152507019408 100644 --- a/build-tools/src/main/resources/software/amazon/awssdk/spotbugs-suppressions.xml +++ b/build-tools/src/main/resources/software/amazon/awssdk/spotbugs-suppressions.xml @@ -547,4 +547,14 @@ + + + + + + + + + + diff --git a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java index dfda06093ab3..7db9e3c598d4 100644 --- a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java +++ b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java @@ -39,6 +39,7 @@ import software.amazon.awssdk.retries.internal.circuitbreaker.TokenBucket; import software.amazon.awssdk.retries.internal.circuitbreaker.TokenBucketStore; import software.amazon.awssdk.utils.Logger; +import software.amazon.awssdk.utils.Pair; import software.amazon.awssdk.utils.ToString; import software.amazon.awssdk.utils.Validate; @@ -86,7 +87,7 @@ public abstract class BaseRetryStrategy implements DefaultAwareRetryStrategy { * @see RetryStrategy#acquireInitialToken(AcquireInitialTokenRequest) */ @Override - public final AcquireInitialTokenResponse acquireInitialToken(AcquireInitialTokenRequest request) { + public AcquireInitialTokenResponse acquireInitialToken(AcquireInitialTokenRequest request) { logAcquireInitialToken(request); DefaultRetryToken token = DefaultRetryToken.builder().scope(request.scope()).build(); return AcquireInitialTokenResponse.create(token, computeInitialBackoff(request)); @@ -98,7 +99,20 @@ public final AcquireInitialTokenResponse acquireInitialToken(AcquireInitialToken * @see RetryStrategy#refreshRetryToken(RefreshRetryTokenRequest) */ @Override - public final RefreshRetryTokenResponse refreshRetryToken(RefreshRetryTokenRequest request) { + public RefreshRetryTokenResponse refreshRetryToken(RefreshRetryTokenRequest request) { + Pair refreshedToken = refreshTokenOrThrow(request); + Duration backoff = computeBackoff(request, refreshedToken.left()); + + logRefreshTokenSuccess(refreshedToken.left(), refreshedToken.right(), backoff); + return RefreshRetryTokenResponseImpl.create(refreshedToken.left(), backoff); + } + + /** + * Attempt to refresh the token for a retry or throws {@link TokenAcquisitionFailedException} if unable to do so. + * + * @return A pair of the refreshed token and the successful acquire response from the token bucket. + */ + protected Pair refreshTokenOrThrow(RefreshRetryTokenRequest request) { DefaultRetryToken token = asDefaultRetryToken(request.token()); // Check if we meet the preconditions needed for retrying. These will throw if the expected condition is not meet. @@ -115,12 +129,8 @@ public final RefreshRetryTokenResponse refreshRetryToken(RefreshRetryTokenReques // All the conditions required to retry were meet, update the internal state before retrying. updateStateForRetry(request); - // Refresh the retry token and compute the backoff delay. - DefaultRetryToken refreshedToken = refreshToken(request, acquireResponse); - Duration backoff = computeBackoff(request, refreshedToken); - - logRefreshTokenSuccess(refreshedToken, acquireResponse, backoff); - return RefreshRetryTokenResponseImpl.create(refreshedToken, backoff); + // Refresh the retry token + return Pair.of(refreshToken(request, acquireResponse), acquireResponse); } /** @@ -343,7 +353,7 @@ private void logAcquireInitialToken(AcquireInitialTokenRequest request) { tokenBucket.currentCapacity(), tokenBucket.maxCapacity())); } - private void logRefreshTokenSuccess(DefaultRetryToken token, AcquireResponse acquireResponse, Duration delay) { + protected void logRefreshTokenSuccess(DefaultRetryToken token, AcquireResponse acquireResponse, Duration delay) { log.debug(() -> String.format("Request attempt %d token acquired " + "(backoff: %dms, cost: %d, capacity: %d/%d)", token.attempt(), delay.toMillis(), diff --git a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java index e63e683a18ab..f5fa5f6345c2 100644 --- a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java +++ b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java @@ -16,22 +16,27 @@ package software.amazon.awssdk.retries.internal; import java.time.Duration; +import java.util.concurrent.CompletableFuture; import java.util.function.Predicate; import software.amazon.awssdk.annotations.SdkInternalApi; import software.amazon.awssdk.retries.AdaptiveRetryStrategy; import software.amazon.awssdk.retries.api.AcquireInitialTokenRequest; +import software.amazon.awssdk.retries.api.AcquireInitialTokenResponse; import software.amazon.awssdk.retries.api.BackoffStrategy; import software.amazon.awssdk.retries.api.RefreshRetryTokenRequest; +import software.amazon.awssdk.retries.api.RefreshRetryTokenResponse; +import software.amazon.awssdk.retries.internal.circuitbreaker.AcquireResponse; import software.amazon.awssdk.retries.internal.circuitbreaker.TokenBucketStore; import software.amazon.awssdk.retries.internal.ratelimiter.RateLimiterTokenBucket; import software.amazon.awssdk.retries.internal.ratelimiter.RateLimiterTokenBucketStore; +import software.amazon.awssdk.utils.CompletableFutureUtils; import software.amazon.awssdk.utils.Logger; +import software.amazon.awssdk.utils.Pair; import software.amazon.awssdk.utils.Validate; @SdkInternalApi public final class DefaultAdaptiveRetryStrategy extends BaseRetryStrategy implements AdaptiveRetryStrategy { - private static final Logger LOG = Logger.loggerFor(DefaultAdaptiveRetryStrategy.class); private final RateLimiterTokenBucketStore rateLimiterTokenBucketStore; @@ -42,13 +47,48 @@ public final class DefaultAdaptiveRetryStrategy } @Override - protected Duration computeInitialBackoff(AcquireInitialTokenRequest request) { - throw new UnsupportedOperationException("TODO"); + public AcquireInitialTokenResponse acquireInitialToken(AcquireInitialTokenRequest request) { + return CompletableFutureUtils.joinLikeSync(acquireInitialTokenAsync(request)); + } + + @Override + public RefreshRetryTokenResponse refreshRetryToken(RefreshRetryTokenRequest request) { + return CompletableFutureUtils.joinLikeSync(refreshRetryTokenAsync(request)); + } + + @Override + public CompletableFuture acquireInitialTokenAsync(AcquireInitialTokenRequest request) { + RateLimiterTokenBucket bucket = rateLimiterTokenBucketStore.tokenBucketForScope(request.scope()); + CompletableFuture acquireResult = bucket.acquireAsync(); + + return acquireResult.thenApply(r -> { + DefaultRetryToken token = DefaultRetryToken.builder().scope(request.scope()).build(); + return AcquireInitialTokenResponse.create(token, Duration.ZERO); + }); + } + + @Override + public CompletableFuture refreshRetryTokenAsync(RefreshRetryTokenRequest request) { + DefaultRetryToken token = (DefaultRetryToken) request.token(); + Pair refreshedToken; + try { + refreshedToken = refreshTokenOrThrow(request); + } catch (Throwable t) { + return CompletableFutureUtils.failedFuture(t); + } + + RateLimiterTokenBucket bucket = rateLimiterTokenBucketStore.tokenBucketForScope(token.scope()); + CompletableFuture acquireResult = bucket.acquireAsync(); + return acquireResult.thenApply(r -> { + Duration backoff = computeBackoff(request, refreshedToken.left()); + logRefreshTokenSuccess(refreshedToken.left(), refreshedToken.right(), backoff); + return RefreshRetryTokenResponse.create(refreshedToken.left(), backoff); + }); } @Override - protected Duration computeBackoff(RefreshRetryTokenRequest request, DefaultRetryToken token) { - throw new UnsupportedOperationException("TODO"); + public void close() { + rateLimiterTokenBucketStore.close(); } @Override diff --git a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStore.java b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStore.java index 804649e5ee23..a2a5f2f966ed 100644 --- a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStore.java +++ b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStore.java @@ -15,10 +15,13 @@ package software.amazon.awssdk.retries.internal.ratelimiter; +import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import software.amazon.awssdk.annotations.SdkInternalApi; +import software.amazon.awssdk.annotations.SdkTestInternalApi; import software.amazon.awssdk.annotations.ToBuilderIgnoreField; import software.amazon.awssdk.utils.SdkAutoCloseable; +import software.amazon.awssdk.utils.ThreadFactoryBuilder; import software.amazon.awssdk.utils.Validate; import software.amazon.awssdk.utils.builder.CopyableBuilder; import software.amazon.awssdk.utils.builder.ToCopyableBuilder; @@ -31,16 +34,26 @@ public final class RateLimiterTokenBucketStore implements ToCopyableBuilder, SdkAutoCloseable { private static final int MAX_ENTRIES = 128; + private static final String THREAD_NAME_PREFIX = "sdk-adaptive-rate-limiter-"; + private static final RateLimiterClock DEFAULT_CLOCK = new SystemClock(); private final LruCache scopeToTokenBucket; private final RateLimiterClock clock; private final ScheduledExecutorService scheduler; + private final boolean closeScheduler; private RateLimiterTokenBucketStore(Builder builder) { - this.clock = Validate.paramNotNull(builder.clock, "clock"); - this.scheduler = Validate.paramNotNull(builder.scheduler, "scheduler"); + this(builder.clock, + resolveScheduler(builder), + builder.scheduler == null); + } + + private RateLimiterTokenBucketStore(RateLimiterClock clock, ScheduledExecutorService scheduler, boolean closeScheduler) { + this.clock = Validate.paramNotNull(clock, "clock"); + this.scheduler = Validate.paramNotNull(scheduler, "scheduler"); + this.closeScheduler = closeScheduler; this.scopeToTokenBucket = LruCache.builder( - x -> new RateLimiterTokenBucket(clock, scheduler)) + x -> new RateLimiterTokenBucket(clock, scheduler)) .maxSize(MAX_ENTRIES) .build(); } @@ -48,13 +61,30 @@ private RateLimiterTokenBucketStore(Builder builder) { @Override public void close() { scopeToTokenBucket.evictAll(); - scheduler.shutdownNow(); + if (closeScheduler) { + scheduler.shutdownNow(); + } } public RateLimiterTokenBucket tokenBucketForScope(String scope) { return scopeToTokenBucket.get(scope); } + @SdkTestInternalApi + ScheduledExecutorService scheduler() { + return scheduler; + } + + private static ScheduledExecutorService resolveScheduler(Builder b) { + if (b.scheduler != null) { + return b.scheduler; + } + return Executors.newSingleThreadScheduledExecutor(new ThreadFactoryBuilder() + .daemonThreads(true) + .threadNamePrefix(THREAD_NAME_PREFIX) + .build()); + } + @Override @ToBuilderIgnoreField("scopeToTokenBucket") public Builder toBuilder() { @@ -83,7 +113,13 @@ public Builder clock(RateLimiterClock clock) { return this; } - public Builder executor(ScheduledExecutorService scheduler) { + /** + * The scheduler used by the {@link RateLimiterTokenBucket rate limter buckets} to perform async notifications. + * The configured scheduler will not be closed when {@link #close() closing} this bucket store. + * + * @return This object for method chaining. + */ + public Builder scheduler(ScheduledExecutorService scheduler) { this.scheduler = scheduler; return this; } diff --git a/core/retries/src/test/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStoreTest.java b/core/retries/src/test/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStoreTest.java index 46d5cf1db13d..b915045e1b5a 100644 --- a/core/retries/src/test/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStoreTest.java +++ b/core/retries/src/test/java/software/amazon/awssdk/retries/internal/ratelimiter/RateLimiterTokenBucketStoreTest.java @@ -18,27 +18,39 @@ 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.verify; import java.util.ArrayList; import java.util.List; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ExecutorService; import java.util.concurrent.ScheduledExecutorService; import org.junit.jupiter.api.Test; public class RateLimiterTokenBucketStoreTest { @Test - void close_closesScheduler() { + void close_schedulerProvided_schedulerNotClosed() { ScheduledExecutorService scheduler = mock(ScheduledExecutorService.class); RateLimiterTokenBucketStore store = RateLimiterTokenBucketStore.builder() .clock(new SystemClock()) - .executor(scheduler) + .scheduler(scheduler) .build(); store.close(); - verify(scheduler).shutdownNow(); + verify(scheduler, never()).shutdownNow(); + verify(scheduler, never()).shutdown(); + } + + @Test + void close_schedulerNotProvidedOnBuilder_schedulerClosed() { + RateLimiterTokenBucketStore store = RateLimiterTokenBucketStore.builder() + .clock(new SystemClock()) + .build(); + + store.close(); + + assertThat(store.scheduler().isShutdown()).isTrue(); } @Test @@ -49,7 +61,7 @@ void close_closesAllCacheEntries() { ScheduledExecutorService scheduler = mock(ScheduledExecutorService.class); RateLimiterTokenBucketStore store = RateLimiterTokenBucketStore.builder() .clock(new SystemClock()) - .executor(scheduler) + .scheduler(scheduler) .build(); List> futures = new ArrayList<>(entries); diff --git a/test/codegen-generated-classes-test/src/main/resources/codegen-resources/json/service-2.json b/test/codegen-generated-classes-test/src/main/resources/codegen-resources/json/service-2.json new file mode 100644 index 000000000000..c0372d42aa46 --- /dev/null +++ b/test/codegen-generated-classes-test/src/main/resources/codegen-resources/json/service-2.json @@ -0,0 +1,38 @@ +{ + "version":"2.0", + "metadata":{ + "apiVersion":"2010-05-08", + "endpointPrefix":"json-service-endpoint", + "globalEndpoint": "json-service.amazonaws.com", + "jsonVersion":"1.0", + "protocol":"json", + "serviceAbbreviation":"Aws Json Service", + "serviceFullName":"Some Service That Uses AWS JSON", + "serviceId":"Aws Json Service", + "signingName": "aws-json-service", + "signatureVersion":"v4", + "uid":"aws-json-service-2010-05-08", + "awsQueryCompatible":{} + }, + "operations":{ + "AllType": { + "name": "APostOperation", + "http": { + "method": "POST", + "requestUri": "/" + }, + "httpChecksumRequired": true + } + }, + "shapes": { + "OneShape": { + "type": "structure", + "members": { + "StringMember": { + "shape": "String" + } + } + }, + "String":{"type":"string"} + } +} diff --git a/test/codegen-generated-classes-test/src/test/java/software/amazon/awssdk/services/retry/AdaptiveRetryRateLimitingTest.java b/test/codegen-generated-classes-test/src/test/java/software/amazon/awssdk/services/retry/AdaptiveRetryRateLimitingTest.java new file mode 100644 index 000000000000..39ded28482b9 --- /dev/null +++ b/test/codegen-generated-classes-test/src/test/java/software/amazon/awssdk/services/retry/AdaptiveRetryRateLimitingTest.java @@ -0,0 +1,236 @@ +/* + * Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. + * + * Licensed under the Apache License, Version 2.0 (the "License"). + * You may not use this file except in compliance with the License. + * A copy of the License is located at + * + * http://aws.amazon.com/apache2.0 + * + * or in the "license" file accompanying this file. This file 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 software.amazon.awssdk.services.retry; + +import static com.github.tomakehurst.wiremock.client.WireMock.aResponse; +import static com.github.tomakehurst.wiremock.client.WireMock.any; +import static com.github.tomakehurst.wiremock.client.WireMock.anyUrl; +import static com.github.tomakehurst.wiremock.core.WireMockConfiguration.options; +import static org.assertj.core.api.Assertions.assertThat; + +import com.github.tomakehurst.wiremock.common.FileSource; +import com.github.tomakehurst.wiremock.extension.Parameters; +import com.github.tomakehurst.wiremock.extension.ResponseDefinitionTransformer; +import com.github.tomakehurst.wiremock.http.Request; +import com.github.tomakehurst.wiremock.http.ResponseDefinition; +import com.github.tomakehurst.wiremock.junit5.WireMockExtension; +import com.github.tomakehurst.wiremock.stubbing.ServeEvent; +import java.net.URI; +import java.time.Duration; +import java.time.Instant; +import java.util.List; +import java.util.concurrent.ConcurrentLinkedDeque; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.stream.Collectors; +import java.util.stream.IntStream; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.RegisterExtension; +import software.amazon.awssdk.auth.credentials.AwsBasicCredentials; +import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider; +import software.amazon.awssdk.awscore.retry.AwsRetryStrategy; +import software.amazon.awssdk.core.SdkSystemSetting; +import software.amazon.awssdk.regions.Region; +import software.amazon.awssdk.retries.AdaptiveRetryStrategy; +import software.amazon.awssdk.services.json.JsonClient; + +/** + * End-to-end testing for testing the correctness of the rate limiting implementation of the adaptive retry strategy. + */ +class AdaptiveRetryRateLimitingTest { + private static final AtomicLong eventCount = new AtomicLong(0L); + + private static final int throttleTps = 20; + private static final double tpsLower = throttleTps * 0.5; + private static final double tpsUpper = throttleTps * 1.5; + + private static final Duration testDuration = Duration.ofSeconds(30); + private static final Duration warmupDuration = Duration.ofSeconds(2); + private static final Duration tailTrim = Duration.ofSeconds(1); + private static final int numWorkerThreads = Runtime.getRuntime().availableProcessors(); + + private static final RollingThrottle throttler = new RollingThrottle(throttleTps); + + private static String newRetriesPropertySave; + + @RegisterExtension + static WireMockExtension wm = WireMockExtension.newInstance() + .options(options().dynamicPort() + .extensions(new ThrottlingTransformer())) + .build(); + + @BeforeAll + static void setup() { + newRetriesPropertySave = System.getProperty(SdkSystemSetting.AWS_NEW_RETRIES_2026.property()); + System.setProperty(SdkSystemSetting.AWS_NEW_RETRIES_2026.property(), "true"); + } + + @AfterAll + static void teardown() { + if (newRetriesPropertySave != null) { + System.setProperty(SdkSystemSetting.AWS_NEW_RETRIES_2026.property(), newRetriesPropertySave); + } else { + System.clearProperty(SdkSystemSetting.AWS_NEW_RETRIES_2026.property()); + } + } + + /** + * The server has a static max TPS of 20 before responding with throttling exceptions. The client sending rate should match + * this within an acceptable margin. + */ + @Test + void staticServerThrottling() throws Exception { + + wm.stubFor(any(anyUrl()).willReturn(aResponse().withTransformers("throttling"))); + + AdaptiveRetryStrategy adaptiveRetryStrategy = AwsRetryStrategy.adaptiveRetryStrategy(true); + + JsonClient client = JsonClient.builder() + .endpointOverride(URI.create(wm.baseUrl())) + .region(Region.US_WEST_2) + .credentialsProvider(StaticCredentialsProvider.create( + AwsBasicCredentials.create("akid", "secret"))) + .overrideConfiguration(o -> o.retryStrategy(adaptiveRetryStrategy)) + .build(); + + Instant startT = Instant.now(); + Instant endAt = startT.plus(testDuration); + AtomicBoolean stop = new AtomicBoolean(false); + + ExecutorService pool = Executors.newFixedThreadPool(numWorkerThreads); + try { + List> futures = IntStream.range(0, numWorkerThreads) + .mapToObj(i -> pool.submit(() -> { + while (!stop.get() && Instant.now().isBefore(endAt)) { + try { + client.allType(r -> { + }); + } catch (Exception e) { + // We measure wire-level TPS, not call success. + // Throttling exceptions surface here once retries are exhausted. + } + } + })) + .collect(Collectors.toList()); + + for (Future f : futures) { + f.get(testDuration.getSeconds() + 30, TimeUnit.SECONDS); + } + } finally { + stop.set(true); + } + Instant endT = Instant.now(); + + double measuredTps = tpsInWindow( + wm.getAllServeEvents(), + startT.plus(warmupDuration), + endT.minus(tailTrim)); + + assertThat(measuredTps) + .as("Expected send rate >= %.1f TPS, but got %.1f TPS. " + + "Is the SDK not converging to the server throttle rate?", + tpsLower, measuredTps) + .isGreaterThanOrEqualTo(tpsLower); + assertThat(measuredTps) + .as("Expected send rate <= %.1f TPS, but got %.1f TPS. " + + "Is the SDK not respecting the server throttle rate limit?", + tpsUpper, measuredTps) + .isLessThanOrEqualTo(tpsUpper); + } + + private static double tpsInWindow(List events, Instant from, Instant to) { + long count = events.stream() + .map(e -> e.getRequest().getLoggedDate().toInstant()) + .filter(t -> !t.isBefore(from) && t.isBefore(to)) + .count(); + double seconds = Duration.between(from, to).toMillis() / 1000.0; + return seconds <= 0 ? 0 : count / seconds; + } + + /** + * Sliding 1-second window admission control. Once {@code TPS} requests have been admitted within the past second, returns + * false until the window slides. + */ + static final class RollingThrottle { + private final int tps; + private final ConcurrentLinkedDeque hits = new ConcurrentLinkedDeque<>(); + + RollingThrottle(int tps) { + this.tps = tps; + } + + synchronized boolean tryAcquire() { + long now = System.nanoTime(); + long cutoff = now - TimeUnit.SECONDS.toNanos(1); + while (!hits.isEmpty() && hits.peekFirst() < cutoff) { + hits.pollFirst(); + } + if (hits.size() >= tps) { + return false; + } + hits.addLast(now); + return true; + } + } + + /** + * WireMock transformer that delegates the 200-vs-throttle decision to {@link #throttler}. + */ + public static final class ThrottlingTransformer extends ResponseDefinitionTransformer { + + private static final String successBody = "{}"; + private static final String throttleBody = + "{\"__type\":\"ThrottlingException\",\"message\":\"Rate exceeded\"}"; + + @Override + public String getName() { + return "throttling"; + } + + @Override + public boolean applyGlobally() { + return false; + } + + @Override + public ResponseDefinition transform(Request request, + ResponseDefinition responseDefinition, + FileSource files, + Parameters parameters) { + eventCount.incrementAndGet(); + + if (throttler.tryAcquire()) { + return aResponse() + .withStatus(200) + .withHeader("Content-Type", "application/x-amz-json-1.0") + .withBody(successBody) + .build(); + } + return aResponse() + .withStatus(400) + .withHeader("Content-Type", "application/x-amz-json-1.0") + .withHeader("x-amzn-ErrorType", "ThrottlingException") + .withBody(throttleBody) + .build(); + } + } +} From 0171182cfdfaf765d73d4afd864e21ad858a4563 Mon Sep 17 00:00:00 2001 From: Dongie Agnir Date: Tue, 21 Jul 2026 14:07:21 -0700 Subject: [PATCH 2/2] Review comments --- .../retries/internal/BaseRetryStrategy.java | 2 +- .../internal/DefaultAdaptiveRetryStrategy.java | 15 ++++++++++----- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java index 7db9e3c598d4..6c6551dc68e2 100644 --- a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java +++ b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/BaseRetryStrategy.java @@ -345,7 +345,7 @@ private String acquisitionFailedMessage(AcquireResponse response) { response.maxCapacity()); } - private void logAcquireInitialToken(AcquireInitialTokenRequest request) { + protected void logAcquireInitialToken(AcquireInitialTokenRequest request) { // Request attempt 1 token acquired (backoff: 0ms, cost: 0, capacity: 500/500) TokenBucket tokenBucket = tokenBucketStore.tokenBucketForScope(request.scope()); log.debug(() -> String.format("Request attempt 1 token acquired " diff --git a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java index f5fa5f6345c2..685bba2d5e46 100644 --- a/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java +++ b/core/retries/src/main/java/software/amazon/awssdk/retries/internal/DefaultAdaptiveRetryStrategy.java @@ -58,6 +58,7 @@ public RefreshRetryTokenResponse refreshRetryToken(RefreshRetryTokenRequest requ @Override public CompletableFuture acquireInitialTokenAsync(AcquireInitialTokenRequest request) { + logAcquireInitialToken(request); RateLimiterTokenBucket bucket = rateLimiterTokenBucketStore.tokenBucketForScope(request.scope()); CompletableFuture acquireResult = bucket.acquireAsync(); @@ -70,19 +71,23 @@ public CompletableFuture acquireInitialTokenAsync(A @Override public CompletableFuture refreshRetryTokenAsync(RefreshRetryTokenRequest request) { DefaultRetryToken token = (DefaultRetryToken) request.token(); - Pair refreshedToken; + Pair refreshResult; try { - refreshedToken = refreshTokenOrThrow(request); + refreshResult = refreshTokenOrThrow(request); } catch (Throwable t) { return CompletableFutureUtils.failedFuture(t); } + DefaultRetryToken refreshedToken = refreshResult.left(); + AcquireResponse acquireResponse = refreshResult.right(); RateLimiterTokenBucket bucket = rateLimiterTokenBucketStore.tokenBucketForScope(token.scope()); CompletableFuture acquireResult = bucket.acquireAsync(); return acquireResult.thenApply(r -> { - Duration backoff = computeBackoff(request, refreshedToken.left()); - logRefreshTokenSuccess(refreshedToken.left(), refreshedToken.right(), backoff); - return RefreshRetryTokenResponse.create(refreshedToken.left(), backoff); + // Note: This is the backoff imposed standard retry strategy, *not* the rate limiter. This must still be honored by + // the caller before sending the request. + Duration backoff = computeBackoff(request, refreshedToken); + logRefreshTokenSuccess(refreshedToken, acquireResponse, backoff); + return RefreshRetryTokenResponse.create(refreshedToken, backoff); }); }