diff --git a/google-cloud-storage/src/main/java/com/google/cloud/storage/DefaultRetryContext.java b/google-cloud-storage/src/main/java/com/google/cloud/storage/DefaultRetryContext.java index 8353058114..095c01487b 100644 --- a/google-cloud-storage/src/main/java/com/google/cloud/storage/DefaultRetryContext.java +++ b/google-cloud-storage/src/main/java/com/google/cloud/storage/DefaultRetryContext.java @@ -27,6 +27,7 @@ import java.util.LinkedList; import java.util.List; import java.util.Locale; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; @@ -142,17 +143,26 @@ public void recordError(T t, OnSuccess onSuccess, OnFailur BackoffDuration backoffDuration = (BackoffDuration) nextBackoff; lastBackoffResult = nextBackoff; - pendingBackoff = - scheduledExecutorService.schedule( - () -> { - try { - onSuccess.onSuccess(); - } finally { - clearPendingBackoff(); - } - }, - backoffDuration.getDuration().toNanos(), - TimeUnit.NANOSECONDS); + try { + pendingBackoff = + scheduledExecutorService.schedule( + () -> { + try { + onSuccess.onSuccess(); + } finally { + clearPendingBackoff(); + } + }, + backoffDuration.getDuration().toNanos(), + TimeUnit.NANOSECONDS); + } catch (RejectedExecutionException e) { + InterruptedBackoffComment comment = + new InterruptedBackoffComment( + "Interrupted backoff -- unretryable error due to executor service shutdown"); + comment.addSuppressed(e); + t.addSuppressed(comment); + onFailure.onFailure(t); + } } else { String msg = String.format( diff --git a/google-cloud-storage/src/main/java/com/google/cloud/storage/RetryContext.java b/google-cloud-storage/src/main/java/com/google/cloud/storage/RetryContext.java index 7afd38d82e..caf8144168 100644 --- a/google-cloud-storage/src/main/java/com/google/cloud/storage/RetryContext.java +++ b/google-cloud-storage/src/main/java/com/google/cloud/storage/RetryContext.java @@ -16,6 +16,8 @@ package com.google.cloud.storage; +import static java.util.Objects.requireNonNull; + import com.google.api.client.util.Sleeper; import com.google.api.core.ApiClock; import com.google.api.core.ApiFuture; @@ -40,6 +42,7 @@ import java.util.concurrent.ScheduledFuture; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; +import org.checkerframework.checker.nullness.qual.NonNull; @InternalApi @InternalExtensionOnly @@ -131,6 +134,16 @@ static BackoffComment of(String message) { } } + final class InterruptedBackoffComment extends Throwable { + InterruptedBackoffComment(@NonNull String message) { + super( + requireNonNull(message, "message must be non null"), + /* cause= */ null, + /* enableSuppression= */ true, + /* writableStackTrace= */ false); + } + } + final class DirectScheduledExecutorService implements ScheduledExecutorService { private static final DirectScheduledExecutorService INSTANCE = new DirectScheduledExecutorService(Sleeper.DEFAULT, NanoClock.getDefaultClock()); diff --git a/google-cloud-storage/src/test/java/com/google/cloud/storage/RetryContextTest.java b/google-cloud-storage/src/test/java/com/google/cloud/storage/RetryContextTest.java index e0fd73dee7..e4d2ad839f 100644 --- a/google-cloud-storage/src/test/java/com/google/cloud/storage/RetryContextTest.java +++ b/google-cloud-storage/src/test/java/com/google/cloud/storage/RetryContextTest.java @@ -18,6 +18,10 @@ import static com.google.cloud.storage.TestUtils.assertAll; import static com.google.common.truth.Truth.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; import com.google.api.core.ApiClock; import com.google.api.core.NanoClock; @@ -31,6 +35,8 @@ import com.google.cloud.RetryHelper; import com.google.cloud.RetryHelper.RetryHelperException; import com.google.cloud.storage.Backoff.Jitterer; +import com.google.cloud.storage.RetryContext.BackoffComment; +import com.google.cloud.storage.RetryContext.InterruptedBackoffComment; import com.google.cloud.storage.RetryContext.OnFailure; import com.google.cloud.storage.RetryContext.OnSuccess; import com.google.cloud.storage.Retrying.RetryingDependencies; @@ -42,6 +48,7 @@ import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executors; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -423,6 +430,24 @@ public void resetAlsoResetsBackoffState() throws Exception { }); } + @Test + public void rejectedExecutionException_funneledToOnFailureHandlerAsSuppressedException() { + ScheduledExecutorService exec = mock(ScheduledExecutorService.class); + RejectedExecutionException alreadyShutdown = new RejectedExecutionException("already shutdown"); + when(exec.schedule(any(Runnable.class), anyLong(), any())).thenThrow(alreadyShutdown); + Throwable t1 = new RuntimeException("{err1}", new Throwable("{err1Cause}")); + RetryContext ctx = + RetryContext.of(exec, maxAttempts(2), Retrying.alwaysRetry(), Jitterer.noJitter()); + + AtomicReference err1 = new AtomicReference<>(); + ctx.recordError(t1, failOnSuccess(), err1::set); + Throwable t = err1.get(); + assertThat(t).isNotNull(); + assertThat(t.getSuppressed()[0]).isInstanceOf(BackoffComment.class); + assertThat(t.getSuppressed()[1]).isInstanceOf(InterruptedBackoffComment.class); + assertThat(t.getSuppressed()[1].getSuppressed()[0]).isSameInstanceAs(alreadyShutdown); + } + private static ApiException apiException(Code code, String message) { return ApiExceptionFactory.createException(message, null, GrpcStatusCode.of(code), false); }