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..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
@@ -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);
}
/**
@@ -335,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 "
@@ -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..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
@@ -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,53 @@ 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) {
+ logAcquireInitialToken(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 refreshResult;
+ try {
+ 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 -> {
+ // 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);
+ });
}
@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();
+ }
+ }
+}