Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import com.google.api.core.ApiFuture;
import com.google.api.core.ApiFutures;
import java.io.InterruptedIOException;
import java.net.SocketTimeoutException;
import java.nio.channels.ClosedByInterruptException;
import java.util.concurrent.Callable;
import org.jspecify.annotations.NullMarked;
Expand Down Expand Up @@ -103,6 +104,11 @@ public ApiFuture<ResponseT> submit(RetryingFuture<ResponseT> retryingFuture) {
sleep(retryingFuture.getAttemptSettings().getRandomizedRetryDelayDuration());
ResponseT response = retryingFuture.getCallable().call();
retryingFuture.setAttemptFuture(ApiFutures.immediateFuture(response));
} catch (SocketTimeoutException e) {
// A connect or read timeout is an InterruptedIOException, but no thread was interrupted.
// Setting the interrupt flag here makes setAttemptFuture throw a new InterruptedException,
// which drops this exception and stops the retry. Let the retry algorithm judge it.
retryingFuture.setAttemptFuture(ApiFutures.<ResponseT>immediateFailedFuture(e));
} catch (InterruptedException | InterruptedIOException | ClosedByInterruptException e) {
Thread.currentThread().interrupt();
retryingFuture.setAttemptFuture(ApiFutures.<ResponseT>immediateFailedFuture(e));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,19 @@
*/
package com.google.api.gax.retrying;

import static com.google.api.gax.retrying.FailingCallable.FAST_RETRY_SETTINGS;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertThrows;

import com.google.api.core.CurrentMillisClock;
import java.io.InterruptedIOException;
import java.net.SocketTimeoutException;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.atomic.AtomicInteger;
import org.junit.jupiter.api.Test;

class DirectRetryingExecutorTest extends AbstractRetryingExecutorTest {

Expand All @@ -45,4 +57,61 @@ protected RetryAlgorithm<String> getAlgorithm(
new TestResultRetryAlgorithm<String>(apocalypseCountDown, apocalypseException),
new ExponentialRetryAlgorithm(retrySettings, CurrentMillisClock.getDefaultClock()));
}

/**
* Runs {@code callable} under an algorithm that retries a {@link SocketTimeoutException} and
* nothing else, so the test sees which exception the algorithm was given.
*/
private RetryingFuture<String> runRetryingTimeouts(Callable<String> callable) {
setUp(false);
RetryAlgorithm<String> algorithm =
new RetryAlgorithm<>(
new BasicResultRetryAlgorithm<String>() {
@Override
public boolean shouldRetry(Throwable prevThrowable, String prevResponse) {
return prevThrowable instanceof SocketTimeoutException;
}
},
new ExponentialRetryAlgorithm(
FAST_RETRY_SETTINGS, CurrentMillisClock.getDefaultClock()));
RetryingExecutorWithContext<String> executor = getExecutor(algorithm);
RetryingFuture<String> future = executor.createFuture(callable, retryingContext);
future.setAttemptFuture(executor.submit(future));
return future;
}

@Test
void testSocketTimeoutReachesTheRetryAlgorithm() throws Exception {
AtomicInteger calls = new AtomicInteger();
try {
RetryingFuture<String> future =
runRetryingTimeouts(
() -> {
if (calls.getAndIncrement() == 0) {
throw new SocketTimeoutException("Read timed out");
}
return "SUCCESS";
});
assertFalse(Thread.currentThread().isInterrupted());
assertEquals("SUCCESS", future.get());
assertEquals(2, calls.get());
} finally {
Thread.interrupted();
}
}

@Test
void testInterruptedIOExceptionStillFailsAsInterrupted() {
try {
RetryingFuture<String> future =
runRetryingTimeouts(
() -> {
throw new InterruptedIOException("interrupted");
});
ExecutionException e = assertThrows(ExecutionException.class, future::get);
assertInstanceOf(InterruptedException.class, e.getCause());
} finally {
Thread.interrupted();
}
}
Comment on lines +104 to +116

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

In testInterruptedIOExceptionStillFailsAsInterrupted, the test asserts that future.get() throws an ExecutionException with an InterruptedException as its cause. However, because DirectRetryingExecutor.submit calls Thread.currentThread().interrupt(), the current thread's interrupt flag is set.

If the interrupt flag is correctly preserved, any subsequent blocking call to future.get() (which delegates to Guava's AbstractFuture.get()) should immediately throw a standalone InterruptedException rather than an ExecutionException wrapping it, unless the interrupt flag was somehow cleared and lost during the future's completion callback execution.

To ensure that the thread's interrupted status is correctly preserved and not lost, we should explicitly assert that the current thread remains interrupted after runRetryingTimeouts completes, and then clear the flag. If the flag is currently being lost, we should investigate why and fix the underlying issue.

References
  1. In Java, do not swallow InterruptedException. When catching it, restore the thread's interrupted status by calling Thread.currentThread().interrupt() and handle the interruption appropriately.

}
Loading