From 7d5bbe9ff5836b0b00f00b91383e01be18095505 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 23 Sep 2026 14:14:06 +0300 Subject: [PATCH 01/12] fix: retry OTLP gRPC responses with trailer status MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- CHANGELOG.md | 4 ++++ .../sender/okhttp/internal/OkHttpGrpcSender.java | 10 ++++++---- .../okhttp/internal/OkHttpGrpcSenderTest.java | 16 ++++++++++++++++ 3 files changed, 26 insertions(+), 4 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 66bf947bf05..f629b3a73cb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +### Exporters + +* OTLP gRPC: Retry responses that report a retryable status in trailers. + ## Version 1.66.0 (2026-09-11) ### API 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 48b990cf5a0..8ba4a202103 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 @@ -366,13 +366,15 @@ public CompletableResultCode shutdown() { /** Whether response is retriable or not. */ public static boolean isRetryable(Response response) { - // We don't check trailers for retry since retryable error codes always come with response - // headers, not trailers, in practice. String grpcStatus = response.header(GRPC_STATUS); if (grpcStatus == null) { - return false; + try { + grpcStatus = response.trailers().get(GRPC_STATUS); + } catch (IOException e) { + return false; + } } - return RetryUtil.retryableGrpcStatusCodes().contains(grpcStatus); + return grpcStatus != null && RetryUtil.retryableGrpcStatusCodes().contains(grpcStatus); } // From grpc-java diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java index f68350ca3b0..fd28417eb2c 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java @@ -32,6 +32,7 @@ import java.util.concurrent.atomic.AtomicReference; import javax.net.ssl.SSLContext; import javax.net.ssl.SSLException; +import okhttp3.Headers; import okhttp3.MediaType; import okhttp3.Protocol; import okhttp3.Request; @@ -67,6 +68,21 @@ void isRetryable_NonRetryableGrpcStatus() { assertFalse(isRetryable); } + @Test + void isRetryable_RetryableGrpcStatusInTrailers() { + Response response = + new Response.Builder() + .request(new Request.Builder().url("http://localhost/").build()) + .protocol(Protocol.HTTP_2) + .code(200) + .body(ResponseBody.create("body", TEXT_PLAIN)) + .message("Retryable") + .trailers(() -> Headers.of(GRPC_STATUS, "14")) + .build(); + + assertTrue(OkHttpGrpcSender.isRetryable(response)); + } + @Test void send_rejectedExecution_callsOnError() { ThreadPoolExecutor executor = From b95dd4da1d18b7b4a4cb5cc8674107ff050805f5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 23 Sep 2026 14:14:44 +0300 Subject: [PATCH 02/12] docs: link OTLP trailer retry changelog MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- CHANGELOG.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index f629b3a73cb..129c92ad4f1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,8 @@ ### Exporters -* OTLP gRPC: Retry responses that report a retryable status in trailers. +* OTLP gRPC: Retry responses that report a retryable status in trailers + ([#8854](https://github.com/open-telemetry/opentelemetry-java/pull/8854)). ## Version 1.66.0 (2026-09-11) From d129565566dee0e269a2a14e377474d15c22efd5 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 23 Sep 2026 14:51:41 +0300 Subject: [PATCH 03/12] fix: preserve response bodies while reading trailers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../okhttp/internal/OkHttpGrpcSender.java | 56 ++++++++++++++++++- .../okhttp/internal/RetryInterceptor.java | 48 +++++++++++++++- .../okhttp/internal/OkHttpGrpcSenderTest.java | 5 +- 3 files changed, 104 insertions(+), 5 deletions(-) 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 8ba4a202103..ef58da3aa24 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 @@ -57,7 +57,9 @@ import okhttp3.Callback; import okhttp3.ConnectionSpec; import okhttp3.Dispatcher; +import okhttp3.Headers; import okhttp3.HttpUrl; +import okhttp3.MediaType; import okhttp3.OkHttpClient; import okhttp3.Protocol; import okhttp3.Request; @@ -66,6 +68,7 @@ import okhttp3.ResponseBody; import okhttp3.TlsVersion; import okio.Buffer; +import okio.BufferedSource; import okio.GzipSource; /** @@ -122,7 +125,10 @@ public OkHttpGrpcSender( if (retryPolicy != null) { clientBuilder.addInterceptor( new RetryInterceptor( - retryPolicy, OkHttpGrpcSender::isRetryable, response -> OptionalLong.empty())); + retryPolicy, + OkHttpGrpcSender::isRetryable, + response -> OptionalLong.empty(), + response -> prepareResponseForRetry(response, maxResponseBodySize))); } boolean isPlainHttp = endpoint.startsWith("http://"); @@ -370,13 +376,59 @@ public static boolean isRetryable(Response response) { if (grpcStatus == null) { try { grpcStatus = response.trailers().get(GRPC_STATUS); - } catch (IOException e) { + } catch (IOException | IllegalStateException e) { return false; } } return grpcStatus != null && RetryUtil.retryableGrpcStatusCodes().contains(grpcStatus); } + private static Response prepareResponseForRetry(Response response, long maxResponseBodySize) + throws IOException { + if (response.header(GRPC_STATUS) != null) { + return response; + } + + ResponseBody body = response.body(); + Buffer buffer = new Buffer(); + long readUpTo = + maxResponseBodySize >= Long.MAX_VALUE - 5 ? Long.MAX_VALUE : maxResponseBodySize + 6; + while (buffer.size() < readUpTo) { + long read = body.source().read(buffer, readUpTo - buffer.size()); + if (read == -1L) { + break; + } + } + + boolean responseBodyTooLarge = buffer.size() > maxResponseBodySize; + Headers trailers = responseBodyTooLarge ? null : response.trailers(); + Buffer replacementBuffer = buffer; + ResponseBody replacementBody = + new ResponseBody() { + @Override + public long contentLength() { + return replacementBuffer.size(); + } + + @Override + public MediaType contentType() { + return body.contentType(); + } + + @Override + public BufferedSource source() { + return replacementBuffer; + } + }; + Response.Builder responseBuilder = response.newBuilder(); + response.close(); + responseBuilder.body(replacementBody); + if (trailers != null) { + responseBuilder.trailers(() -> trailers); + } + return responseBuilder.build(); + } + // From grpc-java /** Unescape the provided ascii to a unicode {@link String}. */ 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 145cf65b70f..24a58e2768b 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 @@ -38,6 +38,7 @@ public final class RetryInterceptor implements Interceptor { private final RetryPolicy retryPolicy; private final Function isRetryable; private final Function retryDelayNanosExtractor; + private final ResponseTransformer responseTransformer; private final Predicate retryExceptionPredicate; private final Sleeper sleeper; private final Supplier randomJitter; @@ -55,7 +56,26 @@ public RetryInterceptor( ? RetryInterceptor::isRetryableException : retryPolicy.getRetryExceptionPredicate(), TimeUnit.NANOSECONDS::sleep, - () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d)); + () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), + response -> response); + } + + // Visible for testing + RetryInterceptor( + RetryPolicy retryPolicy, + Function isRetryable, + Function retryDelayNanosExtractor, + ResponseTransformer responseTransformer) { + this( + retryPolicy, + isRetryable, + retryDelayNanosExtractor, + retryPolicy.getRetryExceptionPredicate() == null + ? RetryInterceptor::isRetryableException + : retryPolicy.getRetryExceptionPredicate(), + TimeUnit.NANOSECONDS::sleep, + () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), + responseTransformer); } // Visible for testing @@ -66,12 +86,32 @@ public RetryInterceptor( Predicate retryExceptionPredicate, Sleeper sleeper, Supplier randomJitter) { + this( + retryPolicy, + isRetryable, + retryDelayNanosExtractor, + retryExceptionPredicate, + sleeper, + randomJitter, + response -> response); + } + + // Visible for testing + RetryInterceptor( + RetryPolicy retryPolicy, + Function isRetryable, + Function retryDelayNanosExtractor, + Predicate retryExceptionPredicate, + Sleeper sleeper, + Supplier randomJitter, + ResponseTransformer responseTransformer) { this.retryPolicy = retryPolicy; this.isRetryable = isRetryable; this.retryDelayNanosExtractor = retryDelayNanosExtractor; this.retryExceptionPredicate = retryExceptionPredicate; this.sleeper = sleeper; this.randomJitter = randomJitter; + this.responseTransformer = responseTransformer; } @Override @@ -108,6 +148,7 @@ public Response intercept(Chain chain) throws IOException { try { response = chain.proceed(chain.request()); if (response != null) { + response = responseTransformer.transform(response); boolean retryable = Boolean.TRUE.equals(isRetryable.apply(response)); if (logger.isLoggable(Level.FINER)) { logger.log( @@ -152,6 +193,11 @@ public Response intercept(Chain chain) throws IOException { throw exception; } + @FunctionalInterface + interface ResponseTransformer { + Response transform(Response response) throws IOException; + } + private static String responseStringRepresentation(Response response) { StringJoiner joiner = new StringJoiner(",", "Response{", "}"); joiner.add("code=" + response.code()); diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java index fd28417eb2c..34c43990a5c 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java @@ -69,18 +69,19 @@ void isRetryable_NonRetryableGrpcStatus() { } @Test - void isRetryable_RetryableGrpcStatusInTrailers() { + void isRetryable_RetryableGrpcStatusInTrailers() throws IOException { Response response = new Response.Builder() .request(new Request.Builder().url("http://localhost/").build()) .protocol(Protocol.HTTP_2) .code(200) - .body(ResponseBody.create("body", TEXT_PLAIN)) + .body(ResponseBody.create("", TEXT_PLAIN)) .message("Retryable") .trailers(() -> Headers.of(GRPC_STATUS, "14")) .build(); assertTrue(OkHttpGrpcSender.isRetryable(response)); + assertThat(response.body().string()).isEmpty(); } @Test From 65032eb8486f816f2838489e9c84c04339e78fef Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 23 Sep 2026 15:17:34 +0300 Subject: [PATCH 04/12] fix: preserve retry response trailers --- .../exporter/sender/okhttp/internal/OkHttpGrpcSender.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) 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 ef58da3aa24..e7446ed3946 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 @@ -401,7 +401,7 @@ private static Response prepareResponseForRetry(Response response, long maxRespo } boolean responseBodyTooLarge = buffer.size() > maxResponseBodySize; - Headers trailers = responseBodyTooLarge ? null : response.trailers(); + Headers trailers = responseBodyTooLarge ? Headers.of() : response.trailers(); Buffer replacementBuffer = buffer; ResponseBody replacementBody = new ResponseBody() { @@ -423,9 +423,7 @@ public BufferedSource source() { Response.Builder responseBuilder = response.newBuilder(); response.close(); responseBuilder.body(replacementBody); - if (trailers != null) { - responseBuilder.trailers(() -> trailers); - } + responseBuilder.trailers(() -> trailers); return responseBuilder.build(); } From 9f3bc14d3c95107f64f45a90e28cb4f0612d2ec4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 23 Sep 2026 15:51:02 +0300 Subject: [PATCH 05/12] fix: preserve trailers across OkHttp versions --- .../exporter/sender/okhttp/internal/OkHttpGrpcSender.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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 e7446ed3946..8b6aacae0e7 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 @@ -423,7 +423,9 @@ public BufferedSource source() { Response.Builder responseBuilder = response.newBuilder(); response.close(); responseBuilder.body(replacementBody); - responseBuilder.trailers(() -> trailers); + if (trailers.size() > 0) { + responseBuilder.trailers(() -> trailers); + } return responseBuilder.build(); } From d0e3ead1750b3e5da9b578b14c4274e2fd69b9ac Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 23 Sep 2026 17:16:27 +0300 Subject: [PATCH 06/12] fix: support trailer retry on older OkHttp MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../sender/okhttp/internal/OkHttpGrpcSender.java | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) 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 8b6aacae0e7..e7cda98f3c9 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 @@ -423,8 +423,13 @@ public BufferedSource source() { Response.Builder responseBuilder = response.newBuilder(); response.close(); responseBuilder.body(replacementBody); - if (trailers.size() > 0) { - responseBuilder.trailers(() -> trailers); + String grpcStatus = trailers.get(GRPC_STATUS); + if (grpcStatus != null) { + responseBuilder.header(GRPC_STATUS, grpcStatus); + } + String grpcMessage = trailers.get(GRPC_MESSAGE); + if (grpcMessage != null) { + responseBuilder.header(GRPC_MESSAGE, grpcMessage); } return responseBuilder.build(); } From a7ee760dfa8f5479d8611ab7f4af1e8e8816895b Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= <72094408+efegokdemir@users.noreply.github.com> Date: Wed, 30 Sep 2026 16:34:10 +0300 Subject: [PATCH 07/12] refactor okhttp grpc retries around resolved responses --- .../okhttp/internal/OkHttpGrpcSender.java | 136 ++++++++---------- .../okhttp/internal/RetryInterceptor.java | 109 ++------------ .../sender/okhttp/internal/RetryState.java | 102 +++++++++++++ .../okhttp/internal/OkHttpGrpcSenderTest.java | 81 ++++++++--- 4 files changed, 231 insertions(+), 197 deletions(-) create mode 100644 exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java 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 e7cda98f3c9..c5449860928 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 @@ -57,18 +57,14 @@ import okhttp3.Callback; import okhttp3.ConnectionSpec; import okhttp3.Dispatcher; -import okhttp3.Headers; import okhttp3.HttpUrl; -import okhttp3.MediaType; import okhttp3.OkHttpClient; import okhttp3.Protocol; import okhttp3.Request; -import okhttp3.RequestBody; import okhttp3.Response; import okhttp3.ResponseBody; import okhttp3.TlsVersion; import okio.Buffer; -import okio.BufferedSource; import okio.GzipSource; /** @@ -90,6 +86,7 @@ public final class OkHttpGrpcSender implements GrpcSender { @Nullable private final Compressor compressor; private final Supplier>> headersSupplier; private final long maxResponseBodySize; + @Nullable private final RetryPolicy retryPolicy; /** Creates a new {@link OkHttpGrpcSender}. */ @SuppressWarnings("TooManyParameters") @@ -122,15 +119,6 @@ public OkHttpGrpcSender( .dispatcher(dispatcher) .callTimeout(Duration.ofMillis(callTimeoutMillis)) .connectTimeout(Duration.ofMillis(connectTimeoutMillis)); - if (retryPolicy != null) { - clientBuilder.addInterceptor( - new RetryInterceptor( - retryPolicy, - OkHttpGrpcSender::isRetryable, - response -> OptionalLong.empty(), - response -> prepareResponseForRetry(response, maxResponseBodySize))); - } - boolean isPlainHttp = endpoint.startsWith("http://"); if (isPlainHttp) { clientBuilder.connectionSpecs(Collections.singletonList(ConnectionSpec.CLEARTEXT)); @@ -164,6 +152,7 @@ public OkHttpGrpcSender( this.headersSupplier = headersSupplier; this.url = HttpUrl.get(endpoint); this.maxResponseBodySize = maxResponseBodySize; + this.retryPolicy = retryPolicy; } @Override @@ -182,8 +171,23 @@ public void send( if (compressor != null) { requestBuilder.addHeader("grpc-encoding", compressor.getEncoding()); } - RequestBody requestBody = new GrpcRequestBody(messageWriter, compressor); - requestBuilder.post(requestBody); + requestBuilder.post(new GrpcRequestBody(messageWriter, compressor)); + + sendAttempt(requestBuilder, messageWriter, onResponse, onError, 0, newRetryState()); + } + + @Nullable + private RetryState newRetryState() { + return retryPolicy == null ? null : new RetryState(retryPolicy); + } + + private void sendAttempt( + Request.Builder requestBuilder, + MessageWriter messageWriter, + Consumer onResponse, + Consumer onError, + int attempt, + @Nullable RetryState retryState) { try { InstrumentationUtil.suppressInstrumentation( @@ -194,12 +198,42 @@ public void send( new Callback() { @Override public void onFailure(Call call, IOException e) { - onError.accept(e); + if (retryState != null + && retryState.canRetry(attempt) + && retryState.shouldRetryOnException(e) + && retryState.backoff(OptionalLong.empty())) { + sendAttempt( + requestBuilder, + messageWriter, + onResponse, + onError, + attempt + 1, + retryState); + } else { + onError.accept(e); + } } @Override public void onResponse(Call call, Response response) { - handleResponse(response, onResponse); + handleResponse( + response, + resolvedResponse -> { + if (retryState != null + && retryState.canRetry(attempt) + && isRetryable(resolvedResponse) + && retryState.backoff(OptionalLong.empty())) { + sendAttempt( + requestBuilder, + messageWriter, + onResponse, + onError, + attempt + 1, + retryState); + } else { + onResponse.accept(resolvedResponse); + } + }); } })); } catch (RejectedExecutionException e) { @@ -207,7 +241,7 @@ public void onResponse(Call call, Response response) { } } - private void handleResponse(Response response, Consumer onResponse) { + void handleResponse(Response response, Consumer onResponse) { try (ResponseBody body = response.body()) { // A gRPC message frame has a 5-byte header: 1 compression-flag byte + 4 message-length // bytes. Read the header first so that the size limit applies to the message payload only, @@ -370,68 +404,10 @@ public CompletableResultCode shutdown() { return CompletableResultCode.ofSuccess(); } - /** Whether response is retriable or not. */ - public static boolean isRetryable(Response response) { - String grpcStatus = response.header(GRPC_STATUS); - if (grpcStatus == null) { - try { - grpcStatus = response.trailers().get(GRPC_STATUS); - } catch (IOException | IllegalStateException e) { - return false; - } - } - return grpcStatus != null && RetryUtil.retryableGrpcStatusCodes().contains(grpcStatus); - } - - private static Response prepareResponseForRetry(Response response, long maxResponseBodySize) - throws IOException { - if (response.header(GRPC_STATUS) != null) { - return response; - } - - ResponseBody body = response.body(); - Buffer buffer = new Buffer(); - long readUpTo = - maxResponseBodySize >= Long.MAX_VALUE - 5 ? Long.MAX_VALUE : maxResponseBodySize + 6; - while (buffer.size() < readUpTo) { - long read = body.source().read(buffer, readUpTo - buffer.size()); - if (read == -1L) { - break; - } - } - - boolean responseBodyTooLarge = buffer.size() > maxResponseBodySize; - Headers trailers = responseBodyTooLarge ? Headers.of() : response.trailers(); - Buffer replacementBuffer = buffer; - ResponseBody replacementBody = - new ResponseBody() { - @Override - public long contentLength() { - return replacementBuffer.size(); - } - - @Override - public MediaType contentType() { - return body.contentType(); - } - - @Override - public BufferedSource source() { - return replacementBuffer; - } - }; - Response.Builder responseBuilder = response.newBuilder(); - response.close(); - responseBuilder.body(replacementBody); - String grpcStatus = trailers.get(GRPC_STATUS); - if (grpcStatus != null) { - responseBuilder.header(GRPC_STATUS, grpcStatus); - } - String grpcMessage = trailers.get(GRPC_MESSAGE); - if (grpcMessage != null) { - responseBuilder.header(GRPC_MESSAGE, grpcMessage); - } - return responseBuilder.build(); + /** Whether a resolved response is retriable or not. */ + static boolean isRetryable(GrpcResponse response) { + return RetryUtil.retryableGrpcStatusCodes() + .contains(Integer.toString(response.getStatusCode().getValue())); } // From grpc-java 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 24a58e2768b..e8af7bf832f 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 @@ -9,14 +9,8 @@ import io.opentelemetry.sdk.common.export.RetryPolicy; import java.io.IOException; -import java.net.ConnectException; -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; import java.util.function.Function; import java.util.function.Predicate; import java.util.function.Supplier; @@ -38,10 +32,7 @@ public final class RetryInterceptor implements Interceptor { private final RetryPolicy retryPolicy; private final Function isRetryable; private final Function retryDelayNanosExtractor; - private final ResponseTransformer responseTransformer; - private final Predicate retryExceptionPredicate; - private final Sleeper sleeper; - private final Supplier randomJitter; + private final RetryState retryState; /** Constructs a new retrier. */ public RetryInterceptor( @@ -53,29 +44,10 @@ public RetryInterceptor( isRetryable, retryDelayNanosExtractor, retryPolicy.getRetryExceptionPredicate() == null - ? RetryInterceptor::isRetryableException + ? RetryState::isRetryableException : retryPolicy.getRetryExceptionPredicate(), - TimeUnit.NANOSECONDS::sleep, - () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), - response -> response); - } - - // Visible for testing - RetryInterceptor( - RetryPolicy retryPolicy, - Function isRetryable, - Function retryDelayNanosExtractor, - ResponseTransformer responseTransformer) { - this( - retryPolicy, - isRetryable, - retryDelayNanosExtractor, - retryPolicy.getRetryExceptionPredicate() == null - ? RetryInterceptor::isRetryableException - : retryPolicy.getRetryExceptionPredicate(), - TimeUnit.NANOSECONDS::sleep, - () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), - responseTransformer); + RetryState::defaultSleeper, + RetryState::defaultRandomJitter); } // Visible for testing @@ -86,32 +58,10 @@ public RetryInterceptor( Predicate retryExceptionPredicate, Sleeper sleeper, Supplier randomJitter) { - this( - retryPolicy, - isRetryable, - retryDelayNanosExtractor, - retryExceptionPredicate, - sleeper, - randomJitter, - response -> response); - } - - // Visible for testing - RetryInterceptor( - RetryPolicy retryPolicy, - Function isRetryable, - Function retryDelayNanosExtractor, - Predicate retryExceptionPredicate, - Sleeper sleeper, - Supplier randomJitter, - ResponseTransformer responseTransformer) { this.retryPolicy = retryPolicy; this.isRetryable = isRetryable; this.retryDelayNanosExtractor = retryDelayNanosExtractor; - this.retryExceptionPredicate = retryExceptionPredicate; - this.sleeper = sleeper; - this.randomJitter = randomJitter; - this.responseTransformer = responseTransformer; + this.retryState = new RetryState(retryPolicy, retryExceptionPredicate, sleeper, randomJitter); } @Override @@ -119,26 +69,15 @@ public Response intercept(Chain chain) throws IOException { Response response = null; 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 = - retryDelayNanos.isPresent() - ? retryDelayNanos.getAsLong() - : (long) (randomJitter.get() * currentBackoffNanos); - nextBackoffNanos = (long) (currentBackoffNanos * retryPolicy.getBackoffMultiplier()); - retryDelayNanos = OptionalLong.empty(); - try { - sleeper.sleep(backoffNanos); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); + if (!retryState.backoff(retryDelayNanos)) { break; // Break out and return response or throw } + retryDelayNanos = OptionalLong.empty(); // Close response from previous attempt if (response != null) { response.close(); @@ -148,7 +87,6 @@ public Response intercept(Chain chain) throws IOException { try { response = chain.proceed(chain.request()); if (response != null) { - response = responseTransformer.transform(response); boolean retryable = Boolean.TRUE.equals(isRetryable.apply(response)); if (logger.isLoggable(Level.FINER)) { logger.log( @@ -170,7 +108,7 @@ public Response intercept(Chain chain) throws IOException { } catch (IOException e) { exception = e; response = null; - boolean retryable = retryExceptionPredicate.test(exception); + boolean retryable = retryState.shouldRetryOnException(exception); if (logger.isLoggable(Level.FINER)) { logger.log( Level.FINER, @@ -193,11 +131,6 @@ public Response intercept(Chain chain) throws IOException { throw exception; } - @FunctionalInterface - interface ResponseTransformer { - Response transform(Response response) throws IOException; - } - private static String responseStringRepresentation(Response response) { StringJoiner joiner = new StringJoiner(",", "Response{", "}"); joiner.add("code=" + response.code()); @@ -210,31 +143,15 @@ private static String responseStringRepresentation(Response response) { } // Visible for testing - boolean shouldRetryOnException(IOException e) { - return retryExceptionPredicate.test(e); + static boolean isRetryableException(IOException e) { + return RetryState.isRetryableException(e); } // Visible for testing - static boolean isRetryableException(IOException e) { - // Known retryable SocketTimeoutException messages: null, "connect timed out", "timeout" - // Known retryable ConnectTimeout messages: "Failed to connect to - // localhost/[0:0:0:0:0:0:0:1]:62611" - // Known retryable UnknownHostException messages: "xxxxxx.com" - // Known retryable SocketException: Socket closed - if (e instanceof SocketTimeoutException) { - return true; - } else if (e instanceof ConnectException) { - return true; - } else if (e instanceof UnknownHostException) { - return true; - } else if (e instanceof SocketException) { - return true; - } - return false; + boolean shouldRetryOnException(IOException e) { + return retryState.shouldRetryOnException(e); } // Visible for testing - interface Sleeper { - void sleep(long delayNanos) throws InterruptedException; - } + interface Sleeper extends RetryState.Sleeper {} } diff --git a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java new file mode 100644 index 00000000000..ffcf261dc0c --- /dev/null +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java @@ -0,0 +1,102 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.exporter.sender.okhttp.internal; + +import io.opentelemetry.sdk.common.export.RetryPolicy; +import java.io.IOException; +import java.net.ConnectException; +import java.net.SocketException; +import java.net.SocketTimeoutException; +import java.net.UnknownHostException; +import java.util.OptionalLong; +import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.TimeUnit; +import java.util.function.Predicate; +import java.util.function.Supplier; + +/** Shared retry policy state and mechanics for OkHttp senders. */ +final class RetryState { + + private final RetryPolicy retryPolicy; + private final Predicate retryExceptionPredicate; + private final Sleeper sleeper; + private final Supplier randomJitter; + private long nextBackoffNanos; + + RetryState(RetryPolicy retryPolicy) { + this( + retryPolicy, + retryPolicy.getRetryExceptionPredicate() == null + ? RetryState::isRetryableException + : retryPolicy.getRetryExceptionPredicate(), + TimeUnit.NANOSECONDS::sleep, + () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d)); + } + + // Visible for testing. + RetryState( + RetryPolicy retryPolicy, + Predicate retryExceptionPredicate, + Sleeper sleeper, + Supplier randomJitter) { + this.retryPolicy = retryPolicy; + this.retryExceptionPredicate = retryExceptionPredicate; + this.sleeper = sleeper; + this.randomJitter = randomJitter; + this.nextBackoffNanos = retryPolicy.getInitialBackoff().toNanos(); + } + + static void defaultSleeper(long delayNanos) throws InterruptedException { + TimeUnit.NANOSECONDS.sleep(delayNanos); + } + + static double defaultRandomJitter() { + return ThreadLocalRandom.current().nextDouble(0.8d, 1.2d); + } + + boolean backoff(OptionalLong retryDelayNanos) { + long currentBackoffNanos = Math.min(nextBackoffNanos, retryPolicy.getMaxBackoff().toNanos()); + long backoffNanos = + retryDelayNanos.isPresent() + ? retryDelayNanos.getAsLong() + : (long) (randomJitter.get() * currentBackoffNanos); + nextBackoffNanos = (long) (currentBackoffNanos * retryPolicy.getBackoffMultiplier()); + try { + sleeper.sleep(backoffNanos); + return true; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } + } + + boolean canRetry(int attempt) { + return attempt + 1 < retryPolicy.getMaxAttempts(); + } + + boolean shouldRetryOnException(IOException exception) { + return retryExceptionPredicate.test(exception); + } + + // Visible for testing. + static boolean isRetryableException(IOException e) { + if (e instanceof SocketTimeoutException) { + return true; + } else if (e instanceof ConnectException) { + return true; + } else if (e instanceof UnknownHostException) { + return true; + } else if (e instanceof SocketException) { + return true; + } + return false; + } + + @FunctionalInterface + interface Sleeper { + void sleep(long delayNanos) throws InterruptedException; + } +} diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java index 34c43990a5c..2e9bf518a8a 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java @@ -44,8 +44,7 @@ class OkHttpGrpcSenderTest { - private static final String GRPC_STATUS = "grpc-status"; - private static final MediaType TEXT_PLAIN = MediaType.get("text/plain"); + private static final MediaType GRPC_MEDIA_TYPE = MediaType.get("application/grpc"); static Set provideRetryableGrpcStatusCodes() { return RetryUtil.retryableGrpcStatusCodes(); @@ -54,7 +53,11 @@ static Set provideRetryableGrpcStatusCodes() { @ParameterizedTest(name = "isRetryable should return true for GRPC status code: {0}") @MethodSource("provideRetryableGrpcStatusCodes") void isRetryable_RetryableGrpcStatus(String retryableGrpcStatus) { - Response response = createResponse(503, retryableGrpcStatus, "Retryable"); + GrpcResponse response = + ImmutableGrpcResponse.create( + GrpcStatusCode.fromValue(Integer.parseInt(retryableGrpcStatus)), + "Retryable", + new byte[0]); boolean isRetryable = OkHttpGrpcSender.isRetryable(response); assertTrue(isRetryable); } @@ -63,25 +66,72 @@ void isRetryable_RetryableGrpcStatus(String retryableGrpcStatus) { void isRetryable_NonRetryableGrpcStatus() { String nonRetryableGrpcStatus = Integer.valueOf(GrpcStatusCode.UNKNOWN.getValue()).toString(); // INVALID_ARGUMENT - Response response = createResponse(503, nonRetryableGrpcStatus, "Non-retryable"); + GrpcResponse response = + ImmutableGrpcResponse.create( + GrpcStatusCode.fromValue(Integer.parseInt(nonRetryableGrpcStatus)), + "Non-retryable", + new byte[0]); boolean isRetryable = OkHttpGrpcSender.isRetryable(response); assertFalse(isRetryable); } @Test - void isRetryable_RetryableGrpcStatusInTrailers() throws IOException { + void handleResponse_resolvesTrailerStatusAndMessageAfterConsumingBody() { + OkHttpGrpcSender sender = createSender(Long.MAX_VALUE); + AtomicReference responseRef = new AtomicReference<>(); + byte[] frame = new byte[] {0, 0, 0, 0, 3, 'o', 'k', '!'}; + Response response = + new Response.Builder() + .request(new Request.Builder().url("http://localhost/").build()) + .protocol(Protocol.HTTP_2) + .code(200) + .body(ResponseBody.create(frame, GRPC_MEDIA_TYPE)) + .message("HTTP message") + .trailers(() -> Headers.of("grpc-status", "14", "grpc-message", "retry%20me")) + .build(); + + sender.handleResponse(response, responseRef::set); + + assertThat(responseRef.get().getStatusCode()).isEqualTo(GrpcStatusCode.UNAVAILABLE); + assertThat(responseRef.get().getStatusDescription()).isEqualTo("retry me"); + assertThat(responseRef.get().getResponseMessage()) + .containsExactly((byte) 'o', (byte) 'k', (byte) '!'); + } + + @Test + void handleResponse_enforcesResponseSizeLimit() { + OkHttpGrpcSender sender = createSender(2); + AtomicReference responseRef = new AtomicReference<>(); + byte[] frame = new byte[] {0, 0, 0, 0, 3, 'o', 'k', '!'}; Response response = new Response.Builder() .request(new Request.Builder().url("http://localhost/").build()) .protocol(Protocol.HTTP_2) .code(200) - .body(ResponseBody.create("", TEXT_PLAIN)) - .message("Retryable") - .trailers(() -> Headers.of(GRPC_STATUS, "14")) + .body(ResponseBody.create(frame, GRPC_MEDIA_TYPE)) + .message("HTTP message") + .header("grpc-status", "0") .build(); - assertTrue(OkHttpGrpcSender.isRetryable(response)); - assertThat(response.body().string()).isEmpty(); + sender.handleResponse(response, responseRef::set); + + assertThat(responseRef.get().getStatusCode()).isEqualTo(GrpcStatusCode.RESOURCE_EXHAUSTED); + assertThat(responseRef.get().getResponseMessage()).isEmpty(); + } + + private static OkHttpGrpcSender createSender(long maxResponseBodySize) { + return new OkHttpGrpcSender( + "http://localhost", + null, + Duration.ofSeconds(10), + Duration.ofSeconds(10), + Collections::emptyMap, + null, + null, + null, + null, + maxResponseBodySize, + null); } @Test @@ -114,17 +164,6 @@ void send_rejectedExecution_callsOnError() { assertThat(responseRef.get()).isNull(); } - private static Response createResponse(int httpCode, String grpcStatus, String message) { - return new Response.Builder() - .request(new Request.Builder().url("http://localhost/").build()) - .protocol(Protocol.HTTP_2) - .code(httpCode) - .body(ResponseBody.create("body", TEXT_PLAIN)) - .message(message) - .header(GRPC_STATUS, grpcStatus) - .build(); - } - @Test void shutdown_CompletableResultCodeShouldWaitForThreads() throws Exception { // This test verifies that shutdown() returns a CompletableResultCode that only From be1cea4ecf092f47ae770a6ebb329e556751fb56 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Wed, 30 Sep 2026 22:33:56 +0300 Subject: [PATCH 08/12] fix(okhttp): avoid retrying local response size errors MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../okhttp/internal/OkHttpGrpcSender.java | 20 +++++++++++++------ 1 file changed, 14 insertions(+), 6 deletions(-) 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 c5449860928..3a2f502ad93 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 @@ -45,6 +45,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeUnit; +import java.util.function.BiConsumer; import java.util.function.Consumer; import java.util.function.Supplier; import java.util.logging.Level; @@ -218,9 +219,10 @@ public void onFailure(Call call, IOException e) { public void onResponse(Call call, Response response) { handleResponse( response, - resolvedResponse -> { + (resolvedResponse, canRetry) -> { if (retryState != null && retryState.canRetry(attempt) + && canRetry && isRetryable(resolvedResponse) && retryState.backoff(OptionalLong.empty())) { sendAttempt( @@ -242,6 +244,10 @@ && isRetryable(resolvedResponse) } void handleResponse(Response response, Consumer onResponse) { + handleResponse(response, (resolvedResponse, ignored) -> onResponse.accept(resolvedResponse)); + } + + private void handleResponse(Response response, BiConsumer onResponse) { try (ResponseBody body = response.body()) { // A gRPC message frame has a 5-byte header: 1 compression-flag byte + 4 message-length // bytes. Read the header first so that the size limit applies to the message payload only, @@ -253,7 +259,8 @@ void handleResponse(Response response, Consumer onResponse) { } catch (IOException e) { logger.log(Level.FINE, "Invalid gRPC response frame", e); onResponse.accept( - ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), new byte[0])); + ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), new byte[0]), + false); return; } @@ -278,7 +285,7 @@ void handleResponse(Response response, Consumer onResponse) { } if (wireBuffer.size() > maxResponseBodySize) { - onResponse.accept(responseMessageTooLarge(maxResponseBodySize)); + onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false); return; } @@ -290,7 +297,7 @@ void handleResponse(Response response, Consumer onResponse) { // Compressed: validate the encoding and decompress with a post-decompression size limit String encoding = response.header("grpc-encoding"); if (!"gzip".equalsIgnoreCase(encoding)) { - onResponse.accept(responseUnsupportedGrpcEncoding(encoding)); + onResponse.accept(responseUnsupportedGrpcEncoding(encoding), false); return; } try { @@ -303,7 +310,7 @@ void handleResponse(Response response, Consumer onResponse) { } } if (decompressedBuffer.size() > maxResponseBodySize) { - onResponse.accept(responseMessageTooLarge(maxResponseBodySize)); + onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false); return; } bodyBytes = decompressedBuffer.readByteArray(); @@ -312,7 +319,8 @@ void handleResponse(Response response, Consumer onResponse) { } } onResponse.accept( - ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), bodyBytes)); + ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), bodyBytes), + true); } } From 3f95082a1ee30636bd59719796f78c8a546d9d5f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Thu, 1 Oct 2026 20:29:03 +0300 Subject: [PATCH 09/12] fix: retry trailer-only grpc responses MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../okhttp/internal/OkHttpGrpcSender.java | 8 +++--- .../okhttp/internal/OkHttpGrpcSenderTest.java | 27 +++++++++++++++++++ 2 files changed, 32 insertions(+), 3 deletions(-) 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 3a2f502ad93..00ed11be0df 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 @@ -247,7 +247,8 @@ void handleResponse(Response response, Consumer onResponse) { handleResponse(response, (resolvedResponse, ignored) -> onResponse.accept(resolvedResponse)); } - private void handleResponse(Response response, BiConsumer onResponse) { + // Visible for testing. + void handleResponse(Response response, BiConsumer onResponse) { try (ResponseBody body = response.body()) { // A gRPC message frame has a 5-byte header: 1 compression-flag byte + 4 message-length // bytes. Read the header first so that the size limit applies to the message payload only, @@ -258,9 +259,10 @@ private void handleResponse(Response response, BiConsumer body.source().skip(4); // message length — we bound reads by EOF instead } catch (IOException e) { logger.log(Level.FINE, "Invalid gRPC response frame", e); + GrpcResponse resolvedResponse = + ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), new byte[0]); onResponse.accept( - ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), new byte[0]), - false); + resolvedResponse, resolvedResponse.getStatusCode() != GrpcStatusCode.UNKNOWN); return; } diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java index 2e9bf518a8a..1f20839d0e0 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java @@ -29,6 +29,7 @@ import java.util.concurrent.SynchronousQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import javax.net.ssl.SSLContext; import javax.net.ssl.SSLException; @@ -98,6 +99,32 @@ void handleResponse_resolvesTrailerStatusAndMessageAfterConsumingBody() { .containsExactly((byte) 'o', (byte) 'k', (byte) '!'); } + @Test + void handleResponse_emptyBodyWithTrailerStatusIsRetryable() { + OkHttpGrpcSender sender = createSender(Long.MAX_VALUE); + AtomicReference responseRef = new AtomicReference<>(); + AtomicBoolean retryableRef = new AtomicBoolean(); + Response response = + new Response.Builder() + .request(new Request.Builder().url("http://localhost/").build()) + .protocol(Protocol.HTTP_2) + .code(200) + .body(ResponseBody.create(new byte[0], GRPC_MEDIA_TYPE)) + .message("HTTP message") + .trailers(() -> Headers.of("grpc-status", "14")) + .build(); + + sender.handleResponse( + response, + (resolvedResponse, canRetry) -> { + responseRef.set(resolvedResponse); + retryableRef.set(canRetry); + }); + + assertThat(responseRef.get().getStatusCode()).isEqualTo(GrpcStatusCode.UNAVAILABLE); + assertThat(retryableRef.get()).isTrue(); + } + @Test void handleResponse_enforcesResponseSizeLimit() { OkHttpGrpcSender sender = createSender(2); From 917e1d8280008774bb59c138936473167b3c7f8d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Thu, 8 Oct 2026 01:33:04 +0300 Subject: [PATCH 10/12] Bound OkHttp retries to each export timeout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../okhttp/internal/OkHttpGrpcSender.java | 80 ++++++++++--------- .../okhttp/internal/RetryInterceptor.java | 17 +++- .../sender/okhttp/internal/RetryState.java | 76 +++++++++++++++++- .../okhttp/internal/RetryInterceptorTest.java | 24 +++++- .../okhttp/internal/RetryStateTest.java | 70 ++++++++++++++++ 5 files changed, 223 insertions(+), 44 deletions(-) create mode 100644 exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java 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 00ed11be0df..781c2af5a43 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 @@ -34,6 +34,7 @@ import io.opentelemetry.sdk.common.export.MessageWriter; import io.opentelemetry.sdk.common.export.RetryPolicy; import java.io.IOException; +import java.net.SocketTimeoutException; import java.nio.ByteBuffer; import java.nio.charset.StandardCharsets; import java.time.Duration; @@ -88,6 +89,7 @@ public final class OkHttpGrpcSender implements GrpcSender { private final Supplier>> headersSupplier; private final long maxResponseBodySize; @Nullable private final RetryPolicy retryPolicy; + private final long timeoutNanos; /** Creates a new {@link OkHttpGrpcSender}. */ @SuppressWarnings("TooManyParameters") @@ -105,6 +107,7 @@ public OkHttpGrpcSender( @Nullable List enabledProtocols) { int callTimeoutMillis = (int) Math.min(timeout.toMillis(), Integer.MAX_VALUE); int connectTimeoutMillis = (int) Math.min(connectTimeout.toMillis(), Integer.MAX_VALUE); + this.timeoutNanos = TimeUnit.MILLISECONDS.toNanos(callTimeoutMillis); Dispatcher dispatcher; if (executorService == null) { @@ -159,6 +162,7 @@ public OkHttpGrpcSender( @Override public void send( MessageWriter messageWriter, Consumer onResponse, Consumer onError) { + RetryState retryState = newRetryState(); Request.Builder requestBuilder = new Request.Builder().url(url); Map> headers = headersSupplier.get(); @@ -174,12 +178,12 @@ public void send( } requestBuilder.post(new GrpcRequestBody(messageWriter, compressor)); - sendAttempt(requestBuilder, messageWriter, onResponse, onError, 0, newRetryState()); + sendAttempt(requestBuilder, messageWriter, onResponse, onError, 0, retryState); } @Nullable private RetryState newRetryState() { - return retryPolicy == null ? null : new RetryState(retryPolicy); + return retryPolicy == null ? null : new RetryState(retryPolicy, timeoutNanos); } private void sendAttempt( @@ -192,16 +196,41 @@ private void sendAttempt( try { InstrumentationUtil.suppressInstrumentation( - () -> - client - .newCall(requestBuilder.build()) - .enqueue( - new Callback() { - @Override - public void onFailure(Call call, IOException e) { + () -> { + Call call = client.newCall(requestBuilder.build()); + if (retryState != null && !retryState.configureCallTimeout(call)) { + onError.accept(new SocketTimeoutException("Call timed out")); + return; + } + call.enqueue( + new Callback() { + @Override + public void onFailure(Call call, IOException e) { + if (retryState != null + && retryState.canRetry(attempt) + && retryState.shouldRetryOnException(e) + && retryState.backoff(OptionalLong.empty())) { + sendAttempt( + requestBuilder, + messageWriter, + onResponse, + onError, + attempt + 1, + retryState); + } else { + onError.accept(e); + } + } + + @Override + public void onResponse(Call call, Response response) { + handleResponse( + response, + (resolvedResponse, canRetry) -> { if (retryState != null && retryState.canRetry(attempt) - && retryState.shouldRetryOnException(e) + && canRetry + && isRetryable(resolvedResponse) && retryState.backoff(OptionalLong.empty())) { sendAttempt( requestBuilder, @@ -211,33 +240,12 @@ public void onFailure(Call call, IOException e) { attempt + 1, retryState); } else { - onError.accept(e); + onResponse.accept(resolvedResponse); } - } - - @Override - public void onResponse(Call call, Response response) { - handleResponse( - response, - (resolvedResponse, canRetry) -> { - if (retryState != null - && retryState.canRetry(attempt) - && canRetry - && isRetryable(resolvedResponse) - && retryState.backoff(OptionalLong.empty())) { - sendAttempt( - requestBuilder, - messageWriter, - onResponse, - onError, - attempt + 1, - retryState); - } else { - onResponse.accept(resolvedResponse); - } - }); - } - })); + }); + } + }); + }); } catch (RejectedExecutionException e) { onError.accept(e); } 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 e8af7bf832f..80995ae3ebd 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 @@ -32,7 +32,9 @@ public final class RetryInterceptor implements Interceptor { private final RetryPolicy retryPolicy; private final Function isRetryable; private final Function retryDelayNanosExtractor; - private final RetryState retryState; + private final Predicate retryExceptionPredicate; + private final Sleeper sleeper; + private final Supplier randomJitter; /** Constructs a new retrier. */ public RetryInterceptor( @@ -61,11 +63,20 @@ public RetryInterceptor( this.retryPolicy = retryPolicy; this.isRetryable = isRetryable; this.retryDelayNanosExtractor = retryDelayNanosExtractor; - this.retryState = new RetryState(retryPolicy, retryExceptionPredicate, sleeper, randomJitter); + this.retryExceptionPredicate = retryExceptionPredicate; + this.sleeper = sleeper; + this.randomJitter = randomJitter; } @Override public Response intercept(Chain chain) throws IOException { + RetryState retryState = + new RetryState( + retryPolicy, + retryExceptionPredicate, + sleeper, + randomJitter, + chain.call().timeout().timeoutNanos()); Response response = null; IOException exception = null; int attempt = 0; @@ -149,7 +160,7 @@ static boolean isRetryableException(IOException e) { // Visible for testing boolean shouldRetryOnException(IOException e) { - return retryState.shouldRetryOnException(e); + return retryExceptionPredicate.test(e); } // Visible for testing diff --git a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java index ffcf261dc0c..ea8103f090a 100644 --- a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java @@ -14,8 +14,10 @@ import java.util.OptionalLong; import java.util.concurrent.ThreadLocalRandom; import java.util.concurrent.TimeUnit; +import java.util.function.LongSupplier; import java.util.function.Predicate; import java.util.function.Supplier; +import okhttp3.Call; /** Shared retry policy state and mechanics for OkHttp senders. */ final class RetryState { @@ -24,6 +26,9 @@ final class RetryState { private final Predicate retryExceptionPredicate; private final Sleeper sleeper; private final Supplier randomJitter; + private final long timeoutNanos; + private final LongSupplier nanoTime; + private final long startTimeNanos; private long nextBackoffNanos; RetryState(RetryPolicy retryPolicy) { @@ -33,7 +38,19 @@ final class RetryState { ? RetryState::isRetryableException : retryPolicy.getRetryExceptionPredicate(), TimeUnit.NANOSECONDS::sleep, - () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d)); + () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), + 0); + } + + RetryState(RetryPolicy retryPolicy, long timeoutNanos) { + this( + retryPolicy, + retryPolicy.getRetryExceptionPredicate() == null + ? RetryState::isRetryableException + : retryPolicy.getRetryExceptionPredicate(), + TimeUnit.NANOSECONDS::sleep, + () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), + timeoutNanos); } // Visible for testing. @@ -42,10 +59,40 @@ final class RetryState { Predicate retryExceptionPredicate, Sleeper sleeper, Supplier randomJitter) { + this(retryPolicy, retryExceptionPredicate, sleeper, randomJitter, 0); + } + + // Visible for testing. + RetryState( + RetryPolicy retryPolicy, + Predicate retryExceptionPredicate, + Sleeper sleeper, + Supplier randomJitter, + long timeoutNanos) { + this( + retryPolicy, + retryExceptionPredicate, + sleeper, + randomJitter, + timeoutNanos, + System::nanoTime); + } + + // Visible for testing. + RetryState( + RetryPolicy retryPolicy, + Predicate retryExceptionPredicate, + Sleeper sleeper, + Supplier randomJitter, + long timeoutNanos, + LongSupplier nanoTime) { this.retryPolicy = retryPolicy; this.retryExceptionPredicate = retryExceptionPredicate; this.sleeper = sleeper; this.randomJitter = randomJitter; + this.timeoutNanos = timeoutNanos; + this.nanoTime = nanoTime; + this.startTimeNanos = nanoTime.getAsLong(); this.nextBackoffNanos = retryPolicy.getInitialBackoff().toNanos(); } @@ -58,21 +105,44 @@ static double defaultRandomJitter() { } boolean backoff(OptionalLong retryDelayNanos) { + long remainingNanos = remainingNanos(); + if (remainingNanos <= 0) { + return false; + } long currentBackoffNanos = Math.min(nextBackoffNanos, retryPolicy.getMaxBackoff().toNanos()); - long backoffNanos = + long requestedBackoffNanos = retryDelayNanos.isPresent() ? retryDelayNanos.getAsLong() : (long) (randomJitter.get() * currentBackoffNanos); + long backoffNanos = Math.min(requestedBackoffNanos, remainingNanos); nextBackoffNanos = (long) (currentBackoffNanos * retryPolicy.getBackoffMultiplier()); try { sleeper.sleep(backoffNanos); - return true; + return remainingNanos() > 0; } catch (InterruptedException e) { Thread.currentThread().interrupt(); return false; } } + long remainingNanos() { + if (timeoutNanos <= 0) { + return Long.MAX_VALUE; + } + return Math.max(0, timeoutNanos - (nanoTime.getAsLong() - startTimeNanos)); + } + + boolean configureCallTimeout(Call call) { + long remainingNanos = remainingNanos(); + if (remainingNanos <= 0) { + return false; + } + if (timeoutNanos > 0) { + call.timeout().timeout(remainingNanos, TimeUnit.NANOSECONDS); + } + return true; + } + boolean canRetry(int attempt) { return attempt + 1 < retryPolicy.getMaxAttempts(); } 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 33968f19308..ec99d5b2a05 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 @@ -105,9 +105,10 @@ public boolean test(IOException e) { @Test void noRetryOnNullResponse() throws IOException { Interceptor.Chain chain = mock(Interceptor.Chain.class); + Request request = new Request.Builder().url(server.httpUri().toString()).build(); when(chain.proceed(any())).thenReturn(null); - when(chain.request()) - .thenReturn(new Request.Builder().url(server.httpUri().toString()).build()); + when(chain.request()).thenReturn(request); + when(chain.call()).thenReturn(client.newCall(request)); assertThatThrownBy( () -> { retrier.intercept(chain); @@ -132,6 +133,25 @@ void noRetry() throws Exception { verifyNoInteractions(sleeper); } + @Test + void backoffResetsForEachCall() throws Exception { + server.enqueue(HttpResponse.of(HttpStatus.INTERNAL_SERVER_ERROR)); + server.enqueue(HttpResponse.of(HttpStatus.OK)); + server.enqueue(HttpResponse.of(HttpStatus.INTERNAL_SERVER_ERROR)); + server.enqueue(HttpResponse.of(HttpStatus.OK)); + when(random.get()).thenReturn(1.0d); + doNothing().when(sleeper).sleep(anyLong()); + + try (Response first = sendRequest(); + Response second = sendRequest()) { + assertThat(first.isSuccessful()).isTrue(); + assertThat(second.isSuccessful()).isTrue(); + } + + verify(sleeper, times(2)).sleep(TimeUnit.SECONDS.toNanos(1)); + verify(sleeper, never()).sleep(TimeUnit.MILLISECONDS.toNanos(1600)); + } + @ParameterizedTest // Test is mostly same for 5 or more attempts since it's the max. We check the backoff timings and // handling of max attempts by checking both. diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java new file mode 100644 index 00000000000..29b022aca67 --- /dev/null +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java @@ -0,0 +1,70 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.exporter.sender.okhttp.internal; + +import static org.assertj.core.api.Assertions.assertThat; + +import io.opentelemetry.sdk.common.export.RetryPolicy; +import java.io.IOException; +import java.time.Duration; +import java.util.OptionalLong; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import okhttp3.Call; +import okhttp3.OkHttpClient; +import okhttp3.Request; +import org.junit.jupiter.api.Test; + +class RetryStateTest { + + @Test + void backoffIsCappedByRemainingTimeout() { + AtomicLong clock = new AtomicLong(); + AtomicLong sleptNanos = new AtomicLong(); + RetryState retryState = + new RetryState( + retryPolicy(), + (IOException e) -> true, + delay -> { + sleptNanos.set(delay); + clock.addAndGet(delay); + }, + () -> 1.0d, + TimeUnit.MILLISECONDS.toNanos(100), + clock::get); + clock.set(TimeUnit.MILLISECONDS.toNanos(40)); + + assertThat(retryState.backoff(OptionalLong.empty())).isFalse(); + assertThat(sleptNanos.get()).isEqualTo(TimeUnit.MILLISECONDS.toNanos(60)); + assertThat(retryState.remainingNanos()).isZero(); + } + + @Test + void callTimeoutUsesRemainingExportBudget() { + AtomicLong clock = new AtomicLong(); + RetryState retryState = + new RetryState( + retryPolicy(), + (IOException e) -> true, + delay -> {}, + () -> 1.0d, + TimeUnit.SECONDS.toNanos(1), + clock::get); + clock.set(TimeUnit.MILLISECONDS.toNanos(400)); + Call call = new OkHttpClient().newCall(new Request.Builder().url("http://localhost/").build()); + + assertThat(retryState.configureCallTimeout(call)).isTrue(); + assertThat(call.timeout().timeoutNanos()).isEqualTo(TimeUnit.MILLISECONDS.toNanos(600)); + } + + private static RetryPolicy retryPolicy() { + return RetryPolicy.builder() + .setInitialBackoff(Duration.ofMillis(300)) + .setMaxBackoff(Duration.ofMillis(300)) + .setMaxAttempts(3) + .build(); + } +} From bbadcdcbde6a4f96bcc8fd25a1dc166028ae1db2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Efe=20G=C3=B6kdemir?= Date: Thu, 8 Oct 2026 17:36:36 +0300 Subject: [PATCH 11/12] test: cover OkHttp retry timeout branches MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Efe Gökdemir --- .../okhttp/internal/RetryStateTest.java | 61 +++++++++++++++++++ 1 file changed, 61 insertions(+) diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java index 29b022aca67..89cb811b873 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java @@ -42,6 +42,24 @@ void backoffIsCappedByRemainingTimeout() { assertThat(retryState.remainingNanos()).isZero(); } + @Test + void backoffDoesNotSleepAfterTimeoutExpires() { + AtomicLong clock = new AtomicLong(); + AtomicLong sleepCalls = new AtomicLong(); + RetryState retryState = + new RetryState( + retryPolicy(), + (IOException e) -> true, + delay -> sleepCalls.incrementAndGet(), + () -> 1.0d, + TimeUnit.MILLISECONDS.toNanos(100), + clock::get); + clock.set(TimeUnit.MILLISECONDS.toNanos(100)); + + assertThat(retryState.backoff(OptionalLong.empty())).isFalse(); + assertThat(sleepCalls).hasValue(0); + } + @Test void callTimeoutUsesRemainingExportBudget() { AtomicLong clock = new AtomicLong(); @@ -60,6 +78,49 @@ void callTimeoutUsesRemainingExportBudget() { assertThat(call.timeout().timeoutNanos()).isEqualTo(TimeUnit.MILLISECONDS.toNanos(600)); } + @Test + void callTimeoutIsUnchangedWithoutExportTimeout() { + RetryState retryState = + new RetryState(retryPolicy(), (IOException e) -> true, delay -> {}, () -> 1.0d, 0); + Call call = new OkHttpClient().newCall(new Request.Builder().url("http://localhost/").build()); + + assertThat(retryState.configureCallTimeout(call)).isTrue(); + assertThat(call.timeout().timeoutNanos()).isZero(); + } + + @Test + void callTimeoutIsRejectedAfterExportBudgetExpires() { + AtomicLong clock = new AtomicLong(); + RetryState retryState = + new RetryState( + retryPolicy(), + (IOException e) -> true, + delay -> {}, + () -> 1.0d, + TimeUnit.SECONDS.toNanos(1), + clock::get); + clock.set(TimeUnit.SECONDS.toNanos(1)); + Call call = new OkHttpClient().newCall(new Request.Builder().url("http://localhost/").build()); + + assertThat(retryState.configureCallTimeout(call)).isFalse(); + } + + @Test + void retryPredicateAndAttemptLimitAreRespected() { + RetryState retryState = + new RetryState( + retryPolicy(), + exception -> exception.getMessage().startsWith("retry"), + delay -> {}, + () -> 1.0d); + + assertThat(retryState.canRetry(0)).isTrue(); + assertThat(retryState.canRetry(1)).isTrue(); + assertThat(retryState.canRetry(2)).isFalse(); + assertThat(retryState.shouldRetryOnException(new IOException("retry this"))).isTrue(); + assertThat(retryState.shouldRetryOnException(new IOException("stop"))).isFalse(); + } + private static RetryPolicy retryPolicy() { return RetryPolicy.builder() .setInitialBackoff(Duration.ofMillis(300)) From 023bcc0bc974813deab7d1028a78bfc650ac094f Mon Sep 17 00:00:00 2001 From: Jack Berg <34418638+jack-berg@users.noreply.github.com> Date: Thu, 8 Oct 2026 16:10:41 -0500 Subject: [PATCH 12/12] Add abstract test case, fix deadlines, and fix interaction with shutdown --- CHANGELOG.md | 5 -- .../AbstractGrpcTelemetryExporterTest.java | 35 +++++++++++++ .../okhttp/internal/OkHttpGrpcSender.java | 37 +++++++++++-- .../okhttp/internal/RetryInterceptor.java | 3 +- .../sender/okhttp/internal/RetryState.java | 52 ++----------------- .../okhttp/internal/OkHttpGrpcSenderTest.java | 48 ----------------- .../okhttp/internal/RetryStateTest.java | 17 +++++- 7 files changed, 87 insertions(+), 110 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 129c92ad4f1..66bf947bf05 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,11 +2,6 @@ ## Unreleased -### Exporters - -* OTLP gRPC: Retry responses that report a retryable status in trailers - ([#8854](https://github.com/open-telemetry/opentelemetry-java/pull/8854)). - ## Version 1.66.0 (2026-09-11) ### API diff --git a/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractGrpcTelemetryExporterTest.java b/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractGrpcTelemetryExporterTest.java index de279046b9b..be728f1c81e 100644 --- a/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractGrpcTelemetryExporterTest.java +++ b/exporters/otlp/testing-internal/src/main/java/io/opentelemetry/exporter/otlp/testing/internal/AbstractGrpcTelemetryExporterTest.java @@ -20,6 +20,7 @@ import com.google.protobuf.Message; import com.google.protobuf.Parser; import com.linecorp.armeria.common.HttpData; +import com.linecorp.armeria.common.HttpHeaders; import com.linecorp.armeria.common.HttpRequest; import com.linecorp.armeria.common.HttpResponse; import com.linecorp.armeria.common.HttpStatus; @@ -123,6 +124,9 @@ public abstract class AbstractGrpcTelemetryExporterTest { private static final ConcurrentLinkedQueue grpcResponses = new ConcurrentLinkedQueue<>(); + private static final ConcurrentLinkedQueue grpcTrailerResponses = + new ConcurrentLinkedQueue<>(); + private static volatile byte[] defaultResponseBytes = new byte[0]; private static final AtomicInteger attempts = new AtomicInteger(); @@ -214,6 +218,10 @@ public HttpResponse serve(ServiceRequestContext ctx, HttpRequest req) throws Exception { httpRequests.add(req); attempts.incrementAndGet(); + HttpResponse trailerResponse = grpcTrailerResponses.poll(); + if (trailerResponse != null) { + return trailerResponse; + } return unwrap().serve(ctx, req); } }); @@ -330,6 +338,7 @@ void tearDown() { void reset() { exportedResourceTelemetry.clear(); grpcResponses.clear(); + grpcTrailerResponses.clear(); attempts.set(0); httpRequests.clear(); grpcEncodingServerAttempts.set(0); @@ -983,6 +992,32 @@ void retryableError(int code) { assertThat(attempts).hasValue(2); } + @ParameterizedTest + @ValueSource(ints = {1, 4, 8, 10, 11, 14, 15}) + @SuppressLogger(GrpcExporter.class) + void retryableErrorInTrailers(int code) { + // grpc-java commits the RPC on initial response headers, preventing trailer-based retries. + assumeThat(exporter.unwrap()) + .extracting("delegate.grpcSender") + .matches(sender -> sender.getClass().getSimpleName().equals("OkHttpGrpcSender")); + + grpcTrailerResponses.add( + HttpResponse.of( + ResponseHeaders.builder(HttpStatus.OK) + .contentType(MediaType.parse("application/grpc+proto")) + .build(), + HttpData.empty(), + HttpHeaders.of("grpc-status", Integer.toString(code)))); + + CompletableResultCode result = + exporter + .export(Collections.singletonList(generateFakeTelemetry())) + .join(10, TimeUnit.SECONDS); + + assertThat(attempts).hasValue(2); + assertThat(result.isSuccess()).isTrue(); + } + @Test @SuppressLogger(GrpcExporter.class) void retryableError_tooManyAttempts() { 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 781c2af5a43..41a06a1b45e 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 @@ -82,6 +82,8 @@ public final class OkHttpGrpcSender implements GrpcSender { private static final String GRPC_STATUS = "grpc-status"; private static final String GRPC_MESSAGE = "grpc-message"; + private final Object shutdownLock = new Object(); + private boolean isShutdown; private final boolean managedExecutor; private final OkHttpClient client; private final HttpUrl url; @@ -183,7 +185,17 @@ public void send( @Nullable private RetryState newRetryState() { - return retryPolicy == null ? null : new RetryState(retryPolicy, timeoutNanos); + return retryPolicy == null + ? null + : new RetryState( + retryPolicy, + retryPolicy.getRetryExceptionPredicate() == null + ? RetryState::isRetryableException + : retryPolicy.getRetryExceptionPredicate(), + RetryState::defaultSleeper, + RetryState::defaultRandomJitter, + timeoutNanos, + System::nanoTime); } private void sendAttempt( @@ -202,11 +214,13 @@ private void sendAttempt( onError.accept(new SocketTimeoutException("Call timed out")); return; } - call.enqueue( + enqueue( + call, new Callback() { @Override public void onFailure(Call call, IOException e) { - if (retryState != null + if (!call.isCanceled() + && retryState != null && retryState.canRetry(attempt) && retryState.shouldRetryOnException(e) && retryState.backoff(OptionalLong.empty())) { @@ -251,12 +265,22 @@ && isRetryable(resolvedResponse) } } + private void enqueue(Call call, Callback callback) { + synchronized (shutdownLock) { + if (!isShutdown) { + call.enqueue(callback); + return; + } + } + call.cancel(); + callback.onFailure(call, new IOException("Canceled")); + } + void handleResponse(Response response, Consumer onResponse) { handleResponse(response, (resolvedResponse, ignored) -> onResponse.accept(resolvedResponse)); } - // Visible for testing. - void handleResponse(Response response, BiConsumer onResponse) { + private void handleResponse(Response response, BiConsumer onResponse) { try (ResponseBody body = response.body()) { // A gRPC message frame has a 5-byte header: 1 compression-flag byte + 4 message-length // bytes. Read the header first so that the size limit applies to the message payload only, @@ -385,6 +409,9 @@ private static String grpcMessage(Response response) { @Override public CompletableResultCode shutdown() { + synchronized (shutdownLock) { + isShutdown = true; + } client.dispatcher().cancelAll(); client.connectionPool().evictAll(); 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 80995ae3ebd..fbb267a086d 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 @@ -76,7 +76,8 @@ public Response intercept(Chain chain) throws IOException { retryExceptionPredicate, sleeper, randomJitter, - chain.call().timeout().timeoutNanos()); + chain.call().timeout().timeoutNanos(), + System::nanoTime); Response response = null; IOException exception = null; int attempt = 0; diff --git a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java index ea8103f090a..3d49777fe01 100644 --- a/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java @@ -31,54 +31,6 @@ final class RetryState { private final long startTimeNanos; private long nextBackoffNanos; - RetryState(RetryPolicy retryPolicy) { - this( - retryPolicy, - retryPolicy.getRetryExceptionPredicate() == null - ? RetryState::isRetryableException - : retryPolicy.getRetryExceptionPredicate(), - TimeUnit.NANOSECONDS::sleep, - () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), - 0); - } - - RetryState(RetryPolicy retryPolicy, long timeoutNanos) { - this( - retryPolicy, - retryPolicy.getRetryExceptionPredicate() == null - ? RetryState::isRetryableException - : retryPolicy.getRetryExceptionPredicate(), - TimeUnit.NANOSECONDS::sleep, - () -> ThreadLocalRandom.current().nextDouble(0.8d, 1.2d), - timeoutNanos); - } - - // Visible for testing. - RetryState( - RetryPolicy retryPolicy, - Predicate retryExceptionPredicate, - Sleeper sleeper, - Supplier randomJitter) { - this(retryPolicy, retryExceptionPredicate, sleeper, randomJitter, 0); - } - - // Visible for testing. - RetryState( - RetryPolicy retryPolicy, - Predicate retryExceptionPredicate, - Sleeper sleeper, - Supplier randomJitter, - long timeoutNanos) { - this( - retryPolicy, - retryExceptionPredicate, - sleeper, - randomJitter, - timeoutNanos, - System::nanoTime); - } - - // Visible for testing. RetryState( RetryPolicy retryPolicy, Predicate retryExceptionPredicate, @@ -138,7 +90,9 @@ boolean configureCallTimeout(Call call) { return false; } if (timeoutNanos > 0) { - call.timeout().timeout(remainingNanos, TimeUnit.NANOSECONDS); + call.timeout() + .timeout(remainingNanos, TimeUnit.NANOSECONDS) + .deadlineNanoTime(startTimeNanos + timeoutNanos); } return true; } diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java index 1f20839d0e0..af5c1fb7c0f 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/OkHttpGrpcSenderTest.java @@ -29,7 +29,6 @@ import java.util.concurrent.SynchronousQueue; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import javax.net.ssl.SSLContext; import javax.net.ssl.SSLException; @@ -99,53 +98,6 @@ void handleResponse_resolvesTrailerStatusAndMessageAfterConsumingBody() { .containsExactly((byte) 'o', (byte) 'k', (byte) '!'); } - @Test - void handleResponse_emptyBodyWithTrailerStatusIsRetryable() { - OkHttpGrpcSender sender = createSender(Long.MAX_VALUE); - AtomicReference responseRef = new AtomicReference<>(); - AtomicBoolean retryableRef = new AtomicBoolean(); - Response response = - new Response.Builder() - .request(new Request.Builder().url("http://localhost/").build()) - .protocol(Protocol.HTTP_2) - .code(200) - .body(ResponseBody.create(new byte[0], GRPC_MEDIA_TYPE)) - .message("HTTP message") - .trailers(() -> Headers.of("grpc-status", "14")) - .build(); - - sender.handleResponse( - response, - (resolvedResponse, canRetry) -> { - responseRef.set(resolvedResponse); - retryableRef.set(canRetry); - }); - - assertThat(responseRef.get().getStatusCode()).isEqualTo(GrpcStatusCode.UNAVAILABLE); - assertThat(retryableRef.get()).isTrue(); - } - - @Test - void handleResponse_enforcesResponseSizeLimit() { - OkHttpGrpcSender sender = createSender(2); - AtomicReference responseRef = new AtomicReference<>(); - byte[] frame = new byte[] {0, 0, 0, 0, 3, 'o', 'k', '!'}; - Response response = - new Response.Builder() - .request(new Request.Builder().url("http://localhost/").build()) - .protocol(Protocol.HTTP_2) - .code(200) - .body(ResponseBody.create(frame, GRPC_MEDIA_TYPE)) - .message("HTTP message") - .header("grpc-status", "0") - .build(); - - sender.handleResponse(response, responseRef::set); - - assertThat(responseRef.get().getStatusCode()).isEqualTo(GrpcStatusCode.RESOURCE_EXHAUSTED); - assertThat(responseRef.get().getResponseMessage()).isEmpty(); - } - private static OkHttpGrpcSender createSender(long maxResponseBodySize) { return new OkHttpGrpcSender( "http://localhost", diff --git a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java index 89cb811b873..84886390a48 100644 --- a/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java @@ -76,16 +76,27 @@ void callTimeoutUsesRemainingExportBudget() { assertThat(retryState.configureCallTimeout(call)).isTrue(); assertThat(call.timeout().timeoutNanos()).isEqualTo(TimeUnit.MILLISECONDS.toNanos(600)); + assertThat(call.timeout().hasDeadline()).isTrue(); + assertThat(call.timeout().deadlineNanoTime()).isEqualTo(TimeUnit.SECONDS.toNanos(1)); + + clock.set(TimeUnit.MILLISECONDS.toNanos(700)); + Call retry = new OkHttpClient().newCall(new Request.Builder().url("http://localhost/").build()); + + assertThat(retryState.configureCallTimeout(retry)).isTrue(); + assertThat(retry.timeout().timeoutNanos()).isEqualTo(TimeUnit.MILLISECONDS.toNanos(300)); + assertThat(retry.timeout().deadlineNanoTime()).isEqualTo(call.timeout().deadlineNanoTime()); } @Test void callTimeoutIsUnchangedWithoutExportTimeout() { RetryState retryState = - new RetryState(retryPolicy(), (IOException e) -> true, delay -> {}, () -> 1.0d, 0); + new RetryState( + retryPolicy(), (IOException e) -> true, delay -> {}, () -> 1.0d, 0, System::nanoTime); Call call = new OkHttpClient().newCall(new Request.Builder().url("http://localhost/").build()); assertThat(retryState.configureCallTimeout(call)).isTrue(); assertThat(call.timeout().timeoutNanos()).isZero(); + assertThat(call.timeout().hasDeadline()).isFalse(); } @Test @@ -112,7 +123,9 @@ void retryPredicateAndAttemptLimitAreRespected() { retryPolicy(), exception -> exception.getMessage().startsWith("retry"), delay -> {}, - () -> 1.0d); + () -> 1.0d, + 0, + System::nanoTime); assertThat(retryState.canRetry(0)).isTrue(); assertThat(retryState.canRetry(1)).isTrue();