Skip to content
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -87,6 +87,7 @@ public final class OkHttpGrpcSender implements GrpcSender {
@Nullable private final Compressor compressor;
private final Supplier<Map<String, List<String>>> headersSupplier;
private final long maxResponseBodySize;
@Nullable private final RetryPolicy retryPolicy;

/** Creates a new {@link OkHttpGrpcSender}. */
@SuppressWarnings("TooManyParameters")
Expand Down Expand Up @@ -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));
Expand Down Expand Up @@ -158,6 +153,7 @@ public OkHttpGrpcSender(
this.headersSupplier = headersSupplier;
this.url = HttpUrl.get(endpoint);
this.maxResponseBodySize = maxResponseBodySize;
this.retryPolicy = retryPolicy;
}

@Override
Expand All @@ -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<GrpcResponse> onResponse,
Consumer<Throwable> onError,
int attempt,
@Nullable RetryState retryState) {

try {
InstrumentationUtil.suppressInstrumentation(
Expand All @@ -188,20 +199,56 @@ 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) {
onError.accept(e);
}
}

private void handleResponse(Response response, Consumer<GrpcResponse> onResponse) {
void handleResponse(Response response, Consumer<GrpcResponse> onResponse) {
handleResponse(response, (resolvedResponse, ignored) -> onResponse.accept(resolvedResponse));
}

// Visible for testing.
void handleResponse(Response response, BiConsumer<GrpcResponse, Boolean> 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,
Expand All @@ -212,8 +259,10 @@ private void handleResponse(Response response, Consumer<GrpcResponse> 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;
}

Expand All @@ -238,7 +287,7 @@ private void handleResponse(Response response, Consumer<GrpcResponse> onResponse
}

if (wireBuffer.size() > maxResponseBodySize) {
onResponse.accept(responseMessageTooLarge(maxResponseBodySize));
onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false);
return;
}

Expand All @@ -250,7 +299,7 @@ private void handleResponse(Response response, Consumer<GrpcResponse> 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 {
Expand All @@ -263,7 +312,7 @@ private void handleResponse(Response response, Consumer<GrpcResponse> onResponse
}
}
if (decompressedBuffer.size() > maxResponseBodySize) {
onResponse.accept(responseMessageTooLarge(maxResponseBodySize));
onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false);
return;
}
bodyBytes = decompressedBuffer.readByteArray();
Expand All @@ -272,7 +321,8 @@ private void handleResponse(Response response, Consumer<GrpcResponse> onResponse
}
}
onResponse.accept(
ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), bodyBytes));
ImmutableGrpcResponse.create(grpcStatus(response), grpcMessage(response), bodyBytes),
true);
}
}

Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -38,9 +32,7 @@ public final class RetryInterceptor implements Interceptor {
private final RetryPolicy retryPolicy;
private final Function<Response, Boolean> isRetryable;
private final Function<Response, OptionalLong> retryDelayNanosExtractor;
private final Predicate<IOException> retryExceptionPredicate;
private final Sleeper sleeper;
private final Supplier<Double> randomJitter;
private final RetryState retryState;

/** Constructs a new retrier. */
public RetryInterceptor(
Expand All @@ -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
Expand All @@ -69,36 +61,23 @@ 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
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();
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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 {}
}
Loading
Loading