From 8cd225adaccd5c7f26a4ea2ac4c02f707602bf43 Mon Sep 17 00:00:00 2001 From: Aditya Vishe Date: Sat, 18 Jul 2026 18:56:40 +0530 Subject: [PATCH] Respect Retry-After in OTLP HTTP senders --- .../exporter/internal/RetryUtil.java | 41 +++++++++++ .../exporter/internal/RetryUtilTest.java | 71 +++++++++++++++++++ .../AbstractHttpTelemetryExporterTest.java | 36 ++++++++++ .../sender/jdk/internal/JdkHttpSender.java | 36 +++++++--- .../okhttp/internal/OkHttpGrpcSender.java | 4 +- .../okhttp/internal/OkHttpHttpSender.java | 9 ++- .../okhttp/internal/RetryInterceptor.java | 18 ++++- .../okhttp/internal/RetryInterceptorTest.java | 45 +++++++++++- 8 files changed, 243 insertions(+), 17 deletions(-) create mode 100644 exporters/common/src/test/java/io/opentelemetry/exporter/internal/RetryUtilTest.java diff --git a/exporters/common/src/main/java/io/opentelemetry/exporter/internal/RetryUtil.java b/exporters/common/src/main/java/io/opentelemetry/exporter/internal/RetryUtil.java index eab9e30cab2..8d6d43ae91c 100644 --- a/exporters/common/src/main/java/io/opentelemetry/exporter/internal/RetryUtil.java +++ b/exporters/common/src/main/java/io/opentelemetry/exporter/internal/RetryUtil.java @@ -6,11 +6,19 @@ package io.opentelemetry.exporter.internal; import io.opentelemetry.sdk.common.export.GrpcStatusCode; +import java.time.Duration; +import java.time.Instant; +import java.time.ZonedDateTime; +import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeParseException; import java.util.Arrays; import java.util.Collections; import java.util.HashSet; +import java.util.OptionalLong; import java.util.Set; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import javax.annotation.Nullable; /** * This class is internal and is hence not for public use. Its APIs are unstable and can change at @@ -47,4 +55,37 @@ public static Set retryableGrpcStatusCodes() { public static Set retryableHttpResponseCodes() { return RETRYABLE_HTTP_STATUS_CODES; } + + /** + * Returns the delay specified by a {@code Retry-After} header, or empty if the value is absent or + * malformed. + */ + public static OptionalLong retryAfterNanos(@Nullable String retryAfter) { + return retryAfterNanos(retryAfter, Instant.now()); + } + + static OptionalLong retryAfterNanos(@Nullable String retryAfter, Instant now) { + if (retryAfter == null) { + return OptionalLong.empty(); + } + + try { + long delaySeconds = Long.parseLong(retryAfter); + if (delaySeconds < 0) { + return OptionalLong.empty(); + } + return OptionalLong.of(TimeUnit.SECONDS.toNanos(delaySeconds)); + } catch (NumberFormatException ignored) { + // Fall through and try the HTTP-date form. + } + + try { + Instant retryAfterInstant = + ZonedDateTime.parse(retryAfter, DateTimeFormatter.RFC_1123_DATE_TIME).toInstant(); + long delayNanos = Duration.between(now, retryAfterInstant).toNanos(); + return OptionalLong.of(Math.max(0, delayNanos)); + } catch (DateTimeParseException | ArithmeticException ignored) { + return OptionalLong.empty(); + } + } } diff --git a/exporters/common/src/test/java/io/opentelemetry/exporter/internal/RetryUtilTest.java b/exporters/common/src/test/java/io/opentelemetry/exporter/internal/RetryUtilTest.java new file mode 100644 index 00000000000..d88e61b8ba6 --- /dev/null +++ b/exporters/common/src/test/java/io/opentelemetry/exporter/internal/RetryUtilTest.java @@ -0,0 +1,71 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.exporter.internal; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.time.Instant; +import java.time.ZoneOffset; +import java.time.ZonedDateTime; +import java.time.format.DateTimeFormatter; +import java.util.OptionalLong; +import java.util.concurrent.TimeUnit; +import org.junit.jupiter.api.Test; + +class RetryUtilTest { + + @Test + void retryAfterNull() { + OptionalLong delayNanos = RetryUtil.retryAfterNanos(null, Instant.EPOCH); + + assertThat(delayNanos).isEmpty(); + } + + @Test + void retryAfterSeconds() { + OptionalLong delayNanos = RetryUtil.retryAfterNanos("30", Instant.EPOCH); + + assertThat(delayNanos).hasValue(TimeUnit.SECONDS.toNanos(30)); + } + + @Test + void retryAfterNegativeSeconds() { + OptionalLong delayNanos = RetryUtil.retryAfterNanos("-1", Instant.EPOCH); + + assertThat(delayNanos).isEmpty(); + } + + @Test + void retryAfterDate() { + Instant now = Instant.parse("2026-07-17T00:00:00Z"); + String retryAfter = + ZonedDateTime.ofInstant(now.plusSeconds(45), ZoneOffset.UTC) + .format(DateTimeFormatter.RFC_1123_DATE_TIME); + + OptionalLong delayNanos = RetryUtil.retryAfterNanos(retryAfter, now); + + assertThat(delayNanos).hasValue(TimeUnit.SECONDS.toNanos(45)); + } + + @Test + void retryAfterPastDateClampsToZero() { + Instant now = Instant.parse("2026-07-17T00:00:00Z"); + String retryAfter = + ZonedDateTime.ofInstant(now.minusSeconds(1), ZoneOffset.UTC) + .format(DateTimeFormatter.RFC_1123_DATE_TIME); + + OptionalLong delayNanos = RetryUtil.retryAfterNanos(retryAfter, now); + + assertThat(delayNanos).hasValue(0L); + } + + @Test + void retryAfterMalformed() { + OptionalLong delayNanos = RetryUtil.retryAfterNanos("bad-value", Instant.EPOCH); + + assertThat(delayNanos).isEmpty(); + } +} diff --git a/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractHttpTelemetryExporterTest.java b/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractHttpTelemetryExporterTest.java index 5f5b13b1adf..95ceb732737 100644 --- a/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractHttpTelemetryExporterTest.java +++ b/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractHttpTelemetryExporterTest.java @@ -729,6 +729,34 @@ void retryableError_tooManyAttempts() { assertThat(attempts).hasValue(2); } + @Test + void retryableError_retryAfterHonored() { + addHttpResponse(502, "0"); + + assertThat( + exporter + .export(Collections.singletonList(generateFakeTelemetry())) + .join(10, TimeUnit.SECONDS) + .isSuccess()) + .isTrue(); + + assertThat(attempts).hasValue(2); + } + + @Test + void retryableError_malformedRetryAfterFallsBack() { + addHttpResponse(503, "not-a-retry-after"); + + assertThat( + exporter + .export(Collections.singletonList(generateFakeTelemetry())) + .join(10, TimeUnit.SECONDS) + .isSuccess()) + .isTrue(); + + assertThat(attempts).hasValue(2); + } + @ParameterizedTest @SuppressLogger(HttpExporter.class) @ValueSource(ints = {400, 401, 403, 500, 501}) @@ -1208,6 +1236,14 @@ private static void addHttpResponse(int code) { httpErrors.add(HttpResponse.of(code)); } + private static void addHttpResponse(int code, String retryAfter) { + httpErrors.add( + HttpResponse.of( + ResponseHeaders.builder(HttpStatus.valueOf(code)) + .add("Retry-After", retryAfter) + .build())); + } + private static void addHttpResponse(int code, AbstractMessageLite bodyMessage) { httpErrors.add( HttpResponse.of( diff --git a/exporters/sender/jdk/src/main/java/io/opentelemetry/exporter/sender/jdk/internal/JdkHttpSender.java b/exporters/sender/jdk/src/main/java/io/opentelemetry/exporter/sender/jdk/internal/JdkHttpSender.java index 6c3610e7345..810f65d6bf2 100644 --- a/exporters/sender/jdk/src/main/java/io/opentelemetry/exporter/sender/jdk/internal/JdkHttpSender.java +++ b/exporters/sender/jdk/src/main/java/io/opentelemetry/exporter/sender/jdk/internal/JdkHttpSender.java @@ -5,6 +5,7 @@ package io.opentelemetry.exporter.sender.jdk.internal; +import io.opentelemetry.exporter.internal.RetryUtil; import io.opentelemetry.sdk.common.CompletableResultCode; import io.opentelemetry.sdk.common.export.Compressor; import io.opentelemetry.sdk.common.export.HttpResponse; @@ -27,6 +28,7 @@ import java.util.List; import java.util.Map; import java.util.Optional; +import java.util.OptionalLong; import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentLinkedQueue; @@ -54,7 +56,7 @@ */ public final class JdkHttpSender implements HttpSender { - private static final Set retryableStatusCodes = Set.of(429, 502, 503, 504); + private static final Set retryableStatusCodes = RetryUtil.retryableHttpResponseCodes(); private static final ThreadLocal threadLocalBaos = ThreadLocal.withInitial(NoCopyByteArrayOutputStream::new); @@ -212,21 +214,30 @@ HttpResponse sendInternal(MessageWriter requestBodyWriter) throws IOException { // If no retry policy, short circuit if (retryPolicy == null) { - return sendRequest(requestBuilder, byteBufferPool); + return toHttpResponse(sendRequest(requestBuilder, byteBufferPool)); } long attempt = 0; long nextBackoffNanos = retryPolicy.getInitialBackoff().toNanos(); HttpResponse httpResponse = null; IOException exception = null; + OptionalLong retryDelayNanos = OptionalLong.empty(); do { if (attempt > 0) { + long remainingNanos = timeout.toNanos() - (System.nanoTime() - startTimeNanos); + if (remainingNanos <= 0) { + break; + } // Compute and sleep for backoff long currentBackoffNanos = Math.min(nextBackoffNanos, retryPolicy.getMaxBackoff().toNanos()); - long backoffNanos = - (long) (ThreadLocalRandom.current().nextDouble(0.8d, 1.2d) * currentBackoffNanos); + long requestedBackoffNanos = + retryDelayNanos.isPresent() + ? retryDelayNanos.getAsLong() + : (long) (ThreadLocalRandom.current().nextDouble(0.8d, 1.2d) * currentBackoffNanos); + long backoffNanos = Math.min(requestedBackoffNanos, remainingNanos); nextBackoffNanos = (long) (currentBackoffNanos * retryPolicy.getBackoffMultiplier()); + retryDelayNanos = OptionalLong.empty(); try { TimeUnit.NANOSECONDS.sleep(backoffNanos); } catch (InterruptedException e) { @@ -243,8 +254,10 @@ HttpResponse sendInternal(MessageWriter requestBodyWriter) throws IOException { exception = null; requestBuilder.timeout(timeout.minusNanos(System.nanoTime() - startTimeNanos)); try { - httpResponse = sendRequest(requestBuilder, byteBufferPool); - boolean retryable = retryableStatusCodes.contains(httpResponse.getStatusCode()); + java.net.http.HttpResponse rawResponse = + sendRequest(requestBuilder, byteBufferPool); + httpResponse = toHttpResponse(rawResponse); + boolean retryable = retryableStatusCodes.contains(rawResponse.statusCode()); if (logger.isLoggable(Level.FINER)) { logger.log( Level.FINER, @@ -258,6 +271,7 @@ HttpResponse sendInternal(MessageWriter requestBodyWriter) throws IOException { if (!retryable) { return httpResponse; } + retryDelayNanos = retryDelayNanos(rawResponse); } catch (IOException e) { exception = e; boolean retryable = retryExceptionPredicate.test(exception); @@ -287,12 +301,10 @@ private static String responseStringRepresentation(HttpResponse response) { return "HttpResponse{code=" + response.getStatusCode() + "}"; } - private HttpResponse sendRequest( + private java.net.http.HttpResponse sendRequest( HttpRequest.Builder requestBuilder, ByteBufferPool byteBufferPool) throws IOException { try { - java.net.http.HttpResponse response = - client.send(requestBuilder.build(), BodyHandlers.ofInputStream()); - return toHttpResponse(response); + return client.send(requestBuilder.build(), BodyHandlers.ofInputStream()); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IllegalStateException(e); @@ -412,6 +424,10 @@ private void resetPool() { } } + private static OptionalLong retryDelayNanos(java.net.http.HttpResponse response) { + return RetryUtil.retryAfterNanos(response.headers().firstValue("Retry-After").orElse(null)); + } + @Override public CompletableResultCode shutdown() { if (managedExecutor) { diff --git a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSender.java b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSender.java index d9c1e490051..85aee7d1679 100644 --- a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSender.java +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSender.java @@ -40,6 +40,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.OptionalLong; import java.util.concurrent.ExecutorService; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; @@ -116,7 +117,8 @@ public OkHttpGrpcSender( .connectTimeout(Duration.ofMillis(connectTimeoutMillis)); if (retryPolicy != null) { clientBuilder.addInterceptor( - new RetryInterceptor(retryPolicy, OkHttpGrpcSender::isRetryable)); + new RetryInterceptor( + retryPolicy, OkHttpGrpcSender::isRetryable, response -> OptionalLong.empty())); } boolean isPlainHttp = endpoint.startsWith("http://"); diff --git a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpHttpSender.java b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpHttpSender.java index b5c611d1bcb..fd91a5fd6c8 100644 --- a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpHttpSender.java +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpHttpSender.java @@ -20,6 +20,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.OptionalLong; import java.util.concurrent.ExecutorService; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; @@ -103,7 +104,9 @@ public OkHttpHttpSender( } if (retryPolicy != null) { - builder.addInterceptor(new RetryInterceptor(retryPolicy, OkHttpHttpSender::isRetryable)); + builder.addInterceptor( + new RetryInterceptor( + retryPolicy, OkHttpHttpSender::isRetryable, OkHttpHttpSender::retryDelayNanos)); } boolean isPlainHttp = endpoint.getScheme().equals("http"); @@ -121,6 +124,10 @@ public OkHttpHttpSender( this.maxResponseBodySize = maxResponseBodySize; } + private static OptionalLong retryDelayNanos(Response response) { + return RetryUtil.retryAfterNanos(response.header("Retry-After")); + } + @Override public void send( MessageWriter messageWriter, Consumer onResponse, Consumer onError) { diff --git a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptor.java b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptor.java index 988c8277c26..145cf65b70f 100644 --- a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptor.java +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptor.java @@ -13,6 +13,7 @@ import java.net.SocketException; import java.net.SocketTimeoutException; import java.net.UnknownHostException; +import java.util.OptionalLong; import java.util.StringJoiner; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; @@ -36,15 +37,20 @@ public final class RetryInterceptor implements Interceptor { private final RetryPolicy retryPolicy; private final Function isRetryable; + private final Function retryDelayNanosExtractor; private final Predicate retryExceptionPredicate; private final Sleeper sleeper; private final Supplier randomJitter; /** Constructs a new retrier. */ - public RetryInterceptor(RetryPolicy retryPolicy, Function isRetryable) { + public RetryInterceptor( + RetryPolicy retryPolicy, + Function isRetryable, + Function retryDelayNanosExtractor) { this( retryPolicy, isRetryable, + retryDelayNanosExtractor, retryPolicy.getRetryExceptionPredicate() == null ? RetryInterceptor::isRetryableException : retryPolicy.getRetryExceptionPredicate(), @@ -56,11 +62,13 @@ public RetryInterceptor(RetryPolicy retryPolicy, Function isR RetryInterceptor( RetryPolicy retryPolicy, Function isRetryable, + Function retryDelayNanosExtractor, Predicate retryExceptionPredicate, Sleeper sleeper, Supplier randomJitter) { this.retryPolicy = retryPolicy; this.isRetryable = isRetryable; + this.retryDelayNanosExtractor = retryDelayNanosExtractor; this.retryExceptionPredicate = retryExceptionPredicate; this.sleeper = sleeper; this.randomJitter = randomJitter; @@ -72,14 +80,19 @@ public Response intercept(Chain chain) throws IOException { IOException exception = null; int attempt = 0; long nextBackoffNanos = retryPolicy.getInitialBackoff().toNanos(); + OptionalLong retryDelayNanos = OptionalLong.empty(); do { if (attempt > 0) { // Compute and sleep for backoff // https://github.com/grpc/proposal/blob/master/A6-client-retries.md#exponential-backoff long currentBackoffNanos = Math.min(nextBackoffNanos, retryPolicy.getMaxBackoff().toNanos()); - long backoffNanos = (long) (randomJitter.get() * currentBackoffNanos); + long backoffNanos = + retryDelayNanos.isPresent() + ? retryDelayNanos.getAsLong() + : (long) (randomJitter.get() * currentBackoffNanos); nextBackoffNanos = (long) (currentBackoffNanos * retryPolicy.getBackoffMultiplier()); + retryDelayNanos = OptionalLong.empty(); try { sleeper.sleep(backoffNanos); } catch (InterruptedException e) { @@ -109,6 +122,7 @@ public Response intercept(Chain chain) throws IOException { if (!retryable) { return response; } + retryDelayNanos = retryDelayNanosExtractor.apply(response); } else { throw new NullPointerException("response cannot be null."); } diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptorTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptorTest.java index eeff92abf3d..33968f19308 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptorTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryInterceptorTest.java @@ -31,6 +31,7 @@ import java.net.SocketTimeoutException; import java.net.UnknownHostException; import java.time.Duration; +import java.util.OptionalLong; import java.util.concurrent.TimeUnit; import java.util.function.Predicate; import java.util.function.Supplier; @@ -92,7 +93,12 @@ public boolean test(IOException e) { retrier = new RetryInterceptor( - retryPolicy, r -> !r.isSuccessful(), retryExceptionPredicate, sleeper, random); + retryPolicy, + r -> !r.isSuccessful(), + response -> OptionalLong.empty(), + retryExceptionPredicate, + sleeper, + random); client = new OkHttpClient.Builder().addInterceptor(retrier).build(); } @@ -178,6 +184,35 @@ public Void answer(InvocationOnMock invocation) throws Throwable { } } + @Test + void retryDelayOverride() throws Exception { + succeedOnAttempt(2); + doNothing().when(sleeper).sleep(anyLong()); + + RetryInterceptor retryInterceptor = + new RetryInterceptor( + RetryPolicy.builder() + .setBackoffMultiplier(1.6) + .setInitialBackoff(Duration.ofSeconds(1)) + .setMaxBackoff(Duration.ofSeconds(2)) + .setMaxAttempts(5) + .setRetryExceptionPredicate(retryExceptionPredicate) + .build(), + r -> !r.isSuccessful(), + response -> OptionalLong.of(123L), + retryExceptionPredicate, + sleeper, + random); + client = new OkHttpClient.Builder().addInterceptor(retryInterceptor).build(); + + try (Response response = sendRequest()) { + assertThat(response.isSuccessful()).isTrue(); + } + + verify(sleeper).sleep(123L); + verifyNoInteractions(random); + } + @Test void connectTimeout() throws Exception { client = connectTimeoutClient(); @@ -289,7 +324,10 @@ private static Stream isRetryableExceptionArgs() { @Test void isRetryableExceptionDefaultBehaviour() { RetryInterceptor retryInterceptor = - new RetryInterceptor(RetryPolicy.getDefault(), OkHttpHttpSender::isRetryable); + new RetryInterceptor( + RetryPolicy.getDefault(), + OkHttpHttpSender::isRetryable, + response -> OptionalLong.empty()); assertThat( retryInterceptor.shouldRetryOnException( new SocketTimeoutException("Connect timed out"))) @@ -305,7 +343,8 @@ void isRetryableExceptionCustomRetryPredicate() { RetryPolicy.builder() .setRetryExceptionPredicate((IOException e) -> e.getMessage().equals("retry")) .build(), - OkHttpHttpSender::isRetryable); + OkHttpHttpSender::isRetryable, + response -> OptionalLong.empty()); assertThat(retryInterceptor.shouldRetryOnException(new IOException("some message"))).isFalse(); assertThat(retryInterceptor.shouldRetryOnException(new IOException("retry"))).isTrue();