diff --git a/CHANGELOG.md b/CHANGELOG.md index 66bf947bf05..129c92ad4f1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,11 @@ ## 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/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..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 @@ -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; @@ -61,7 +62,6 @@ import okhttp3.OkHttpClient; import okhttp3.Protocol; import okhttp3.Request; -import okhttp3.RequestBody; import okhttp3.Response; import okhttp3.ResponseBody; import okhttp3.TlsVersion; @@ -87,6 +87,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") @@ -119,12 +120,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())); - } - boolean isPlainHttp = endpoint.startsWith("http://"); if (isPlainHttp) { clientBuilder.connectionSpecs(Collections.singletonList(ConnectionSpec.CLEARTEXT)); @@ -158,6 +153,7 @@ public OkHttpGrpcSender( this.headersSupplier = headersSupplier; this.url = HttpUrl.get(endpoint); this.maxResponseBodySize = maxResponseBodySize; + this.retryPolicy = retryPolicy; } @Override @@ -176,8 +172,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( @@ -188,12 +199,43 @@ 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, 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) { @@ -201,7 +243,12 @@ public void onResponse(Call call, Response response) { } } - private void handleResponse(Response response, Consumer onResponse) { + void handleResponse(Response response, Consumer onResponse) { + handleResponse(response, (resolvedResponse, ignored) -> onResponse.accept(resolvedResponse)); + } + + // 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, @@ -212,8 +259,10 @@ private void handleResponse(Response response, Consumer onResponse 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])); + resolvedResponse, resolvedResponse.getStatusCode() != GrpcStatusCode.UNKNOWN); return; } @@ -238,7 +287,7 @@ private void handleResponse(Response response, Consumer onResponse } if (wireBuffer.size() > maxResponseBodySize) { - onResponse.accept(responseMessageTooLarge(maxResponseBodySize)); + onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false); return; } @@ -250,7 +299,7 @@ private 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 { @@ -263,7 +312,7 @@ private void handleResponse(Response response, Consumer onResponse } } if (decompressedBuffer.size() > maxResponseBodySize) { - onResponse.accept(responseMessageTooLarge(maxResponseBodySize)); + onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false); return; } bodyBytes = decompressedBuffer.readByteArray(); @@ -272,7 +321,8 @@ private void handleResponse(Response response, Consumer onResponse } } onResponse.accept( - ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), bodyBytes)); + ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), bodyBytes), + true); } } @@ -364,15 +414,10 @@ public CompletableResultCode shutdown() { return CompletableResultCode.ofSuccess(); } - /** 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; - } - return RetryUtil.retryableGrpcStatusCodes().contains(grpcStatus); + /** 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 145cf65b70f..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,9 +32,7 @@ public final class RetryInterceptor implements Interceptor { private final RetryPolicy retryPolicy; private final Function isRetryable; private final Function retryDelayNanosExtractor; - private final Predicate retryExceptionPredicate; - private final Sleeper sleeper; - private final Supplier randomJitter; + private final RetryState retryState; /** Constructs a new retrier. */ public RetryInterceptor( @@ -52,10 +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)); + RetryState::defaultSleeper, + RetryState::defaultRandomJitter); } // Visible for testing @@ -69,9 +61,7 @@ public RetryInterceptor( this.retryPolicy = retryPolicy; this.isRetryable = isRetryable; this.retryDelayNanosExtractor = retryDelayNanosExtractor; - this.retryExceptionPredicate = retryExceptionPredicate; - this.sleeper = sleeper; - this.randomJitter = randomJitter; + this.retryState = new RetryState(retryPolicy, retryExceptionPredicate, sleeper, randomJitter); } @Override @@ -79,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(); @@ -129,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, @@ -164,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 f68350ca3b0..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,9 +29,11 @@ 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; +import okhttp3.Headers; import okhttp3.MediaType; import okhttp3.Protocol; import okhttp3.Request; @@ -43,8 +45,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(); @@ -53,7 +54,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); } @@ -62,11 +67,100 @@ 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 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_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", + null, + Duration.ofSeconds(10), + Duration.ofSeconds(10), + Collections::emptyMap, + null, + null, + null, + null, + maxResponseBodySize, + null); + } + @Test void send_rejectedExecution_callsOnError() { ThreadPoolExecutor executor = @@ -97,17 +191,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