diff --git a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/Messages.java b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/Messages.java index 0a56d6fd46..af0d03d5f0 100644 --- a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/Messages.java +++ b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/Messages.java @@ -4,5 +4,8 @@ public class Messages { // INFO messages public static final String WAITING_MS_BEFORE_RETRYING_WITH_TIMEOUT_OF_MS = "Waiting: {} ms before retrying with timeout of: {} ms"; + public static final String RATE_LIMITED_BY_CC_WAITING_S = "CC returned 429 with Retry-After: {} s. Waiting {} s (capped) before retrying."; + public static final String RATE_LIMITED_BY_CC_NO_HEADER_WAITING_MS = "CC returned 429 without Retry-After header. Waiting {} ms before retrying."; + public static final String RANDOM_WAIT_BEFORE_RETRY_MS = "Waiting {} ms (randomized) before retrying failed CC operation."; } diff --git a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationException.java b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationException.java index 2f75d0ce85..6a2b2b60b6 100644 --- a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationException.java +++ b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationException.java @@ -8,6 +8,7 @@ public class CloudOperationException extends CloudException { private final HttpStatus statusCode; private final String statusText; private final String description; + private final Long retryAfterSeconds; public CloudOperationException(HttpStatus statusCode) { this(statusCode, statusCode.getReasonPhrase()); @@ -22,10 +23,16 @@ public CloudOperationException(HttpStatus statusCode, String statusText, String } public CloudOperationException(HttpStatus statusCode, String statusText, String description, Throwable cause) { + this(statusCode, statusText, description, cause, null); + } + + public CloudOperationException(HttpStatus statusCode, String statusText, String description, Throwable cause, + Long retryAfterSeconds) { super(getExceptionMessage(statusCode, statusText, description), cause); this.statusCode = statusCode; this.statusText = statusText; this.description = description; + this.retryAfterSeconds = retryAfterSeconds; } private static String getExceptionMessage(HttpStatus statusCode, String statusText, String description) { @@ -47,4 +54,8 @@ public String getDescription() { return description; } + public Long getRetryAfterSeconds() { + return retryAfterSeconds; + } + } diff --git a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandler.java b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandler.java index 006e6051f4..8b439225f5 100644 --- a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandler.java +++ b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandler.java @@ -1,24 +1,26 @@ package org.cloudfoundry.multiapps.controller.client.facade.rest; +import java.io.IOException; +import java.util.List; +import java.util.Map; +import java.util.stream.Collectors; + import com.fasterxml.jackson.databind.ObjectMapper; import org.cloudfoundry.multiapps.controller.client.facade.CloudOperationException; import org.cloudfoundry.multiapps.controller.client.facade.util.CloudUtil; +import org.springframework.http.HttpHeaders; import org.springframework.http.HttpStatus; import org.springframework.http.client.ClientHttpResponse; import org.springframework.web.client.DefaultResponseErrorHandler; import org.springframework.web.client.RestClientException; -import java.io.IOException; -import java.util.List; -import java.util.Map; -import java.util.stream.Collectors; - public class CloudControllerResponseErrorHandler extends DefaultResponseErrorHandler { private static CloudOperationException getException(ClientHttpResponse response) throws IOException { HttpStatus statusCode = HttpStatus.valueOf(response.getStatusCode() .value()); String statusText = response.getStatusText(); + Long retryAfterSeconds = extractRetryAfterSeconds(response, statusCode); ObjectMapper mapper = new ObjectMapper(); // can reuse, share globally @@ -26,12 +28,29 @@ private static CloudOperationException getException(ClientHttpResponse response) try { @SuppressWarnings("unchecked") Map responseBody = mapper.readValue(response.getBody(), Map.class); String description = getTrimmedDescription(responseBody); - return new CloudOperationException(statusCode, statusText, description); + return new CloudOperationException(statusCode, statusText, description, null, retryAfterSeconds); } catch (IOException e) { // Fall through. Handled below. } } - return new CloudOperationException(statusCode, statusText); + return new CloudOperationException(statusCode, statusText, null, null, retryAfterSeconds); + } + + private static Long extractRetryAfterSeconds(ClientHttpResponse response, HttpStatus statusCode) { + if (statusCode != HttpStatus.TOO_MANY_REQUESTS) { + return null; + } + String headerValue = response.getHeaders() + .getFirst(HttpHeaders.RETRY_AFTER); + if (headerValue == null) { + return null; + } + try { + long parsed = Long.parseLong(headerValue); + return parsed >= 0 ? parsed : null; + } catch (NumberFormatException _) { + return null; + } } private static String getTrimmedDescription(Map responseBody) { diff --git a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutor.java b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutor.java index 54bae27e32..fcb9e7d9e3 100644 --- a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutor.java +++ b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutor.java @@ -7,7 +7,11 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import java.util.concurrent.ThreadLocalRandom; import java.util.function.Function; +import java.util.function.LongConsumer; +import java.util.function.LongSupplier; +import java.util.function.Supplier; import org.apache.commons.collections4.SetUtils; import org.cloudfoundry.multiapps.common.util.MiscUtil; @@ -23,7 +27,8 @@ public class ResilientCloudOperationExecutor extends ResilientOperationExecutor private static final Set DEFAULT_STATUSES_TO_IGNORE = Set.of(HttpStatus.GATEWAY_TIMEOUT, HttpStatus.REQUEST_TIMEOUT, HttpStatus.INTERNAL_SERVER_ERROR, HttpStatus.BAD_GATEWAY, - HttpStatus.SERVICE_UNAVAILABLE); + HttpStatus.SERVICE_UNAVAILABLE, + HttpStatus.TOO_MANY_REQUESTS); private static final int DEFAULT_TIMEOUT_RETRY_WAIT_TIME_IN_MILLIS = 30 * 1000; // 30 seconds @@ -31,7 +36,16 @@ public class ResilientCloudOperationExecutor extends ResilientOperationExecutor Duration.ofMinutes(8), 3, Duration.ofMinutes(15)); + private static final long RATE_LIMIT_RETRY_AFTER_CAP_IN_SECONDS = 120L; + private static final long RATE_LIMIT_FALLBACK_WAIT_IN_MILLIS = 60_000L; + private static final long RANDOM_RETRY_MIN_WAIT_IN_MILLIS = 30_000L; + private static final long RANDOM_RETRY_MAX_WAIT_IN_MILLIS = 90_000L; + private Set additionalStatusesToIgnore = Collections.emptySet(); + private LongConsumer sleeper = MiscUtil::sleep; + private LongSupplier randomDelaySupplier = () -> ThreadLocalRandom.current() + .nextLong(RANDOM_RETRY_MIN_WAIT_IN_MILLIS, + RANDOM_RETRY_MAX_WAIT_IN_MILLIS + 1); @Override public ResilientCloudOperationExecutor withRetryCount(long retryCount) { @@ -48,6 +62,30 @@ public ResilientCloudOperationExecutor withStatusesToIgnore(HttpStatus... status return this; } + ResilientCloudOperationExecutor withSleeper(LongConsumer sleeper) { + this.sleeper = sleeper; + return this; + } + + ResilientCloudOperationExecutor withRandomDelaySupplier(LongSupplier randomDelaySupplier) { + this.randomDelaySupplier = randomDelaySupplier; + return this; + } + + @Override + public T execute(Supplier operation) { + for (long i = 1; i < retryCount; i++) { + try { + return operation.get(); + } catch (RuntimeException e) { + handle(e); + long waitMillis = computeWaitMillis(e, i); + sleeper.accept(waitMillis); + } + } + return operation.get(); + } + public T executeWithExponentialBackoff(Function operation) { int waitTimeBetweenRetriesInMillis = DEFAULT_TIMEOUT_RETRY_WAIT_TIME_IN_MILLIS; int retryIndex = 1; @@ -87,6 +125,34 @@ protected void handle(CloudOperationException e) { e); } + private long computeWaitMillis(RuntimeException e, long attemptIndex) { + if (isRateLimitException(e)) { + return computeRateLimitWaitMillis((CloudOperationException) e); + } + if (attemptIndex >= 2) { + long waitMillis = randomDelaySupplier.getAsLong(); + LOGGER.info(Messages.RANDOM_WAIT_BEFORE_RETRY_MS, waitMillis); + return waitMillis; + } + return waitTimeBetweenRetriesInMillis; + } + + private boolean isRateLimitException(RuntimeException e) { + return e instanceof CloudOperationException cloudException + && cloudException.getStatusCode() == HttpStatus.TOO_MANY_REQUESTS; + } + + private long computeRateLimitWaitMillis(CloudOperationException e) { + Long retryAfterSeconds = e.getRetryAfterSeconds(); + if (retryAfterSeconds == null) { + LOGGER.info(Messages.RATE_LIMITED_BY_CC_NO_HEADER_WAITING_MS, RATE_LIMIT_FALLBACK_WAIT_IN_MILLIS); + return RATE_LIMIT_FALLBACK_WAIT_IN_MILLIS; + } + long cappedSeconds = Math.clamp(retryAfterSeconds, 1L, RATE_LIMIT_RETRY_AFTER_CAP_IN_SECONDS); + LOGGER.info(Messages.RATE_LIMITED_BY_CC_WAITING_S, retryAfterSeconds, cappedSeconds); + return cappedSeconds * 1000L; + } + private boolean shouldRetry(CloudOperationException e) { return getStatusesToIgnore().contains(e.getStatusCode()); } diff --git a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientOperationExecutor.java b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientOperationExecutor.java index bc536a45cc..2adb49cbb3 100644 --- a/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientOperationExecutor.java +++ b/multiapps-controller-client/src/main/java/org/cloudfoundry/multiapps/controller/client/util/ResilientOperationExecutor.java @@ -16,7 +16,7 @@ public class ResilientOperationExecutor { private static final long DEFAULT_RETRY_COUNT = 3; private static final long DEFAULT_WAIT_TIME_BETWEEN_RETRIES_IN_MILLIS = 5000; - private long waitTimeBetweenRetriesInMillis = DEFAULT_WAIT_TIME_BETWEEN_RETRIES_IN_MILLIS; + protected long waitTimeBetweenRetriesInMillis = DEFAULT_WAIT_TIME_BETWEEN_RETRIES_IN_MILLIS; protected long retryCount = DEFAULT_RETRY_COUNT; public ResilientOperationExecutor withRetryCount(long retryCount) { diff --git a/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationExceptionTest.java b/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationExceptionTest.java index 8319f74929..65b2faf73d 100644 --- a/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationExceptionTest.java +++ b/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/CloudOperationExceptionTest.java @@ -42,4 +42,32 @@ void testCauseIsPropagatedToSuper() { Assertions.assertSame(cause, e.getCause()); } + + @Test + void testRetryAfterSecondsNullByDefault() { + CloudOperationException e = new CloudOperationException(HttpStatus.TOO_MANY_REQUESTS); + + Assertions.assertNull(e.getRetryAfterSeconds()); + } + + @Test + void testRetryAfterSecondsStoredWhenProvided() { + CloudOperationException e = new CloudOperationException(HttpStatus.TOO_MANY_REQUESTS, + HttpStatus.TOO_MANY_REQUESTS.getReasonPhrase(), + null, null, 60L); + + Assertions.assertEquals(60L, e.getRetryAfterSeconds()); + } + + @Test + void testFourArgConstructorLeavesRetryAfterSecondsNull() { + Throwable cause = new RuntimeException("boom"); + + CloudOperationException e = new CloudOperationException(HttpStatus.TOO_MANY_REQUESTS, + HttpStatus.TOO_MANY_REQUESTS.getReasonPhrase(), "throttled", cause); + + Assertions.assertNull(e.getRetryAfterSeconds()); + Assertions.assertSame(cause, e.getCause()); + } } + diff --git a/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandlerTest.java b/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandlerTest.java index 0fb08bc9cd..ddf3219cf4 100644 --- a/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandlerTest.java +++ b/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/facade/rest/CloudControllerResponseErrorHandlerTest.java @@ -1,27 +1,26 @@ package org.cloudfoundry.multiapps.controller.client.facade.rest; -import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.fail; - import java.io.ByteArrayInputStream; -import java.io.IOException; import java.io.InputStream; import java.nio.charset.StandardCharsets; +import org.cloudfoundry.multiapps.controller.client.facade.CloudOperationException; import org.junit.jupiter.api.Test; import org.springframework.http.HttpHeaders; import org.springframework.http.HttpStatus; import org.springframework.http.client.ClientHttpResponse; import org.springframework.web.client.RestClientException; -import org.cloudfoundry.multiapps.controller.client.facade.CloudOperationException; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; class CloudControllerResponseErrorHandlerTest { private final CloudControllerResponseErrorHandler handler = new CloudControllerResponseErrorHandler(); @Test - void testWithV2Error() throws IOException { + void testWithV2Error() { HttpStatus statusCode = HttpStatus.UNPROCESSABLE_ENTITY; ClientHttpResponseMock response = new ClientHttpResponseMock(statusCode, getClass().getResourceAsStream("v2-error.json")); CloudOperationException expectedException = new CloudOperationException(statusCode, @@ -31,7 +30,7 @@ void testWithV2Error() throws IOException { } @Test - void testWithV3Error() throws IOException { + void testWithV3Error() { HttpStatus statusCode = HttpStatus.BAD_REQUEST; ClientHttpResponseMock response = new ClientHttpResponseMock(statusCode, getClass().getResourceAsStream("v3-error.json")); CloudOperationException expectedException = new CloudOperationException(statusCode, @@ -41,7 +40,7 @@ void testWithV3Error() throws IOException { } @Test - void testWithInvalidError() throws IOException { + void testWithInvalidError() { HttpStatus statusCode = HttpStatus.BAD_REQUEST; ClientHttpResponseMock response = new ClientHttpResponseMock(statusCode, toInputStream("blabla")); CloudOperationException expectedException = new CloudOperationException(statusCode, statusCode.getReasonPhrase()); @@ -49,34 +48,92 @@ void testWithInvalidError() throws IOException { } @Test - void testWithEmptyResponse() throws IOException { + void testWithEmptyResponse() { HttpStatus statusCode = HttpStatus.BAD_REQUEST; ClientHttpResponseMock response = new ClientHttpResponseMock(statusCode, toInputStream("{ }")); CloudOperationException expectedException = new CloudOperationException(statusCode, statusCode.getReasonPhrase()); testWithError(response, expectedException); } - private void testWithError(ClientHttpResponseMock response, CloudOperationException expectedException) throws IOException { - try { - handler.handleError(response); - fail("Expected an exception"); - } catch (CloudOperationException exception) { - assertEquals(expectedException.getStatusCode(), exception.getStatusCode()); - assertEquals(expectedException.getStatusText(), exception.getStatusText()); - assertEquals(expectedException.getDescription(), exception.getDescription()); - } + @Test + void testWith429AndRetryAfterHeader() { + HttpHeaders headers = new HttpHeaders(); + headers.add(HttpHeaders.RETRY_AFTER, "60"); + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.TOO_MANY_REQUESTS, toInputStream("{}"), headers); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(HttpStatus.TOO_MANY_REQUESTS, e.getStatusCode()); + assertEquals(60L, e.getRetryAfterSeconds()); + } + + @Test + void testWith429WithoutRetryAfterHeader() { + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.TOO_MANY_REQUESTS, toInputStream("{}")); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(HttpStatus.TOO_MANY_REQUESTS, e.getStatusCode()); + assertNull(e.getRetryAfterSeconds()); + } + + @Test + void testWith429AndZeroRetryAfterHeaderIsAccepted() { + HttpHeaders headers = new HttpHeaders(); + headers.add(HttpHeaders.RETRY_AFTER, "0"); + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.TOO_MANY_REQUESTS, toInputStream("{}"), headers); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(0L, e.getRetryAfterSeconds()); + } + + @Test + void testWith429AndNegativeRetryAfterHeaderIsIgnored() { + HttpHeaders headers = new HttpHeaders(); + headers.add(HttpHeaders.RETRY_AFTER, "-1"); + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.TOO_MANY_REQUESTS, toInputStream("{}"), headers); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertNull(e.getRetryAfterSeconds()); } @Test - void testWithNonClientOrServerError() throws IOException { + void testWith429AndNonNumericRetryAfterHeader() { + HttpHeaders headers = new HttpHeaders(); + headers.add(HttpHeaders.RETRY_AFTER, "not-a-number"); + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.TOO_MANY_REQUESTS, toInputStream("{}"), headers); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(HttpStatus.TOO_MANY_REQUESTS, e.getStatusCode()); + assertNull(e.getRetryAfterSeconds()); + } + + @Test + void testNon429DoesNotSetRetryAfterSeconds() { + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.INTERNAL_SERVER_ERROR, toInputStream("{}")); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(HttpStatus.INTERNAL_SERVER_ERROR, e.getStatusCode()); + assertNull(e.getRetryAfterSeconds()); + } + + @Test + void testWith429RetainsRetryAfterAndDescriptionWhenBodyPresent() { + HttpHeaders headers = new HttpHeaders(); + headers.add(HttpHeaders.RETRY_AFTER, "30"); + ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.TOO_MANY_REQUESTS, + toInputStream("{\"description\":\"rate limit exceeded\"}"), headers); + CloudOperationException e = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(HttpStatus.TOO_MANY_REQUESTS, e.getStatusCode()); + assertEquals(30L, e.getRetryAfterSeconds()); + assertEquals("rate limit exceeded", e.getDescription()); + } + + private void testWithError(ClientHttpResponseMock response, CloudOperationException expectedException) { + CloudOperationException exception = assertThrows(CloudOperationException.class, () -> handler.handleError(response)); + assertEquals(expectedException.getStatusCode(), exception.getStatusCode()); + assertEquals(expectedException.getStatusText(), exception.getStatusText()); + assertEquals(expectedException.getDescription(), exception.getDescription()); + } + + @Test + void testWithNonClientOrServerError() { HttpStatus statusCode = HttpStatus.PERMANENT_REDIRECT; ClientHttpResponseMock response = new ClientHttpResponseMock(HttpStatus.PERMANENT_REDIRECT, null); - try { - handler.handleError(response); - fail("Expected an exception"); - } catch (RestClientException exception) { - assertEquals("Unknown status code [" + statusCode + "]", exception.getMessage()); - } + RestClientException exception = assertThrows(RestClientException.class, () -> handler.handleError(response)); + assertEquals("Unknown status code [" + statusCode + "]", exception.getMessage()); } private InputStream toInputStream(String string) { @@ -87,10 +144,16 @@ private static class ClientHttpResponseMock implements ClientHttpResponse { private final HttpStatus statusCode; private final InputStream body; + private final HttpHeaders headers; public ClientHttpResponseMock(HttpStatus statusCode, InputStream body) { + this(statusCode, body, new HttpHeaders()); + } + + public ClientHttpResponseMock(HttpStatus statusCode, InputStream body, HttpHeaders headers) { this.statusCode = statusCode; this.body = body; + this.headers = headers; } @Override @@ -100,7 +163,7 @@ public InputStream getBody() { @Override public HttpHeaders getHeaders() { - throw new UnsupportedOperationException(); + return headers; } @Override @@ -126,3 +189,4 @@ public void close() { } } + diff --git a/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutorTest.java b/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutorTest.java index 42ef66f1b9..fd87bbcbf0 100644 --- a/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutorTest.java +++ b/multiapps-controller-client/src/test/java/org/cloudfoundry/multiapps/controller/client/util/ResilientCloudOperationExecutorTest.java @@ -1,27 +1,37 @@ package org.cloudfoundry.multiapps.controller.client.util; +import java.util.ArrayList; +import java.util.List; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Supplier; +import java.util.stream.Stream; import org.cloudfoundry.multiapps.controller.client.facade.CloudOperationException; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; import org.springframework.http.HttpStatus; class ResilientCloudOperationExecutorTest { private ResilientCloudOperationExecutor executor; + private List sleepCalls; + private AtomicInteger attempts; @BeforeEach void setUp() { + sleepCalls = new ArrayList<>(); + attempts = new AtomicInteger(); executor = new ResilientCloudOperationExecutor().withRetryCount(3) - .withWaitTimeBetweenRetriesInMillis(0); + .withWaitTimeBetweenRetriesInMillis(0) + .withSleeper(sleepCalls::add); } @Test void testRetriesOnDefaultIgnoredStatuses() { - AtomicInteger attempts = new AtomicInteger(); Supplier operation = () -> { if (attempts.incrementAndGet() < 2) { throw new CloudOperationException(HttpStatus.BAD_GATEWAY); @@ -37,7 +47,6 @@ void testRetriesOnDefaultIgnoredStatuses() { @Test void testThrowsImmediatelyOnNonIgnoredStatus() { - AtomicInteger attempts = new AtomicInteger(); Supplier operation = () -> { attempts.incrementAndGet(); throw new CloudOperationException(HttpStatus.NOT_FOUND); @@ -50,7 +59,6 @@ void testThrowsImmediatelyOnNonIgnoredStatus() { @Test void testWithStatusesToIgnoreAddsAdditionalRetryableStatuses() { - AtomicInteger attempts = new AtomicInteger(); Supplier operation = () -> { if (attempts.incrementAndGet() < 2) { throw new CloudOperationException(HttpStatus.NOT_FOUND); @@ -73,4 +81,119 @@ void testFluentBuildersReturnSameTypeForChaining() { .withStatusesToIgnore(HttpStatus.I_AM_A_TEAPOT); Assertions.assertNotNull(chained); } + + static Stream retryAfterHeaderCappingCases() { + return Stream.of(Arguments.of("below cap", 2L, 2_000L), + Arguments.of("above cap", 300L, 120_000L), + Arguments.of("at cap", 120L, 120_000L)); + } + + @ParameterizedTest(name = "{0}: retryAfter={1}s → sleep={2}ms") + @MethodSource("retryAfterHeaderCappingCases") + void testRetryAfterHeaderCapping(String description, long retryAfterSeconds, long expectedSleepMillis) { + Supplier operation = () -> { + if (attempts.incrementAndGet() == 1) { + throw new CloudOperationException(HttpStatus.TOO_MANY_REQUESTS, + HttpStatus.TOO_MANY_REQUESTS.getReasonPhrase(), + null, null, retryAfterSeconds); + } + return "ok"; + }; + + executor.execute(operation); + + Assertions.assertEquals(1, sleepCalls.size()); + Assertions.assertEquals(expectedSleepMillis, sleepCalls.get(0)); + } + + @Test + void testRateLimitFallbackWhenRetryAfterAbsent() { + Supplier operation = () -> { + if (attempts.incrementAndGet() == 1) { + throw new CloudOperationException(HttpStatus.TOO_MANY_REQUESTS, + HttpStatus.TOO_MANY_REQUESTS.getReasonPhrase(), + null, null, null); + } + return "ok"; + }; + + executor.execute(operation); + + Assertions.assertEquals(1, sleepCalls.size()); + Assertions.assertEquals(60_000L, sleepCalls.get(0)); + } + + @Test + void testNon429UsesExistingFixedDelay() { + executor = new ResilientCloudOperationExecutor().withRetryCount(3) + .withWaitTimeBetweenRetriesInMillis(5_000L) + .withSleeper(sleepCalls::add); + Supplier operation = () -> { + if (attempts.incrementAndGet() == 1) { + throw new CloudOperationException(HttpStatus.BAD_GATEWAY); + } + return "ok"; + }; + + executor.execute(operation); + + Assertions.assertEquals(1, sleepCalls.size()); + Assertions.assertEquals(5_000L, sleepCalls.get(0)); + } + + @Test + void testRandomDelayAppliedForNon429AfterFirstRetry() { + long deterministicDelay = 45_000L; + executor = new ResilientCloudOperationExecutor().withRetryCount(4) + .withWaitTimeBetweenRetriesInMillis(0) + .withSleeper(sleepCalls::add) + .withRandomDelaySupplier(() -> deterministicDelay); + Supplier operation = () -> { + if (attempts.incrementAndGet() < 3) { + throw new CloudOperationException(HttpStatus.INTERNAL_SERVER_ERROR); + } + return "ok"; + }; + + executor.execute(operation); + + Assertions.assertEquals(2, sleepCalls.size()); + Assertions.assertEquals(0L, sleepCalls.get(0)); + Assertions.assertEquals(deterministicDelay, sleepCalls.get(1)); + } + + @Test + void testRetryableOperationThatNeverSucceedsIsEventuallyRethrown() { + executor = new ResilientCloudOperationExecutor().withRetryCount(3) + .withWaitTimeBetweenRetriesInMillis(0) + .withSleeper(sleepCalls::add) + .withRandomDelaySupplier(() -> 0L); + Supplier operation = () -> { + attempts.incrementAndGet(); + throw new CloudOperationException(HttpStatus.BAD_GATEWAY); + }; + + Assertions.assertThrows(CloudOperationException.class, () -> executor.execute(operation)); + Assertions.assertEquals(3, attempts.get()); + Assertions.assertEquals(2, sleepCalls.size()); + } + + @Test + void testRunnableOverloadIsRetriedThroughOverriddenExecute() { + executor = new ResilientCloudOperationExecutor().withRetryCount(3) + .withWaitTimeBetweenRetriesInMillis(7_000L) + .withSleeper(sleepCalls::add); + Runnable operation = () -> { + if (attempts.incrementAndGet() < 2) { + throw new CloudOperationException(HttpStatus.SERVICE_UNAVAILABLE); + } + }; + + executor.execute(operation); + + Assertions.assertEquals(2, attempts.get()); + Assertions.assertEquals(1, sleepCalls.size()); + Assertions.assertEquals(7_000L, sleepCalls.get(0)); + } } +