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 48b990cf5a0..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 @@ -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; @@ -45,6 +46,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 +63,6 @@ import okhttp3.OkHttpClient; import okhttp3.Protocol; import okhttp3.Request; -import okhttp3.RequestBody; import okhttp3.Response; import okhttp3.ResponseBody; import okhttp3.TlsVersion; @@ -81,12 +82,16 @@ 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; @Nullable private final Compressor compressor; private final Supplier>> headersSupplier; private final long maxResponseBodySize; + @Nullable private final RetryPolicy retryPolicy; + private final long timeoutNanos; /** Creates a new {@link OkHttpGrpcSender}. */ @SuppressWarnings("TooManyParameters") @@ -104,6 +109,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) { @@ -119,12 +125,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,11 +158,13 @@ public OkHttpGrpcSender( this.headersSupplier = headersSupplier; this.url = HttpUrl.get(endpoint); this.maxResponseBodySize = maxResponseBodySize; + this.retryPolicy = retryPolicy; } @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(); @@ -176,32 +178,109 @@ 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, retryState); + } + + @Nullable + private RetryState newRetryState() { + return retryPolicy == null + ? null + : new RetryState( + retryPolicy, + retryPolicy.getRetryExceptionPredicate() == null + ? RetryState::isRetryableException + : retryPolicy.getRetryExceptionPredicate(), + RetryState::defaultSleeper, + RetryState::defaultRandomJitter, + timeoutNanos, + System::nanoTime); + } + + private void sendAttempt( + Request.Builder requestBuilder, + MessageWriter messageWriter, + Consumer onResponse, + Consumer onError, + int attempt, + @Nullable RetryState retryState) { try { InstrumentationUtil.suppressInstrumentation( - () -> - client - .newCall(requestBuilder.build()) - .enqueue( - new Callback() { - @Override - public void onFailure(Call call, IOException e) { - onError.accept(e); - } - - @Override - public void onResponse(Call call, Response response) { - handleResponse(response, onResponse); - } - })); + () -> { + Call call = client.newCall(requestBuilder.build()); + if (retryState != null && !retryState.configureCallTimeout(call)) { + onError.accept(new SocketTimeoutException("Call timed out")); + return; + } + enqueue( + call, + new Callback() { + @Override + public void onFailure(Call call, IOException e) { + if (!call.isCanceled() + && 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) + && 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 onResponse) { + 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)); + } + + 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, @@ -212,8 +291,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 +319,7 @@ private void handleResponse(Response response, Consumer onResponse } if (wireBuffer.size() > maxResponseBodySize) { - onResponse.accept(responseMessageTooLarge(maxResponseBodySize)); + onResponse.accept(responseMessageTooLarge(maxResponseBodySize), false); return; } @@ -250,7 +331,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 +344,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 +353,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); } } @@ -327,6 +409,9 @@ private static String grpcMessage(Response response) { @Override public CompletableResultCode shutdown() { + synchronized (shutdownLock) { + isShutdown = true; + } client.dispatcher().cancelAll(); client.connectionPool().evictAll(); @@ -364,15 +449,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..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 @@ -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; @@ -52,10 +46,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 @@ -76,29 +70,26 @@ public RetryInterceptor( @Override public Response intercept(Chain chain) throws IOException { + RetryState retryState = + new RetryState( + retryPolicy, + retryExceptionPredicate, + sleeper, + randomJitter, + chain.call().timeout().timeoutNanos(), + System::nanoTime); 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 +120,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 +155,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 retryExceptionPredicate.test(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..3d49777fe01 --- /dev/null +++ b/exporters/sender/okhttp/src/main/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryState.java @@ -0,0 +1,126 @@ +/* + * 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.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 { + + private final RetryPolicy retryPolicy; + 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, + 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(); + } + + 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 remainingNanos = remainingNanos(); + if (remainingNanos <= 0) { + return false; + } + long currentBackoffNanos = Math.min(nextBackoffNanos, retryPolicy.getMaxBackoff().toNanos()); + 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 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) + .deadlineNanoTime(startTimeNanos + timeoutNanos); + } + return true; + } + + 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..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 @@ -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; @@ -43,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(); @@ -53,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); } @@ -62,11 +66,53 @@ 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) '!'); + } + + 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 +143,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 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..84886390a48 --- /dev/null +++ b/exporters/sender/okhttp/src/test/java/io/opentelemetry/exporter/sender/okhttp/internal/RetryStateTest.java @@ -0,0 +1,144 @@ +/* + * 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 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(); + 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)); + 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, 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 + 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, + 0, + System::nanoTime); + + 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)) + .setMaxBackoff(Duration.ofMillis(300)) + .setMaxAttempts(3) + .build(); + } +}