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 1/9] 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 2/9] 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 3/9] 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 4/9] 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 5/9] 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 6/9] 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 7/9] 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 8/9] 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 9/9] 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);